HDFS stores each block three times by default, so an HBase cluster holding 100 TB of table data buys 300 TB of disk. Erasure coding cuts that to about 150 TB with the common Reed-Solomon 6+3 policy, while still surviving the loss of any three pieces of a stripe. For large, cold, rarely rewritten tables that is the single biggest lever on the cost of a cluster.
HBase is not a simple HDFS client, though. It keeps a write-ahead log that needs durability calls erasure-coded files do not support. It rewrites its files constantly through flushes and compactions. It moves files between directories for snapshots, bulk loads and archiving, and it reads small blocks at random. The mechanics of erasure coding itself, meaning striping, parity and reconstruction, are covered in HDFS erasure coding. This article covers what is specific to HBase: which directories may be erasure coded, how to turn it on safely for a table, how existing data migrates, and what changes on the read and write paths.
What HBase keeps under its root directory
Everything an HBase cluster stores lives under one HDFS root, usually /hbase. Each subtree has its own write pattern, and that pattern decides whether it can be erasure coded.
| Path | What it holds | Write pattern | Erasure coding |
|---|---|---|---|
/hbase/WALs | live write-ahead logs, one set per RegionServer | appended continuously, hflush or hsync on every sync | no |
/hbase/oldWALs | logs kept for replication or cleanup | moved from WALs | no (they were written replicated) |
/hbase/MasterData | the master's local store and its own WAL | small, synced writes | no |
/hbase/data/ns/table | regions, column families, HFiles | written once by flush or compaction, then only read | yes, per table |
/hbase/archive | HFiles replaced by compaction but still held by snapshots or replication | moved, not rewritten | keeps whatever layout the file already had |
/hbase/.tmp | files being built by compaction | written then renamed into place | follow the target table (see below) |
The rule that falls out of the table: erasure coding belongs on table data directories, not on the root. A policy on /hbase is inherited by every new file below it, and log directories are exactly where striping does harm.
The WAL rule: no hflush, no erasure coding
A put is acknowledged only after its WAL entry is durable enough to survive a RegionServer crash. HBase gets that guarantee from hflush, which pushes the data to every DataNode in the pipeline, or hsync, which also forces it to disk. The write pipeline behind this is described in HDFS blocks and replication, and HBase's durability levels in HBase WAL durability.
Erasure-coded output streams do not implement hflush or hsync. A striped file is written as cells spread across nine DataNodes, and parity can only be computed once a stripe's data is known, so there is no cheap point at which a half-written stripe is durable. A striped WAL would lose acknowledged writes. HBASE-19369 (fixed in 2.0.2, 2.1.1 and 2.2.0) addressed this by creating WAL files through HDFS's file-builder API so they can be kept replicated even when their directory carries an erasure-coding policy. Treat that as a safety net, not a design: check which WAL writer your version uses and whether the fix covers it, and keep policies off the log directories anyway. If a RegionServer ever complains that its WAL stream lacks hflush or hsync, fix the directory policy; do not disable the check.
The storage arithmetic, and why small files spoil it
With RS-6-3-1024k, data is cut into 1 MB cells. Six data cells and three parity cells form a stripe, so a large file costs 1.5 times its size, against 3 times for triple replication. RS-3-2-1024k costs about 1.67 times and needs fewer machines. The policy names encode exactly this: codec, data units, parity units, cell size.
Parity cells are full size even when the stripe is not. Take a 2 MB HFile under RS-6-3-1024k. It fills two data cells, and the stripe still needs three 1 MB parity cells, so the file stores 5 MB: 2.5 times its size, not 1.5. A file under 1 MB, one cell plus three parity cells, is close to four times. Erasure coding only pays off when HFiles are large, and in HBase that means tables whose regions have been major compacted into big files. Tables with frequent small flushes and many small store files see little saving until compaction consolidates them.
Topology matters too. RS-6-3 needs at least nine DataNodes to place one cell of each stripe per node. To survive the loss of a whole rack, the HDFS guidance is at least (data + parity) / parity racks, which for 6+3 is three racks, with the nodes spread evenly across them.
Turning it on for a table
Start in HDFS. Policies must be enabled before anything can use them, and you should check that the cluster can host the one you pick.
hdfs ec -listPolicies # shows each policy and whether it is enabled
hdfs ec -enablePolicy -policy RS-6-3-1024k
hdfs ec -verifyClusterSetup -policy RS-6-3-1024k # enough DataNodes and racks?From HBase 2.6.0 and 3.0.0 (HBASE-28216), the policy is a table setting, ERASURE_CODING_POLICY. HBase checks that the policy exists, is enabled, and that the cluster can actually write with it, by writing a test file, before it applies the policy to the table's data directory. The policy then lives in the table's schema instead of a hand-set directory attribute. The Java API is shown below. The shell accepts table attributes the same way as other descriptor settings, but check the exact syntax against your version's shell help before scripting it.
// HBase 2.6+ Java API: set the policy, then major compact to rewrite existing files.
try (Connection conn = ConnectionFactory.createConnection(conf);
Admin admin = conn.getAdmin()) {
TableName tn = TableName.valueOf("events_2024");
TableDescriptor current = admin.getDescriptor(tn);
TableDescriptor updated = TableDescriptorBuilder.newBuilder(current)
.setErasureCodingPolicy("RS-6-3-1024k")
.build();
admin.modifyTable(updated);
admin.majorCompact(tn);
}On older HBase releases, the only option is to set the policy directly on the table directory with HDFS tools. That works, because new files simply inherit it, but HBase does not know about it, so a table recreated by a script comes back replicated.
hdfs ec -setPolicy -path /hbase/data/default/events_2024 -policy RS-6-3-1024k
hdfs ec -getPolicy -path /hbase/data/default/events_2024
Why nothing changes until you compact
An erasure-coding policy on a directory applies to files created in it afterwards. Existing HFiles keep their replicated layout. HDFS does not convert files in place, and a rename keeps the layout a file was written with. HBase's own write paths happen to give you the conversion for free: every flush and every compaction writes brand-new HFiles, so a major compaction of the table rewrites all of its data under the new policy. That is why the HBase release note says the policy only takes effect after a major compaction. How to schedule and throttle that rewrite is covered in HBase compaction.
Two paths get around the policy, and both catch people out:
- Bulk load. HFiles built by a MapReduce or Spark job are moved into the region directories. They keep the layout of the directory where they were written, usually a replicated staging area, until a later compaction rewrites them.
- Archive and snapshots. When compaction replaces a file that a snapshot still references, the old file moves to
/hbase/archiveand keeps its replicated layout. During migration, a table with snapshots can briefly use more disk: the old replicated files in the archive plus the new erasure-coded ones. Expire or export snapshots before migrating a large table, or budget for the overlap.
Check the result, not the setting. hdfs ec -getPolicy on the directory only tells you what new files will get. hdfs fsck on the table path reports block groups for erasure-coded files and ordinary replicated blocks for the rest, so it shows how much of the table has actually been converted.
What changes on the read and write paths
Locality disappears. With replication, a compacted region's blocks usually have a replica on the RegionServer's own DataNode, and reads are local. With striping, each HFile block is spread across nine DataNodes, so almost every read crosses the network. Region locality metrics stop meaning much for these tables, and the network becomes part of the read-latency budget.
Random reads touch one or two cells. A get reads a 64 KB HFile block with a positional read. With 1 MB cells that is usually inside one cell on one DataNode, sometimes across two, so a healthy read costs about the same number of round trips as before, only remote. A read from a lost cell is a degraded read: the client fetches six other cells of the stripe and decodes, which is several times the network traffic for that block. Read latency percentiles get worse exactly when DataNodes fail.
Caching matters more. Because every miss is remote and a miss during a failure is expensive, a larger block cache, usually an off-heap BucketCache, recovers much of the lost locality for hot index and data blocks.
Writes fan out. Flushes and compactions write to nine DataNodes in parallel and compute parity on the client, which is the RegionServer. That costs CPU on the RegionServers, so use the native ISA-L coder where your Hadoop build provides it, and throttle compactions so the migration does not starve serving.
Which tables to convert
| Table profile | Recommendation | Why |
|---|---|---|
| Large, cold, append-mostly history (events by month, audit logs) | convert with RS-6-3 or RS-3-2 | big compacted files, few reads, maximum saving |
| Hot serving table with strict p99 latency | keep replicated | locality and degraded-read latency matter more than disk |
| Small tables or many small store files | keep replicated | parity overhead on small files removes most of the saving |
| Tables with long-lived snapshots | convert after snapshot cleanup | the archive holds replicated copies during migration |
| Bulk-loaded tables | convert and compact after each load | loaded files arrive replicated |
A common and safe pattern is time-based: write current data to a replicated table, and roll finished periods into erasure-coded tables that are major compacted once and then mostly read.
Failure modes
- Policy on the root or WAL directory. Depending on version and WAL writer, WAL creation fails or relies on the builder fix to stay replicated, and every other new file under the root is striped. Remove the policy; files already written keep their layout, so roll the logs afterwards.
- Too few DataNodes or racks. Writes fail or lose rack tolerance. HBase 2.6's validation catches missing nodes at setting time, but a cluster that later shrinks can still fall below the policy's needs.
- Setting without compaction. Weeks later, disk use has not moved. Check with
fsckand run the major compaction. - Reconstruction storms. Losing a DataNode triggers background reconstruction, which reads six cells for every one it rebuilds and competes with serving reads. Tune the reconstruction throttles before the first failure, not during it.
- Replicated peers. The policy is per cluster. A replication peer or DR cluster needs its own decision and its own compaction.
What to do next
- Measure HFile size distribution per table with
hdfs dfs -duand a file count, and shortlist tables whose files are mostly well above one stripe (6 MB for RS-6-3-1024k). - Confirm DataNode and rack counts with
hdfs ec -verifyClusterSetupfor the policy you choose, and enable it. - Check that
/hbase,/hbase/WALsand/hbase/MasterDatahave no erasure-coding policy. - Clean up or export snapshots on the first candidate table, set
ERASURE_CODING_POLICY(HBase 2.6+) or the directory policy, and run a throttled major compaction. - Verify with
hdfs fsck, compare disk use before and after, and watch p99 read latency, network traffic and RegionServer CPU for a week. - Kill a DataNode in a test cluster and measure degraded-read latency and reconstruction time before relying on it in production.