skip to content

Where does HDFS put the three replicas of a block under its default rack-aware policy?

level: middleimportance: should knowfreq 38%

answer

  1. one copy near, two copies far
  2. never fewer than two failure domains
  3. the third copy does not need a new rack
  4. one crossing of the expensive link
  5. the cluster has to be told the layout

basics

~20 s

With 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 s

The 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
xml
<property>
  <name>net.topology.script.file.name</name>
  <value>/etc/hadoop/conf/rack-topology.sh</value>
</property>

go deeper

for a junior

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.

for a middle

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.

for a senior

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.

for a principal

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

context