How does an HDFS client write one block through a DataNode replication pipeline?
answer
- the client sends the bytes only once
- a chain, not three parallel sends
- packets go one way, confirmations the other
- one queue holds what is not yet confirmed
- the file can close before all copies exist
basics
~20 sThe client asks the NameNode for a block and a list of DataNodes, then streams packets to the first DataNode, which forwards to the second, which forwards to the third. Acknowledgements travel back up the chain before packets are released.
solid answer
~50 sOn `create()`, the NameNode records the new file; then for each block the client calls for an allocation and gets back an ordered list of DataNodes. The client's output stream buffers bytes, cuts them into packets, and pushes them into a **data queue**. Packets go to the first DataNode only — that node persists the packet and immediately forwards it to the second, which forwards to the third, so the replication cost is paid as a **chain, not a fan-out** from the client. Each sent packet sits in an **ack queue** until an acknowledgement travels back up the pipeline from all DataNodes, at which point it is discarded. When the block is full the pipeline is torn down and a new one is set up for the next block. The file close succeeds once at least `dfs.namenode.replication.min` replicas — one by default — are durably written; the NameNode then arranges any remaining replicas asynchronously.
code
text · 6 linesclient --packet--> DN1 --packet--> DN2 --packet--> DN3
<---ack--- <---ack--- <---ack---
data queue : [p7][p8][p9] <- not yet sent
ack queue : [p4][p5][p6] <- sent, awaiting full ack
^ replayed here if a DataNode drops outgo deeper
Know the shape: the NameNode hands out a block and a list of DataNodes, and the data flows client to DataNode to DataNode to DataNode rather than the client sending three copies itself.
Walk the mechanics end to end — lease, block allocation, packets, data queue and ack queue, pipeline teardown per block — and explain why the chain design keeps the client's uplink from being the bottleneck.
Demonstrate the failure path: pipeline recovery, replaying the ack queue, the new generation stamp that invalidates the stale replica, and completing on minimum replication with the NameNode healing afterwards. Tie hflush and hsync to real durability incidents.
Frame it as a durability policy decision. Minimum replication of one trades a window of single-copy data for write availability, and the choice interacts with rack layout, node failure rates and how long re-replication takes on your hardware.
## Why a pipeline and not a fan-out HDFS writes every block to several DataNodes. The naive design would have the client send the block three times, once per replica — which triples the client's outbound bandwidth and makes the client the bottleneck. Instead HDFS chains the DataNodes: the client sends each packet **once**, to the head of a pipeline, and each DataNode forwards to the next as it receives. Every link carries roughly one block's worth of traffic, the client's uplink carries exactly one, and the cost is spread across the cluster fabric. ## Step by step **1. Create the file.** The client calls `create()` on the `DistributedFileSystem`. The NameNode checks that the path does not already exist and that the caller has permission, then records a new, empty file in the namespace and grants the client a **lease** — HDFS allows exactly one writer per file, and the lease must be renewed or the write is reclaimed. **2. Allocate a block.** The client writes into a `DFSOutputStream`, which buffers bytes locally. When a new block is needed the client asks the NameNode to allocate one; the NameNode returns a fresh block ID and an ordered list of target DataNodes chosen by the replica placement policy. **3. Build the pipeline.** The client opens a connection to the first DataNode, which connects to the second, which connects to the third. Setup acknowledgements flow back, and the pipeline is live. **4. Stream packets.** The stream splits data into packets (64 KB by default), each carrying checksums for its chunks. Packets are appended to the **data queue** and sent to the head DataNode. Each DataNode writes the packet to its local block file and forwards it downstream in the same motion, so the three writes overlap in time rather than happening one after another. **5. Acknowledge.** Every sent packet moves to the **ack queue**. The tail DataNode acknowledges to the middle, which acknowledges to the head, which acknowledges to the client. Only when a packet is acknowledged by every node in the pipeline is it removed from the ack queue — which is precisely what makes recovery possible. **6. Close.** After the last packet, the client flushes the queues and calls `complete()` on the NameNode. The file is complete once the NameNode sees at least `dfs.namenode.replication.min` replicas — one by default — for every block. Bringing the block up to the full replication factor happens afterwards, driven by the NameNode. When a block fills, the pipeline for it is closed and the whole allocate-and-build sequence repeats for the next block, generally with a different set of DataNodes. ## What happens when a DataNode in the pipeline dies This is the part interviewers push on, because it explains why the ack queue exists: 1. The pipeline is closed. 2. Packets still sitting in the ack queue — sent but not fully acknowledged — are pushed back to the **front** of the data queue, so no data is lost. 3. The surviving DataNodes' partial replica gets a new **generation stamp**, which lets the NameNode recognise the stale replica on the failed node as obsolete and schedule it for deletion when that node returns. 4. The failed node is removed from the pipeline and writing continues to the remaining nodes. Depending on policy the client may add a replacement node to restore the pipeline's width. 5. The NameNode later notices the block is under-replicated and instructs a DataNode to copy it to a fresh machine. The write only fails outright if the pipeline cannot keep the configured minimum number of replicas. ## Visibility and durability semantics HDFS is **write-once, read-many with append**: a file has a single writer, there is no random overwrite in place, and only appending to the end is supported. While a block is being written, other readers cannot generally see the newest bytes. `hflush()` pushes buffered data to all DataNodes in the pipeline and makes it visible to new readers, but does not guarantee it has reached disk; `hsync()` additionally forces the DataNodes to sync to their storage. Closing the file does both and finalises the block. That distinction is what bites people who use HDFS as a streaming sink: an unclosed file can appear to have zero length to a reader, and a client that dies without closing leaves a file whose last block is recovered through **lease recovery**, truncating to the last consistently acknowledged length. ## Where the client sits If the writing client is itself running on a DataNode — the usual case for a job task — the placement policy puts the first replica on that same machine, so the head of the pipeline is a local write with no network hop at all. A client outside the cluster gets a randomly chosen (and not overloaded) DataNode as the head instead, and pays a network hop for every byte.
- Why are unacknowledged packets kept in a separate ack queue rather than discarded once sent?Because a DataNode can fail mid-pipeline. If that happens, every packet that was sent but not acknowledged by all nodes is moved back to the front of the data queue and replayed down the rebuilt pipeline. Without that retained copy the client would have no way to recover those bytes, and the write would have to fail rather than heal.
- If dfs.replication is 3, can a write succeed with only one replica on disk?Yes. Completion is gated by `dfs.namenode.replication.min`, which defaults to 1, not by the target replication factor. The file closes once the minimum is met, and the NameNode then schedules copies to reach three. That is a deliberate availability-over-durability trade for the write path, and it is why a cluster losing nodes can still accept writes.
- What is a lease in the HDFS write path, and what happens if the writing client crashes?A lease is the NameNode's grant of exclusive write access to one file, renewed by heartbeat from the client. If the client dies, the lease expires and the NameNode runs lease recovery: it reconciles the last block's length across the DataNodes that hold it, truncates to a consistent acknowledged length, finalises the block and closes the file so others can read or overwrite it.
Think of a bucket brigade rather than three separate deliveries: the writer hands each bucket to one person, it travels down the line while the next bucket is already being handed over, and a shout back up the line confirms it reached the end.
saying these in an interview costs you the question
- Saying the client sends the block separately to all three DataNodes
- Claiming the write fails unless all three replicas complete
- Believing HDFS supports random in-place updates to a file
- Thinking data is visible to readers before hflush or close
- Confusing an HDFS write pipeline with a NameNode checkpoint