When an application writes a file to HDFS, the bytes do not go to the NameNode and they do not fan out from the client to three servers. They flow through a chain: the client streams to the first DataNode, which writes and forwards to the second, which writes and forwards to the third, while acknowledgements travel back up the same chain. This design, the write pipeline, decides HDFS write throughput, what "durable" means after a call returns, and how a write survives a DataNode dying halfway through a block.

This article walks through the pipeline from first principles: why a chain, how a block is allocated, what packets and checksums look like, how replica placement shapes the chain, what hflush and hsync actually promise, how pipeline recovery and lease recovery work, and what goes wrong in production. Configuration defaults quoted here come from the current hdfs-default.xml.

Writing one block: allocate, set up the chain, stream packets down, acks upHDFS clientDFSOutputStreamNameNodelease, addBlock, placement1. create / addBlock2. block id, GS, 3 targetsLease managerone writer per fileData queueAck queueDataNode 1writer's node, rack ADataNode 2rack BDataNode 3rack B, other node3. packetsforwardforwardackack4. acksPipeline recoverydrop bad node, bump generation stamp, resend ack queueEvery byte crosses each link once; only one hop (DataNode 1 to 2) crosses racks.A packet leaves the ack queue only when all live DataNodes in the chain have acknowledged it.
The HDFS write pipeline: the NameNode allocates, the client streams packets down a chain of DataNodes, and acks return up the chain.

Why a chain

Suppose a client must store 128 MB with three replicas. If it sent three copies itself, it would push 384 MB through its own network interface, and the client is usually a busy compute node whose uplink is the scarcest resource in the cluster. A primary-replica star moves that problem to the primary, which must send two copies. In a chain, every participant receives each byte once and sends it at most once, so no link carries more than one copy and the work is spread evenly. When thousands of tasks write at the same time, that even spread is what maximises total cluster throughput.

The price is latency and fragility. A packet is acknowledged only after it has travelled the whole chain and the acks have come back, and one slow or dead DataNode stalls everyone behind it. HDFS accepts that because its workload is large sequential writes, where throughput matters more than per-packet latency, and it limits the latency cost with cut-through forwarding: each DataNode forwards a packet downstream as it receives it, in parallel with writing it to disk, rather than waiting to store the whole block.

Why a chain

Suppose a client must store 128 MB with three replicas. If it sent three copies itself, it would push 384 MB through its own network interface, and the client is usually a busy compute node whose uplink is the scarcest resource in the cluster. A primary-replica star moves that problem to the primary, which must send two copies. In a chain, every participant receives each byte once and sends it at most once, so no link carries more than one copy and the work is spread evenly. When thousands of tasks write at the same time, that even spread is what maximises total cluster throughput.

The price is latency and fragility. A packet is acknowledged only after it has travelled the whole chain and the acks have come back, and one slow or dead DataNode stalls everyone behind it. HDFS accepts that because its workload is large sequential writes, where throughput matters more than per-packet latency, and it limits the latency cost with cut-through forwarding: each DataNode forwards a packet downstream as it receives it, in parallel with writing it to disk, rather than waiting to store the whole block.

A write, step by step

  1. Create. The client calls create on the NameNode, which checks permissions and quotas, adds the file to the namespace as under construction, and grants the client a lease, the right to be the file's only writer.
  2. Allocate. When the first bytes arrive, and again at every block boundary (128 MB by default, dfs.blocksize), the client calls addBlock. The NameNode picks target DataNodes with its placement policy, assigns a block id and an initial generation stamp, and returns the ordered list. It never handles data.
  3. Set up. The client opens a connection to the first DataNode and sends a write-block request naming the whole chain. Each DataNode connects to the next and passes the request on; once the last replies, setup acks flow back and the chain is open.
  4. Stream. The client cuts data into packets and sends them down the chain. Each DataNode forwards the packet, writes data and checksums to local disk, and returns an ack that includes the status of everything downstream.
  5. Finish the block. A final empty packet marks the end of the block. DataNodes finalise the replica and report it to the NameNode. The client either allocates the next block or, on close, calls complete, which succeeds once enough replicas have been reported (at least dfs.namenode.replication.min, default 1).
Configuration conf = new Configuration();
FileSystem fs = FileSystem.get(URI.create("hdfs://nn1:8020"), conf);

