Coping With Node Failures
On this page
2.5 Coping With Node Failures
Recall first
A DataNode stops sending heartbeats while a map task on it is reading a block. What detects the failure, what restores the data, and what happens to the task?
First principles
Failure is normal
A large cluster contains many disks, machines, processes, and network links. The design assumes that components fail and uses detection, redundancy, and re-execution rather than assuming perfect hardware. HDFS and MapReduce provide different but complementary recovery mechanisms.
HDFS recovery
- Detect: each DataNode sends periodic heartbeats to the NameNode. Missing heartbeats cause the NameNode to mark a node unavailable; block reports describe which blocks a live node stores.
- Avoid: clients and tasks stop using a node considered dead.
- Recover replicas: if a block now has fewer live replicas than its replication factor, the NameNode schedules re-replication from an available replica to another DataNode.
- Detect corruption: HDFS checksums block data. If a read fails validation, the client can fetch the block from another replica and the bad replica can be handled by recovery procedures.
These are documented in HDFS robustness. Replication improves availability but does not make a write free: it consumes storage and network bandwidth.
MapReduce recovery
The MapReduce framework monitors tasks. If a map or reduce attempt fails, it launches another attempt, often on a different node. A failed map’s intermediate output may disappear with the worker, so reducers fetch output again from the replacement map attempt. A failed reducer can refetch map partitions and recompute its output.
Speculative execution addresses a slow but still alive “straggler”: the framework may run another attempt on a different node and accept one successful result. This helps when slowness is caused by a sick disk or machine, but wastes resources if the task is genuinely slow because of skew or expensive input.
Task code should therefore be deterministic and safe to retry. Do not use an external side effect such as “charge a card” inside a map function unless the side effect is idempotent or separately coordinated. Hadoop’s output committer helps make task output visible only after a successful attempt, but application-level side effects remain the programmer’s responsibility (MapReduce Tutorial).
Worked failure trace
A block B has replicas on A, B, and C. A map on A emits intermediate data, then A crashes.
- HDFS still serves
Bfrom B or C. - The NameNode detects A’s missing heartbeat and later schedules a replacement replica.
- The failed map attempt is marked failed; a new attempt reads B or C, preferably locally/rack-locally, and regenerates its intermediate partition.
- Reducers discard the failed attempt’s unavailable output and fetch the successful attempt’s output.
The job may slow down, but it can finish without manual repair.
Exercise — revealed answer
Exercise: Does replication alone recover a failed reduce task?
Answer: No. HDFS replication protects input/output blocks, while MapReduce task re-execution recomputes the failed reduce output. The two mechanisms protect different layers.
Exam lens
Organize the answer as failure type → detector → response: DataNode failure → missing heartbeat → stop scheduling and re-replicate; task failure → task monitoring → retry; straggler → speculative duplicate → keep one successful attempt. Mention checksums and idempotence for extra accuracy.
Rapid revision checklist
- Explain heartbeat-based detection.
- Explain block re-replication.
- Explain task retry and lost map output.
- Define speculative execution and its cost.
- Explain why deterministic/idempotent tasks matter.
Key takeaways
- HDFS repairs stored data; MapReduce redoes computation.
- Replicas provide alternate bytes, while retries provide alternate attempts.
- Speculation improves tail latency but can waste resources and does not fix data skew automatically.
Sources
- Apache HDFS Design — robustness.
- Apache Hadoop MapReduce Tutorial — task execution and attempts.
- Apache Job API.
- Syllabus-aligned supplement: the layered failure trace and idempotence warning connect HDFS and MapReduce official mechanisms to the syllabus’s “coping with node failures” topic.