Read amplification is the ratio between the work a database does to answer a read and the data it actually returns. In a B-tree the ratio is small and stable: a handful of page reads per lookup. In an LSM-tree store such as HBase it is variable by design. HBase makes writes cheap by appending them to a memstore and flushing immutable sorted files, and pays for it on reads, which must consult every file that might contain the key and then skip over older versions and delete markers until compaction cleans them up.
This page treats read amplification as something you measure and bound rather than a vague performance complaint. It breaks it into three kinds, builds a small cost model for a single-row Get, shows how the default compaction settings cap the worst case, explains why deletes, versions and filters inflate scans, and ends with a procedure and checklist for diagnosing and fixing it. The mechanisms themselves have their own pages: bloom filters, the block cache and compaction. Here we connect them to one number.
Three kinds of read amplification
It helps to split the ratio into factors that have different causes and different fixes.
| kind | definition | driven by | main levers |
|---|---|---|---|
| file | store files consulted per read | flush rate vs compaction rate | compaction, blooms, time-range pruning |
| block | HFile blocks read per row returned | block size, cache hit rate, file count | block cache, block size, encoding |
| cell | cells scanned per cell returned | versions, tombstones, expired cells, filters | VERSIONS, TTL, major compaction, row-key design |
They multiply. A Get that consults 10 files, reads one uncached 64 KB block per file and scans 50 obsolete versions to find one live cell is slow for three independent reasons, and fixing only one of them helps less than you would expect.
Where it comes from: the read path in one picture
Each region holds one Store per column family, and each Store is a memstore plus zero or more HFiles. A Get is executed as a single-row Scan. For each family it needs, the region server builds a StoreScanner that opens a scanner on the memstore and on every HFile that could contain the requested keys, and merges them through a heap ordered by key and then by timestamp, newest first.
Before a file joins the heap, HBase tries to exclude it. The file's first and last keys rule it out if the row falls outside them; its time range, recorded in file metadata, rules it out if the Scan's time range does not overlap; and its bloom filter rules it out, with a small false-positive rate, if the row (or row and column) was never written to it. Every file that survives costs a seek: an index lookup, usually served from cached index blocks, and a data-block read that is either a block-cache hit or a read from HDFS.
A cost model for a single-row Get
Let F be the number of HFiles in the store, r the number of files that genuinely contain the row, b the bloom false-positive rate and h the block-cache hit rate for data blocks. Ignoring index blocks, which are small and hot, the expected number of data blocks read is r + (F − r)·b, and the expected number of disk reads is that times (1 − h). Without a bloom filter, b is effectively 1 for files whose key range covers the row, so the cost becomes F blocks.
def get_cost(files, files_with_row, bloom_fp, cache_hit):
"""Expected data-block reads and disk reads for one single-row Get."""
negatives = files - files_with_row
data_blocks = files_with_row + negatives * bloom_fp # blooms skip most negatives
disk_reads = data_blocks * (1 - cache_hit)
return data_blocks, disk_reads
for files in (3, 8, 16):
print(files, "files:", get_cost(files, files_with_row=2, bloom_fp=0.01, cache_hit=0.9),
"no bloom:", get_cost(files, files_with_row=2, bloom_fp=1.0, cache_hit=0.9))With r = 2, b = 1% and a 90% hit rate, a store of 3, 8 or 16 files costs between 2.01 and 2.14 data-block reads and about 0.2 disk reads per Get, because the bloom absorbs the extra files. Without blooms the same Gets cost 3, 8 and 16 blocks and 0.3, 0.8 and 1.6 disk reads. That is why ROW blooms are the default for new column families, and why losing them, through an explicit NONE or through access patterns they cannot serve, turns file count directly into latency.
The model also shows the limits of blooms. They only help point lookups. A range Scan across many rows cannot use a row bloom, so every file whose key range overlaps the scan range participates, and file count feeds straight into seek cost. The prefix bloom type ROWPREFIX_FIXED_LENGTH, configured with RowPrefixBloomFilter.prefix_length, helps scans whose start and stop rows share a fixed-length prefix, which is common with salted or entity-prefixed keys.
How compaction bounds file amplification
Every memstore flush creates a file; compaction merges files. The defaults in hbase-default.xml define the control loop. Minor compaction becomes eligible once a store has hbase.hstore.compactionThreshold = 3 files, merges at most hbase.hstore.compaction.max = 10 files at a time, and selects files using a size ratio of 1.2. If a store reaches hbase.hstore.blockingStoreFiles = 16 files, flushes of that region are blocked until compaction catches up or hbase.hstore.blockingWaitTime = 90,000 ms passes, which throttles writes.
In other words, the defaults give you a soft ceiling of about 16 files per store, enforced by pushing back on writers. The important consequence is order of pain: reads degrade first, as the file count climbs from 3 toward 16, and writes stall only at the ceiling. A cluster that is "fine on writes" can be well into read-amplified territory. Raising blockingStoreFiles to avoid write stalls simply moves the ceiling up and makes reads worse; the real fix is compaction throughput or fewer, larger flushes.
Major compaction, which by default runs every hbase.hregion.majorcompaction = 604,800,000 ms (7 days) with 0.5 jitter, rewrites all files in a store into one and is the only point at which delete markers and excess versions are physically removed. Compaction costs write amplification, which the write amplification article covers; every tuning choice here is a trade between the two.
Cell amplification: versions, tombstones, TTL and filters
A Delete in HBase writes a marker, not an erasure. Until a major compaction removes both the marker and the cells it covers, every read of that row or range has to read the marker, read the covered cells and discard them. The same applies to versions beyond the family's VERSIONS limit and to cells older than the TTL: they are hidden from results immediately but still read. See TTL and versions for the semantics.
The classic trap is a queue or sliding-window pattern: rows are written with increasing keys, processed, and deleted, and a consumer repeatedly scans from the start of the range. Each scan walks through every deleted row since the last major compaction before reaching live data, so the scan gets slower every hour even though the table looks small. Setting KEEP_DELETED_CELLS or a large VERSIONS for audit purposes makes this worse by design.
Filters are the other source. A server-side filter reduces what is returned, not what is read: SingleColumnValueFilter on a non-key column reads every row in the range to decide which to drop, and the region server's filteredReadRequestCount metric climbs while readRequestCount tells you little. If a filter discards most rows, the row key is designed for a different query than the one you are running.
Measuring it
Measure each kind with the tool that sees it. For file amplification, the region server exposes storeFileCount and compactionQueueLength; per-region file counts appear in the master UI and in the shell's detailed status. For block amplification, watch blockCacheHitCount and blockCacheExpressHitPercent alongside HDFS read throughput. For cell amplification, the client-side scan metrics are the most direct signal, because they report rows scanned and rows filtered on the server for one specific scan:
Scan scan = new Scan()
.withStartRow(Bytes.toBytes("sensor42#"))
.withStopRow(Bytes.toBytes("sensor42$"))
.addFamily(Bytes.toBytes("d"))
.setTimeRange(fromMillis, toMillis) // lets the server skip whole HFiles by timestamp
.setCacheBlocks(false); // one-off analytic scans should not evict hot blocks
scan.setScanMetricsEnabled(true);
long returned = 0;
try (ResultScanner rs = table.getScanner(scan)) {
for (Result r : rs) returned++;
Map<String, Long> m = rs.getScanMetrics().getMetricsMap();
long scanned = m.getOrDefault("ROWS_SCANNED", 0L);
long filtered = m.getOrDefault("ROWS_FILTERED", 0L);
System.out.printf("returned=%d scanned=%d filtered=%d amp=%.1f regions=%d%n",
returned, scanned, filtered, scanned / Math.max(1.0, returned),
m.getOrDefault("REGIONS_SCANNED", 0L));
}A ratio of ROWS_SCANNED to rows returned near 1 means the scan is efficient; a ratio of 50 means the server does 50 times the necessary work, usually because of tombstones, filters or a key design that does not match the query. Newer HBase releases add block-level scan metrics as well; check which your version reports. Collect these for your top query shapes in a test harness, not only when there is an incident.
Worked example: a slowing time-series table
An events table stores sensor readings keyed by sensorId#timestamp in one family with three versions retained and no TTL. Each sensor reports once a minute, and a retention job deletes readings older than 30 days every hour. The "last hour for sensor 42" query scans from sensor42# and filters by timestamp on the client side. Over the week since the last major compaction its p99 latency climbs steadily, although the data volume is stable.
The measurements tell the story. storeFileCount per store sits between 12 and 15, because flushes of many small regions outpace minor compaction. The scan metrics report about 53,000 rows scanned for 60 returned. Roughly 43,000 of those are the sensor's live readings from the past 30 days, walked because the start row is unbounded; about 10,000 are readings deleted during the week, still read with their tombstones because no major compaction has run. The tombstones explain the week-long rise; the start row explains why the ratio was never close to 1. The block cache hit rate is low because nightly analytic scans evict hot blocks.
The fixes map one-to-one onto the three kinds. The team replaces the retention job with a 30-day TTL and VERSIONS => 1, so expired data stops being written as tombstones and disappears at compaction; changes the query to start at sensor42# plus the timestamp one hour ago, with a matching time range so older files are skipped, which alone brings the scan to about 60 rows; sets setCacheBlocks(false) on the analytic scans; and merges undersized regions so each store flushes larger files less often. A major compaction then purges the accumulated tombstones, so other range scans over the sensor stop paying for them.
# how many files does each store have, and how big are they?
status 'detailed' # per-region storefiles, sizes, read counts
# inspect one HFile: metadata (bloom type, time range, entry count) and stats
hbase hfile -f hdfs:///hbase/data/default/events/<region>/d/<hfile> -m -s
# add a row bloom, or a fixed-length row-prefix bloom for prefix-shaped Gets
alter 'events', {NAME => 'd', BLOOMFILTER => 'ROW'}
alter 'events', {NAME => 'd', BLOOMFILTER => 'ROWPREFIX_FIXED_LENGTH',
CONFIGURATION => {'RowPrefixBloomFilter.prefix_length' => '10'}}
# bound versions and expire old cells; purge tombstones with a major compaction
alter 'events', {NAME => 'd', VERSIONS => 1, TTL => 2592000}
major_compact 'events'
Trade-offs and failure modes
- ROWCOL blooms help Gets for specific columns in wide rows but cost more memory and help nothing for whole-row reads. Choose from the access pattern.
- Too many column families multiply stores, flushes and files, and a Get that touches all families pays file amplification in each.
- Aggressive major compaction removes tombstones promptly but spends disk and network bandwidth; on large clusters, schedule it off-peak instead of disabling it entirely and forgetting.
- Block size: the 64 KB default suits mixed workloads. Smaller blocks cut wasted bytes on random Gets but enlarge the index; larger blocks favour scans.
- Cache pollution: full-table scans with block caching on evict the hot working set; an off-heap BucketCache enlarges the cache but does not fix a pollution pattern.
- Raising blockingStoreFiles hides compaction debt and converts it into read latency.
What to do next
- Chart
storeFileCountandcompactionQueueLengthper region server and alert when stores approach the blocking threshold. - Enable scan metrics for your top five query shapes and record rows scanned per row returned.
- Confirm every point-lookup family has a ROW or suitable prefix bloom, using
hbase hfile -mon real files. - Replace delete-based retention with TTL where possible and set
VERSIONSto what you actually read. - Pass time ranges and tight start and stop rows on scans; disable block caching for analytic scans.
- Schedule major compactions deliberately and verify afterwards that the scanned-to-returned ratio dropped.