When would you choose HDFS erasure coding over three-way replication in Hadoop 3?
answer
- parity instead of whole copies
- same safety, roughly half the disk
- nothing whole sits on one machine
- recovery costs arithmetic, not a copy
- best for data nobody reads
basics
~20 sErasure coding suits large, cold, rarely-read files: a Reed-Solomon scheme such as RS-6-3 tolerates three failures at about 1.5x storage overhead instead of replication's 3x. It costs read locality, CPU on recovery, and does not support append.
solid answer
~50 sThree-way replication burns 200% extra storage to survive two node losses. HDFS erasure coding, added in Hadoop 3.0, instead stripes a block group across DataNodes and stores parity cells alongside — with the system default policy `RS-6-3-1024k`, six data cells plus three parity cells give the same tolerance of three simultaneous failures at roughly 1.5x overhead. The catch is that nothing sits whole on one node, so a read gathers cells from six machines and **data locality disappears**; any missing cell must be reconstructed with Reed-Solomon arithmetic, which is CPU- and network-expensive and makes degraded reads slow. Erasure-coded files also do not support `append`, `hflush` or `hsync`, and small files are inefficient because a file shorter than a full stripe still pays for parity. So: apply it per directory to large, write-once, infrequently-scanned archive data, and leave hot working sets on replication.
code
bash · 5 lineshdfs ec -listPolicies
hdfs ec -setPolicy -path /warehouse/cold -policy RS-6-3-1024k
hdfs ec -getPolicy -path /warehouse/cold
# existing files are unaffected -- rewrite to convert
hadoop distcp /warehouse/hot/2024 /warehouse/cold/2024go deeper
Recall that Hadoop 3 added an alternative to copying every block three times: parity-based erasure coding that stores the same durability for far less disk. Knowing it exists is enough at this level.
Explain the arithmetic of a policy like RS-6-3 — six data cells, three parity, roughly 1.5x overhead, three tolerable failures — and that HDFS stripes the group across nodes rather than keeping blocks contiguous.
Show the operational tradeoffs you would weigh: lost locality on reads, CPU and network cost of reconstruction, degraded-read latency during failures, unsupported append and hflush, and the minimum DataNode count a wide policy needs.
Own it as a storage-tiering policy. Decide which datasets age into an erasure-coded tier and when, model the disk saving against the extra CPU, network and recovery risk, and compare the whole thing against simply moving cold data to object storage.
## The storage bill replication hands you HDFS's default durability mechanism is brute-force copying: replication factor 3 means every byte is stored three times, a 200% overhead. It is simple, it gives you read locality for free (any of three machines can serve the block), and recovery is a plain copy. For a petabyte of cold logs nobody has queried in two years, it is also an enormous waste of disk. ## What erasure coding does instead Hadoop 3.0 introduced **HDFS Erasure Coding**. Rather than copying a block, it takes a group of data cells, computes parity cells over them with Reed–Solomon arithmetic, and spreads the whole group across DataNodes. Any subset of cells equal in size to the original data is enough to reconstruct everything. Policies are named for their parameters. `RS-6-3-1024k` means six data cells and three parity cells, with a 1024 KB cell size; it is the system default policy in Hadoop 3. Nine cells store what six cells' worth of data needs, giving roughly 1.5x overhead, and any three of the nine can be lost without data loss — the same failure tolerance as three-way replication at half the storage cost. Other shipped policies include `RS-3-2-1024k`, `RS-10-4-1024k` and `XOR-2-1-1024k`, trading width against overhead and failure tolerance. Policies other than the system default must be explicitly enabled with `hdfs ec -enablePolicy` before they can be applied. HDFS's implementation is **striped**, not contiguous: the cells of a block group are interleaved across many DataNodes rather than each node holding a whole contiguous block. ## What you give up **Data locality.** With replication, a task can be scheduled on a machine that already holds the block and read it from local disk. With striping, no machine holds a whole block — a read pulls cells from at least six DataNodes over the network. For scan-heavy hot data this is a real throughput and network-cost regression, and it is the single biggest reason not to erasure-code a working set. **Expensive reconstruction.** Losing a replica means copying one file between two nodes. Losing an erasure-coded cell means reading the surviving cells of the group and running Reed–Solomon decode to rebuild it — significant CPU (helped by hardware-accelerated ISA-L, when available) and much more network traffic per byte recovered. A cluster with steady disk failures pays this continuously, and a **degraded read** — one that hits a missing cell and reconstructs on the fly — is far slower than a healthy one. **Restricted write semantics.** Erasure-coded files are write-once in a stricter sense than replicated ones. `append` and `truncate` are not supported, and neither are `hflush` and `hsync`, so an EC directory cannot be a streaming sink that needs durable incremental visibility. `setReplication` is meaningless on an EC file. **Small files are wasteful.** A file shorter than one full stripe still needs a full set of parity cells, so its effective overhead approaches or exceeds replication's. Erasure coding pays off only on files comfortably larger than the stripe width times the cell size — many megabytes at minimum. It also does nothing for the NameNode: an EC block group still costs namespace objects, and in fact tracks more internal blocks than a replicated one. **More DataNodes required.** A `RS-6-3` policy needs at least nine DataNodes to place a group properly, and to survive a rack failure it wants those spread across enough racks. Small clusters cannot use wide policies safely. ## How you apply it Erasure coding is set **per directory**. You choose a policy on a path, and files created under it inherit that policy; existing files are unaffected, so converting old data means rewriting it — typically with `distcp` into an EC-enabled directory. A single cluster routinely runs both schemes: hot partitions on replication, aged partitions moved into an erasure-coded archive tree by a lifecycle job. ``` hdfs ec -listPolicies hdfs ec -enablePolicy -policy RS-3-2-1024k hdfs ec -setPolicy -path /warehouse/cold -policy RS-6-3-1024k hdfs ec -getPolicy -path /warehouse/cold ``` ## The decision rule Use erasure coding when the data is **large, immutable, and rarely read**: compliance archives, raw event history past its query window, backups of curated tables. Keep replication for anything read repeatedly by jobs that benefit from locality, anything written incrementally, anything small, and anything where a degraded read during a node failure would breach a latency expectation. A useful framing for an interview: replication buys read performance and operational simplicity with disk; erasure coding buys disk with CPU, network and flexibility. On a cluster where storage cost dominates and the data is cold, that is an excellent trade. On a hot working set it is a bad one, and applying it cluster-wide is a classic over-correction.
- Why does erasure coding hurt data locality for processing jobs?Because HDFS erasure coding is striped: a block group's cells are spread across at least as many DataNodes as the policy is wide, so no single machine holds a contiguous block. A task cannot be scheduled next to its input the way it can with replication — every read pulls cells over the network from several nodes, and any missing cell triggers an on-the-fly reconstruction that is slower still.
- Is erasure coding a solution to the HDFS small-files problem?No, it makes small files worse on both axes. A file shorter than one stripe still needs a full set of parity cells, so its effective storage overhead can approach or exceed replication's. And the NameNode still tracks namespace objects per file — an EC block group actually involves more internal blocks than a replicated one, so master memory pressure does not improve.
- How do you convert an existing replicated HDFS directory to erasure coding?You rewrite the data. Setting a policy on a path affects only files created afterwards, so the usual pattern is to create an EC-enabled target directory, `distcp` the data into it, verify, then delete the original. Because erasure-coded files do not support append or truncate, make sure nothing downstream was relying on incremental writes to those paths before you switch.
Replication is photocopying a document three times. Erasure coding is tearing it into six pieces and adding three cleverly computed patches, so any six of the nine scraps rebuild the original — cheaper to store, but you have to fetch pieces from six drawers and do arithmetic to read it.
saying these in an interview costs you the question
- Claiming erasure coding also fixes the small-files problem
- Assuming erasure-coded files still support append or hflush
- Applying erasure coding cluster-wide including hot working sets
- Thinking EC preserves data locality the way replication does
- Saying erasure coding tolerates fewer failures than three-way replication