Where does HDFS put the three replicas of a block under its default rack-aware policy?
answer
- one copy near, two copies far
- never fewer than two failure domains
- the third copy does not need a new rack
- one crossing of the expensive link
- the cluster has to be told the layout
basics
~20 sWith the default policy, the first replica goes on the writing client's own node (or a random DataNode if the client is off-cluster), the second on a node in a different rack, and the third on another node in that same second rack.
solid answer
~50 sThe default placement policy balances three goals: write cost, read locality and fault tolerance. Replica one lands on the local node when the writer runs on a DataNode, which makes the first hop of the write pipeline free; otherwise the NameNode picks a random DataNode that is not overloaded or nearly full. Replica two goes to a node in a **different rack**, so losing an entire rack — a switch, a power feed — never loses the block. Replica three goes to a *different node in that same remote rack*, because putting it in a third rack would buy little extra safety while costing another cross-rack copy. The net effect: exactly one block's worth of traffic crosses the rack boundary during a write, and reads can usually be satisfied locally or within a rack. HDFS does not discover racks by itself — you must supply a topology script via `net.topology.script.file.name`, and without it every node is assumed to be in one default rack.
code
xml · 4 lines<property>
<name>net.topology.script.file.name</name>
<value>/etc/hadoop/conf/rack-topology.sh</value>
</property>go deeper
Recall the shape of the rule — one copy local, the others in a second rack — and that the point is to survive losing a whole rack rather than just a single machine.
Explain why the third replica shares the second's rack: only one copy crosses the expensive inter-rack link while two failure domains are still covered. Mention that rack topology comes from a configured script, not discovery.
Show the operational consequences — an unconfigured topology script silently collapsing all replicas into one rack, re-replication honouring the same constraints after node loss, and the balancer being unable to break rack diversity.
Own the mapping from physical or cloud failure domains onto HDFS racks, and the cross-zone bandwidth and cost implications of that mapping. Decide when placement policy alone is enough versus when erasure coding or a different storage tier is the better answer.
## The problem being solved Three replicas can be placed in many ways, and the choices trade off against each other: - **All three in one rack** — cheap to write (no cross-rack traffic), fast to read, and catastrophic when the rack's top-of-rack switch dies: the block is unreachable. - **All three in different racks** — maximally safe, but the write pushes two full copies across the aggregation layer, which is the scarcest bandwidth in a classic Hadoop network. HDFS's default policy, implemented by `BlockPlacementPolicyDefault`, sits deliberately in between. ## The default rule for replication factor 3 1. **First replica**: the node the writer is running on, if that node is a DataNode with room. If the client is outside the cluster, the NameNode chooses a random DataNode, avoiding nodes that are too full or too busy. 2. **Second replica**: a node on a **different rack** from the first. 3. **Third replica**: a **different node on the same rack as the second**. Beyond three, further replicas are scattered across the cluster with a cap on how many may share a rack, so a very high replication factor does not pile up in one place. Read the pattern as: *never fewer than two racks, never more cross-rack copying than necessary.* ## Why this specific shape **Write cost.** The write pipeline is a chain — client to DN1 to DN2 to DN3. With this layout, exactly one link (DN1 to DN2) crosses a rack boundary; the DN2-to-DN3 hop stays inside the second rack, where bandwidth is cheap. Only one block's worth of data traverses the expensive inter-rack path. **Rack-failure tolerance.** Because at least two racks always hold a copy, the loss of any single rack — switch failure, power distribution unit, a maintenance window — leaves the block readable. This is the whole point of rack awareness, and it is why a cluster with no topology script configured is quietly less safe than its operators think. **Read locality.** When a client opens a block, the NameNode returns the replica locations sorted by network distance: same node first, then same rack, then remote rack. Schedulers use the same information to place tasks on a machine that already holds the input block, so a scan reads from local disk and never touches the network. **Balance.** The policy also avoids nodes above a utilisation threshold and spreads replicas so no single node becomes a hot spot, though genuine imbalance across a cluster is corrected by the separate **balancer** tool rather than by placement alone. ## Rack awareness is configuration, not discovery HDFS has no way to learn your physical layout. You provide a **topology script** — an executable named by `net.topology.script.file.name` — that takes hostnames or IP addresses on the command line and prints a rack path such as `/dc1/rack12` for each. The NameNode calls it when a DataNode registers and caches the result. An alternative is a static table mapping via `net.topology.table.file.name`. If you configure nothing, every node resolves to `/default-rack`. Everything still works, but the second replica is no longer guaranteed to be in a different failure domain: the policy believes it already is. That is a common and dangerous misconfiguration — the cluster looks healthy and passes every test right up until a rack goes dark. In cloud deployments the same mechanism is used to map availability zones or placement groups onto "racks", so the policy tolerates a zone failure instead of a physical rack failure. ## What happens after placement The policy runs again whenever a block becomes **under-replicated** — a DataNode declared dead, a disk failure, a corrupt replica detected by the block scanner. The NameNode picks a source and a target and instructs a DataNode-to-DataNode copy, honouring the same rack constraints and throttled so recovery does not starve running jobs. Conversely, if a block ends up **over-replicated** (a dead node returns), the NameNode deletes surplus replicas, preferring to remove ones that would not reduce rack diversity. The **balancer** is a separate concern: it moves replicas between DataNodes to even out disk utilisation, and it too must preserve the rack invariants, so it can never fix imbalance by collapsing every copy of a block into one rack. ## Changing the policy The placement policy is pluggable via `dfs.block.replicator.classname`. Alternatives ship for specific topologies — for example, a policy that spreads replicas across node groups for virtualised clusters where several "nodes" share physical hardware. Most clusters never change it.
- Why is the third replica placed in the same rack as the second rather than in a third rack?Because the extra safety is marginal and the cost is not. Two racks already survive any single rack failure; a third rack would add another cross-rack copy on every write, consuming the scarcest bandwidth in the network. Keeping replicas two and three together means the second pipeline hop stays inside a rack, so exactly one block's worth of data crosses the aggregation layer.
- What is the risk of running a multi-rack HDFS cluster with no topology script configured?Every DataNode resolves to /default-rack, so the placement policy believes it has already achieved rack diversity when it has not. All three replicas can land behind a single top-of-rack switch. The cluster reports full replication and looks perfectly healthy until that rack loses power or its switch fails, at which point blocks become unreadable.
- How does the NameNode decide which replica an HDFS client reads from?It returns the block's locations sorted by network distance from the requesting client: a replica on the same node first, then one in the same rack, then a remote rack. The client tries them in that order and falls through on failure or checksum mismatch. This is the same distance calculation that lets a scheduler place a task where its input block already lives.
Keep one copy of an important document in your own desk, one in a filing cabinet in another building, and a second copy in that same other building — so a fire in either building never destroys the last copy, but you only carry it across town once.
saying these in an interview costs you the question
- Saying all three replicas are always placed in three different racks
- Assuming HDFS detects rack topology automatically
- Claiming the first replica is always remote from the writer
- Thinking the balancer can override rack diversity constraints
- Confusing replica placement with the erasure-coding stripe layout