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.

Advertisement

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.

Source processorListSFTP, ConsumeKafkaConnection queueback pressureTransformrecord reader/writerSink processorPutHDFS, PublishKafkaFlowFile repositorywrite-ahead log of attributesContent repositoryappend-only claims on diskProvenance repositoryevent history per FlowFilestatebyteseventsCluster coordinatorelected; replicates flow editsPrimary noderuns primary-only processorsEvery node runs the same flow on its own data; the repositories are local to each node
A NiFi node: processors exchange FlowFiles through queued connections. Attributes and queue state are journalled in the FlowFile repository, payload bytes live in the content repository, and every action is recorded in the provenance repository. In a cluster each node holds its own repositories.

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.

Advertisement

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.

SymptomWhat it meansFirst response
Queue full, upstream stoppedDownstream is slower than the sourceFind the slowest processor by task time; add concurrent tasks only if its target can take more load
Queue swappingBacklog exceeds the in-memory thresholdTreat as an incident; check the sink, not the heap
Failure queue growingBad data or a rejecting sinkInspect provenance for a sample; fix routing, then replay
Content repository near fullArchive retention or a stalled queue holding claimsLower 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:

  1. ListSFTP on the primary node only, so the listing is not duplicated across the cluster. It keeps state on which files it has already listed.
  2. A load-balanced connection (round robin) to FetchSFTP, so every node fetches a share of the files in parallel.
  3. CompressContent in decompress mode, then UpdateAttribute to derive ingest.date from the filename with the Expression Language, for example ${filename:substringBefore('_')}.
  4. ConvertRecord with a CSV reader and a Parquet writer sharing a schema from a schema registry controller service. Rows that do not parse go to failure, which routes to a quarantine directory, not to a dead end.
  5. PutHDFS with directory /data/vendor/dt=${ingest.date} and a conflict resolution of fail, so a replayed file cannot silently overwrite good data.
  6. PublishKafka sending 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 modeCausePrevention
Node out of memoryProcessor reads whole content into heap, or millions of tiny FlowFilesStream content; use record processors; split carefully
Content repository fills and node stopsLong downstream outage, oversized archiveSize for the outage; cap archive by percentage; alert early
Duplicate files at the sinkAt-least-once replay after a crashIdempotent sink paths, conflict policy, dedupe downstream
Listing runs on every nodeListing processor not set to primary-onlyPrimary-only scheduling plus a load-balanced connection
Data stuck after node lossRepositories are node-localDurable disks; restore the node rather than replace it
Silent data lossFlowFile expiration or auto-terminated failure relationshipsRoute 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.
Key takeaway: NiFi moves data as FlowFiles whose attributes are journalled in a write-ahead log and whose content sits in an append-only repository, so routing is cheap and every action is recorded in provenance. Sessions make each processor step transactional and the guarantee at-least-once, and back pressure on durable queues pushes slowdowns back to the sources instead of losing data. Run it with record-oriented processors, primary-only listings, separate durable repository disks, monitored failure paths and idempotent sinks.