try (FSDataOutputStream out = fs.create(new Path("/logs/app/2026-10-05.log"),
        true,                    // overwrite
        4096,                    // client buffer
        (short) 3,               // replication
        128L * 1024 * 1024)) {   // block size
    for (LogRecord r : records) {
        out.write(r.toBytes());
        if (r.isCommitPoint()) {
            out.hsync();         // durable on all pipeline DataNodes before we continue
        }
    }
}                                // close(): flush, end-of-block packet, complete()

Packets, checksums and the two queues

Inside the client, DFSOutputStream splits the byte stream into chunks of 512 bytes (dfs.bytes-per-checksum) and computes a checksum for each, CRC32C on modern releases. Chunks are grouped into packets of about 64 KB (dfs.client-write-packet-size, default 65,536). Every packet carries a header with its sequence number, its offset in the block, and a last-packet flag, followed by the checksums and then the data.

Two queues drive the stream. The data queue holds packets waiting to be sent; a streamer thread takes them, sends them to the first DataNode, and moves them to the ack queue. A response thread reads acks and removes packets from the ack queue only when every DataNode in the chain reported success. That ack queue is the recovery ledger: anything still in it when the chain breaks is pushed back to the front of the data queue and resent on the rebuilt chain. The number of packets in flight is capped, so a slow pipeline pushes back on the writer by blocking write.

On each DataNode a receiver thread reads packets from upstream (in the normal case the last node in the chain verifies checksums, so corruption is caught before the ack), mirrors the packet downstream and writes it locally. A responder thread pairs downstream acks with local status and sends the combined ack upstream. Each concurrent block read or write uses a transfer thread, bounded by dfs.datanode.max.transfer.threads (default 4,096).

Placement shapes the pipeline

The default placement policy puts the first replica on the writer's own DataNode if the client runs on one, otherwise on a random node; the second on a node in a different rack; and the third on a different node in the same rack as the second. The pipeline is ordered to match, so data crosses racks only once, between the first and second DataNodes, and the third copy travels over the rack's own switch. The block still survives the loss of an entire rack.

All of this depends on rack awareness being configured. Without a topology script, every node is in /default-rack, the policy cannot tell racks apart, and three replicas can land in one rack, so one switch failure loses data. The NameNode also skips nodes that are decommissioning, short of space or overloaded with transfer threads, which is why a busy cluster can produce pipelines that look odd. See rack awareness and blocks and replication for the policy in more detail.

What durable means: hflush, hsync and close

CallWhat it guaranteesTypical cost
write()bytes are in the client buffer; nothing morememory copy
hflush()data has reached all pipeline DataNodes and is visible to new readersone pipeline round trip
hsync()as hflush, and DataNodes have forced the data to diskround trip plus an fsync on each node
close()all data acknowledged and the file completed at the NameNodeflush, final packet, NameNode RPC

The most common misunderstanding is treating an ack, or hflush, as durable. An ack means the DataNodes have the data in memory or the operating system's page cache. If a whole rack loses power before the cache is written, data acknowledged by hflush can disappear. Write-ahead logs such as HBase's need hsync for commits that must survive power loss, and they pay for it in latency. Bulk writers that can regenerate their output usually rely on close alone.

Pipeline recovery

If a DataNode in the chain fails, or stops acknowledging within the socket timeout, the client stops streaming and rebuilds. It removes the bad node, asks the NameNode for a new generation stamp for the block, and sets up a new chain with the surviving nodes, which bring their partial replicas into line at the last acknowledged offset. The packets left in the ack queue are resent, and streaming resumes. The new generation stamp is the fence: if the failed DataNode comes back and reports its replica, the stamp is stale and the NameNode discards that copy instead of mixing it with the real one.

Whether the client also adds a replacement DataNode is controlled by dfs.client.block.write.replace-datanode-on-failure.enable (default true) and its .policy (default DEFAULT). Under the default policy, with replication of three or more, a replacement is added when the surviving count falls to half the replication factor or less, or when the stream was flushed or appended. The newcomer has to be brought up to date by copying the partial replica from a survivor before streaming resumes. On small clusters there may be no spare node, and the write fails with an error about failing to replace a bad DataNode; .best-effort (default false) lets it continue with fewer replicas instead, accepting reduced durability until the NameNode re-replicates the finished block.

Leases and lease recovery

The lease is what makes a single, ordered pipeline safe: only one client may write a file at a time, so there is no need to merge concurrent writers. The client renews its lease in the background. If the client dies, the lease expires: after the soft limit another client may claim the file, and after the hard limit the NameNode recovers it itself. The hard limit was hard-coded at one hour for many years; HDFS-14758 made it shorter and configurable in Hadoop 3.3.0 and several backport releases, so check your version's hdfs-default.xml rather than trusting an old number.

