Apache NiFi is a dataflow engine: you draw a graph of processors on a canvas, and NiFi moves data along it, durably, with back pressure, retries and a full audit trail of what happened to every piece of data. It sits in the same part of a Hadoop estate as Flume and Sqoop, the edge where files, messages and database rows arrive, but it is general-purpose rather than tied to one source or sink, and it is designed to be changed while it runs.
A NiFi flow is a live system you edit in place, not a compiled job. Running one well means knowing what is on disk, what a processor promises when it finishes, and where data waits when a sink slows down. This article explains those mechanics, builds a worked flow into HDFS, and ends with operational rules.
What NiFi is for
Most ingest problems share a shape: data arrives from many places in many formats at uneven rates, needs light routing and transformation, and must reach a few destinations without loss even when one is down for an hour. Written as code, each pipeline needs its own queue, retry policy, dead-letter path, metrics and audit log. NiFi moves those concerns into the platform, so a pipeline reduces to processors and wiring.
NiFi is not a stream processor for stateful joins and windows (use Flink or Spark), nor a broker with long retention and replay for many consumers (use Kafka). The healthy architecture is NiFi at the edge doing collection, routing, conversion and enrichment, handing data to Kafka or storage.
FlowFiles and the three repositories
The unit of data in NiFi is the FlowFile: a small map of string attributes (filename, path, uuid, mime.type and whatever processors add) plus a pointer to its content. The content itself is not carried around. It lives in the content repository, an append-only store on local disk where many small payloads are packed into larger files called resource claims. Each FlowFile references a slice of a claim by offset and length.
That separation is why NiFi can route a 2 GB file between ten processors cheaply: routing, attribute updates and cloning copy pointers, not bytes. When a processor modifies content, the new bytes are appended as a new claim and the old one is released once unreferenced, optionally into an archive (bounded by time and disk percentage in nifi.properties) so provenance can replay it.
The FlowFile repository is a write-ahead log of FlowFile state: which queue each FlowFile is in and what its attributes are. It is small but its latency gates throughput, because every session commit writes to it. On restart NiFi replays it to rebuild every queue exactly as it was. The provenance repository records an event for every action: RECEIVE, CREATE, ATTRIBUTES_MODIFIED, CONTENT_MODIFIED, ROUTE, FORK, JOIN, SEND, DROP and others. It is indexed so you can search by attribute or FlowFile UUID and see the full lineage of one record, including which upstream file it was split from.
Put the three repositories on separate disks: content I/O is large and sequential, FlowFile repository I/O is small synchronous writes, and provenance indexing is seek-heavy.
Processors, sessions and scheduling
A processor does its work in onTrigger, which NiFi calls on a scheduler thread. The processor pulls FlowFiles from its input queue through a ProcessSession, reads or rewrites content, sets attributes and transfers each FlowFile to a named relationship such as success or failure. Nothing becomes visible to the rest of the flow until the session commits. If the processor throws, the session rolls back and the input FlowFiles return to the queue untouched. This is the transactional unit of NiFi, and it is what makes the delivery guarantee at-least-once: a crash after the downstream write but before the commit replays that input on restart.
Here is a minimal custom processor that tags a FlowFile with a SHA-256 of its content. It shows the contract: get, read through a callback, add an attribute, transfer to exactly one relationship.
@Tags({"hash", "content"})
@CapabilityDescription("Adds a content.sha256 attribute to each FlowFile.")
public class HashContent extends AbstractProcessor {
public static final Relationship REL_SUCCESS = new Relationship.Builder()
.name("success").description("Hashed FlowFiles").build();
public static final Relationship REL_FAILURE = new Relationship.Builder()
.name("failure").description("Content could not be read").build();
@Override
public Set<Relationship> getRelationships() {
return Set.of(REL_SUCCESS, REL_FAILURE);
}
@Override
public void onTrigger(ProcessContext context, ProcessSession session) {
FlowFile flowFile = session.get();
if (flowFile == null) {
return; // queue was empty for this trigger
}
try {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
session.read(flowFile, in -> {
byte[] buf = new byte[64 * 1024];
int n;
while ((n = in.read(buf)) != -1) {
digest.update(buf, 0, n); // streamed: never load the whole file
}
});
String hex = HexFormat.of().formatHex(digest.digest());
flowFile = session.putAttribute(flowFile, "content.sha256", hex);
session.transfer(flowFile, REL_SUCCESS);
} catch (Exception e) {
getLogger().error("Hashing failed for {}", flowFile, e);
session.transfer(flowFile, REL_FAILURE); // route, do not drop
}
// AbstractProcessor commits the session after onTrigger returns.
}
}Two habits matter. Stream content through the callback rather than into a byte array, or one oversized file takes down the JVM for every flow on the node. And route bad data to a relationship instead of throwing: an exception rolls back and retries the same input forever.
Scheduling is per processor: a run schedule, a number of concurrent tasks, and for some processors a run duration that batches many FlowFiles into one commit. All threads come from one shared timer-driven pool. NiFi 2.x removed the event-driven pool, so 1.x flows using event-driven scheduling must be switched to timer-driven.
Connections and back pressure
Every connection between two processors is a durable queue with two back-pressure thresholds: an object count and a data size (recent releases default new connections to 10,000 FlowFiles and 1 GB). When either is exceeded, NiFi stops scheduling the upstream processor. Pressure therefore propagates backwards hop by hop until it reaches a source processor, which stops pulling from SFTP or Kafka. Data waits where it was, on durable storage, instead of being dropped or piling up in memory.
Queued FlowFile records live in heap, so beyond a swap threshold NiFi writes batches of them to swap files. Swapping keeps the node alive but is a signal to investigate, not a setting to raise.
Connections can also carry prioritisers (oldest first, by a priority attribute and so on) and a FlowFile expiration that drops waiting data. Expiration suits data that is worthless when late, such as live telemetry; elsewhere it silently turns back pressure into loss.
| Symptom | What it means | First response |
|---|---|---|
| Queue full, upstream stopped | Downstream is slower than the source | Find the slowest processor by task time; add concurrent tasks only if its target can take more load |
| Queue swapping | Backlog exceeds the in-memory threshold | Treat as an incident; check the sink, not the heap |
| Failure queue growing | Bad data or a rejecting sink | Inspect provenance for a sample; fix routing, then replay |
| Content repository near full | Archive retention or a stalled queue holding claims | Lower archive retention; drain or expire the stalled queue |
Worked example: vendor files into HDFS
A vendor drops gzipped CSV files on an SFTP server every few minutes. You need them as Parquet in HDFS, partitioned by date, with bad rows quarantined and a Kafka notification when each file lands. The flow:
ListSFTPon the primary node only, so the listing is not duplicated across the cluster. It keeps state on which files it has already listed.- A load-balanced connection (round robin) to
FetchSFTP, so every node fetches a share of the files in parallel. CompressContentin decompress mode, thenUpdateAttributeto deriveingest.datefrom the filename with the Expression Language, for example${filename:substringBefore('_')}.ConvertRecordwith a CSV reader and a Parquet writer sharing a schema from a schema registry controller service. Rows that do not parse go tofailure, which routes to a quarantine directory, not to a dead end.PutHDFSwith directory/data/vendor/dt=${ingest.date}and a conflict resolution of fail, so a replayed file cannot silently overwrite good data.PublishKafkasending a small JSON notification with the HDFS path and row count attribute.
Record-oriented processors are the key choice. One FlowFile per row would mean millions of FlowFile records, each with repository and provenance overhead; a record reader and writer streams rows through one FlowFile. Likewise merge small inputs with MergeRecord before HDFS, because tiny files are a NameNode problem (see the HDFS small-files problem).
Because delivery is at-least-once, deterministic output names plus fail-on-conflict make a replayed file visible instead of silent; Kafka consumers should deduplicate on the path.
Clustering and load-balanced connections
A NiFi cluster is a set of nodes that all run the same flow. One node is elected cluster coordinator; it accepts flow edits from the UI or API and replicates them to every node, and it tracks node heartbeats. One node is elected primary node, the only one that runs primary-only processors such as listings. Elections and cluster state traditionally use ZooKeeper; 2.x also offers a Kubernetes-native option, so check which your deployment uses before planning failover.
Data is not shared between nodes: a FlowFile lives on one node until a load-balanced connection moves it. Round robin spreads work, partition by attribute sends equal attribute values to the same node, and single node funnels everything to one node. If a node dies, the FlowFiles on its disks wait for it to return; they are not re-processed elsewhere. That is the key fact for recovery planning, and why repository disks must be durable volumes.
Parameters, versioning and NiFi 2.x
Configuration that differs between environments, such as hostnames, paths and credentials, belongs in parameter contexts, referenced in properties as #{hdfs.base.dir}. Sensitive parameters are encrypted at rest and never exported. NiFi 2.x removed the older variable registry, so parameters are now the only mechanism; migrating a 1.x flow means converting variables to parameters before the upgrade.
Process groups are the unit of versioning: commit a group to a flow registry, and another environment imports the same version bound to its own parameter context. Treat it like code: review diffs, promote versions, and avoid production hand edits, which the UI flags as local changes.
NiFi 2.x also requires Java 21 and adds a native Python processor API. Python processors run in separate processes managed by NiFi, which adds overhead; for high-volume paths a record-oriented Java processor remains the efficient choice.
Operating NiFi
Size the node around disks first. The content repository needs room for the peak backlog you are willing to absorb during a downstream outage, plus archive. If the sink can be down for two hours and inbound is 50 MB/s, that is about 360 GB of queued content per cluster before NiFi has to push back on the sources. Heap holds queued FlowFile records and attributes, not content, so it scales with FlowFile count and attribute size, which is another argument for record-oriented flows and short attributes.
Monitor queue depth and back-pressure state per connection, processor task time and bulletins (NiFi's error notices), repository disk usage and JVM health. Alert on sustained back pressure at the edge, growing failure queues and repository usage, not on CPU.
Secure it like the data plane it is. Run NiFi over TLS with real user authentication and per-process-group policies, keep sensitive values in sensitive parameters, and integrate with Apache Atlas if you need provenance-level lineage in your governance catalogue.
Failure modes and trade-offs
| Failure mode | Cause | Prevention |
|---|---|---|
| Node out of memory | Processor reads whole content into heap, or millions of tiny FlowFiles | Stream content; use record processors; split carefully |
| Content repository fills and node stops | Long downstream outage, oversized archive | Size for the outage; cap archive by percentage; alert early |
| Duplicate files at the sink | At-least-once replay after a crash | Idempotent sink paths, conflict policy, dedupe downstream |
| Listing runs on every node | Listing processor not set to primary-only | Primary-only scheduling plus a load-balanced connection |
| Data stuck after node loss | Repositories are node-local | Durable disks; restore the node rather than replace it |
| Silent data loss | FlowFile expiration or auto-terminated failure relationships | Route failures somewhere visible; use expiration only by design |
The central trade-off is visibility against per-item overhead: every FlowFile pays for repository writes and provenance, cheap per large file and ruinous per tiny record. Against Flume, NiFi has richer routing and live editing but more operational surface. For database change capture, a dedicated tool such as Debezium CDC to Kafka beats polling with NiFi.
What to do next
- Draw your current ingest paths and mark where NiFi would hold data during a sink outage; size the content repository for that duration.
- Put the content, FlowFile and provenance repositories on separate durable volumes.
- Convert per-row flows to record readers and writers; check that no processor buffers whole files in heap.
- Mark every listing or polling source primary-only and follow it with a load-balanced connection.
- Give every failure relationship a destination you monitor; remove FlowFile expiration where loss is not intended.
- Make sinks idempotent and version each process group in a registry with parameter contexts.
- If you are on 1.x, inventory variables and event-driven processors before planning the 2.x upgrade.