Lease recovery asks the DataNodes holding the last block to agree on its length, settles on the longest length all of them have, bumps the generation stamp, and closes the file. Until then, a file left open by a crashed writer reports a shorter length to readers and can block a job that wants to read or append. Operators can force it:

# find files still open for write under a path
hdfs fsck /logs/app -openforwrite

# ask the NameNode to recover the lease and close the file
hdfs debug recoverLease -path /logs/app/2026-10-05.log -retries 5

Worked example: a DataNode dies mid-block

A task on DataNode dn07 (rack A) writes a 300 MB file with three replicas. The file needs three blocks: 128 MB, 128 MB and 44 MB. For the first block the NameNode returns dn07, dn22 (rack B) and dn25 (rack B). The client sends roughly 2,000 packets of about 64 KB each. The first hop is a local socket, so the rack uplink carries each packet once from dn07 to dn22, and dn25 receives it over rack B's own switch.

Halfway through the second block, dn25 has a disk failure and stops acknowledging. The client times out, finds about 40 packets in its ack queue, drops dn25 and gets a new generation stamp. Two replicas survive out of three, which is more than half the replication factor, so under the default policy the client continues with two DataNodes rather than finding a third. It resends the 40 packets and finishes the block. After close, the NameNode sees the second block under-replicated and schedules a third copy in the background. The application saw a pause of a few seconds and no error. Had the job called hflush before the failure, the policy would have added a replacement node during the write, because readers may already depend on that data.

Failure modes

  • Slow DataNode: one node with a failing disk or saturated network throttles every pipeline that includes it. DataNode logs warn about slow writes to mirror or disk; cluster-wide slow-peer and slow-disk reporting helps find the culprit.
  • Transfer thread exhaustion: too many concurrent writers or many small files exhaust dfs.datanode.max.transfer.threads, and pipeline setup fails with errors that look like network problems.
  • Replacement failure on small clusters: with three or four DataNodes, a single failure can leave no candidate, so writes fail; tune the policy deliberately rather than by accident.
  • Files left open: crashed writers leave files under construction; readers see short lengths until lease recovery runs.
  • Durability assumed from hflush: data lost on correlated power failure because nobody called hsync.
  • No rack awareness: placement silently puts replicas in one rack.
  • Erasure-coded directories: files there do not use this replication pipeline at all; a striped stream writes data and parity cells to many DataNodes in parallel, with different failure behaviour. See HDFS erasure coding.

Operating and tuning

Monitor DataNode write throughput and latency per node, transfer thread usage, pipeline recovery counts in client and DataNode logs, under-replicated blocks, and files open for write. Rising recovery counts concentrated on one node point at hardware. Keep client and DataNode timeouts consistent across the cluster so failures are detected quickly but long fsyncs do not trigger false recoveries. Raise transfer threads on nodes that serve heavy concurrent writes, and fix the small-files pattern rather than tuning around it. The DataNode internals behind these signals are covered in the DataNode and HDFS tuning.

Trade-offs

ChoiceGainCost
Chain replicationeven network load, high aggregate throughputlatency of the full chain; one slow node stalls it
Ack before fsyncfast writesacked data can be lost on correlated power loss
hsync on commitsreal durabilityfsync latency on every DataNode
Single writer via leasesimple, ordered protocolno concurrent writers; recovery after crashes
Continue with fewer replicaswrite survives on small clusterstemporarily reduced durability
Erasure coding insteadabout half the storage of three replicasmore CPU and network, different write path

What to do next

  1. Confirm rack awareness: check that hdfs dfsadmin -printTopology shows real racks.
  2. List which writers need durability on commit and make sure they call hsync, not just hflush.
  3. Review the replace-datanode-on-failure policy for each cluster size, especially small ones.
  4. Alert on files open for write longer than your longest job, and script lease recovery.
  5. Track pipeline recoveries per DataNode to catch failing disks early.
  6. Check transfer thread usage on busy nodes before raising it.
  7. Look up your release's lease hard limit instead of assuming one hour.
Key takeaway: The HDFS write pipeline trades per-packet latency for even network load: the client streams packets down a rack-aware chain, acks come back up, and the ack queue plus generation stamps let a write survive a failed DataNode. Know that an ack is not an fsync, configure rack awareness, and watch slow nodes and open files.