HBase read performance is rarely a problem of averages. A well-run cluster serves most Gets in a millisecond or two, straight from the block cache. The complaints come from the tail: the one request in a hundred that takes 200 ms, which a user-facing service then multiplies by fanning each page view out to many rows. Tuning that tail means finding which layer the slow requests are waiting in, and each layer has different causes and fixes.
This article is a diagnosis procedure. It follows a Get from the client down to the disk, names the metric that shows whether each layer is at fault, and gives the fixes. The mechanics of the store itself, the per-Get cost model over memstore and HFiles, are covered in the read amplification article and are only summarised here.
Why the tail is the metric
Suppose a page needs 20 rows that live in different regions, fetched in parallel. The page is as slow as its slowest row. If each Get independently exceeds its own p99 one time in a hundred, the probability that at least one of the 20 does is 1 - 0.9920, about 18 percent. The backend's p99 has become roughly the page's p82. At 100 rows it is 63 percent. That is why HBase read tuning targets p99 and p999, and why a small fraction of slow servers, cold regions or pauses dominates the user experience.
It also tells you what to measure. Per-request latency percentiles, from both the client and the RegionServers, broken down by table and by server. Averages hide exactly the thing you are looking for.
The layers a read passes through
The client finds the region from its cached copy of hbase:meta and sends an RPC to the RegionServer. The request waits in a call queue until a handler thread is free. The handler opens scanners over the memstore and the relevant store files, skipping files by time range and bloom filter, and reads blocks: from the on-heap L1 cache for index and bloom blocks, from the BucketCache for data blocks if configured, or from HDFS on a miss. HDFS reads either the local disk directly (short-circuit) or goes over the network to a DataNode. Every layer adds to the tail independently, so the first job is to work out which one.
Step 1: split the latency
Start with three numbers per RegionServer, taken over the same window as the client's complaint. RegionServer IPC metrics report queueCallTime (time waiting for a handler) and processCallTime (time on a handler) as histograms with percentile fields, and the server metrics report Get_99th_percentile. Compare them with the client's own measured p99.
| Observation | The tail lives in | Go to |
|---|---|---|
| Client p99 far above server p99 | Network, retries, meta lookups, client GC | Step 5 |
| queueCallTime p99 high, processCallTime normal | Handler saturation | Step 2 |
| processCallTime p99 high on all servers | Store layout or cache | Step 3 |
| processCallTime p99 high on a few servers | Locality, disk, GC on those hosts | Step 4 |
Always break the numbers down by server before by table. Averaging across a hundred RegionServers hides the three that cause most of the tail. A simple heat map of per-server p99 over the last day is usually the fastest diagnostic you have.
Step 2: handlers and call queues
A RegionServer has a fixed pool of handler threads, hbase.regionserver.handler.count, 30 by default. If every handler is busy, a Get waits in the queue behind whatever is running, and long scans are the usual offender: a handful of big scans can occupy every handler while millisecond Gets queue behind them. The metrics numActiveHandler and numCallsInGeneralQueue show it directly.
The fix is to separate the traffic. HBase can split handlers into several queues and dedicate a share to reads, and within reads a share to scans, so short Gets never wait behind long scans.
<!-- hbase-site.xml on the RegionServers -->
<property><name>hbase.regionserver.handler.count</name><value>60</value></property>
<!-- spread handlers over several queues instead of one shared queue -->
<property><name>hbase.ipc.server.callqueue.handler.factor</name><value>0.1</value></property>
<!-- 60% of queues/handlers serve reads, the rest writes -->
<property><name>hbase.ipc.server.callqueue.read.ratio</name><value>0.6</value></property>
<!-- of the read handlers, 30% serve scans, 70% serve Gets -->
<property><name>hbase.ipc.server.callqueue.scan.ratio</name><value>0.3</value></property>Don't just raise the handler count. More handlers mean more concurrent requests competing for the same CPU, heap and disks, and past the point where those saturate you trade queue time for processing time and more garbage collection. Raise it in steps while watching both histograms. If a queue grows because the disks are the limit, the fix is lower down. Large scans should also set sensible caching and result-size limits so each RPC finishes quickly. The scan article covers that sizing.
Step 3: the store and the cache
When processing time is high everywhere, look at how much work each Get does. Three numbers answer it. Store files per store: each extra HFile a Get cannot skip is another index lookup and possibly another disk read, so a storeFileCount climbing because compaction is behind shows up directly in p99. Bloom filters: with BLOOMFILTER => 'ROW', the default, a Get skips files that cannot contain the row, and a table created with blooms off pays for every file. The bloom filter article explains when ROWCOL helps. Cache hit ratio: watch blockCacheExpressHitPercent, the hit ratio for requests that asked to be cached, and the per-tier hit counts.
Hit ratio is about working set. If the hot data for a table is 400 GB per RegionServer and the cache is 30 GB, no tuning of eviction will fix it. Either add an off-heap or file-backed BucketCache sized to the working set, or accept disk reads and make them fast. A useful estimate: hot bytes are roughly rows read per hour times average row size times blocks per row. Remember that the default 64 KB block size means a 1 KB random read pulls a 64 KB block into cache. For point-read tables, a 16 KB block size and a DATA_BLOCK_ENCODING such as FAST_DIFF or ROW_INDEX_V1 let the same cache hold more useful rows. Smaller blocks mean a larger index, so measure before changing a large table.
Step 4: locality, short-circuit and hedged reads
HBase stores its files in HDFS, and a RegionServer reads fastest when the blocks of its regions sit on its own DataNode. The metric percentFilesLocal reports this. Locality drops whenever regions move, after a restart, a balancer run or a failure, because the files stay where they were written. Until a major compaction rewrites them locally, every cache miss goes over the network to another DataNode. A server with locality at 0.4 after a rolling restart is a classic source of tail latency.
Three settings shape the storage layer's tail. Short-circuit reads (dfs.client.read.shortcircuit with a dfs.domain.socket.path) let the RegionServer read local block files directly instead of streaming through the DataNode process. HBase-level checksums, on by default through hbase.regionserver.checksum.verify, avoid a second read of the HDFS checksum file. Hedged reads start a second read against another replica if the first is slow, and take whichever answers first.
<property><name>dfs.client.read.shortcircuit</name><value>true</value></property>
<property><name>dfs.domain.socket.path</name><value>/var/lib/hadoop-hdfs/dn_socket</value></property>
<!-- hedged reads: >0 threads enables them -->
<property><name>dfs.client.hedged.read.threadpool.size</name><value>20</value></property>
<!-- wait this long for the first replica before asking another -->
<property><name>dfs.client.hedged.read.threshold.millis</name><value>30</value></property>Hedging cuts the tail caused by one slow disk or busy DataNode, at the cost of extra I/O. Set the threshold near the p99 of healthy reads so only the slow few are hedged, and watch the hedgedReads and hedgedReadWins counters. If wins are a large share of hedges, something below is persistently slow and hedging is masking it. Find the disk.
Step 5: the client side
Two defaults cause long tails that no server tuning can fix. The first is retry behaviour. A client whose RPC timeout is long and whose retry count is high will wait many seconds behind a RegionServer that has just died, until the region is reassigned. Set hbase.rpc.timeout for a single attempt and hbase.client.operation.timeout for the whole operation to values that fit your service's budget, and fail fast instead of hanging. The second is meta lookups. A client that is recreated per request loses its region location cache. Share one Connection per process.
For reads that can tolerate slightly stale data, region replicas cut the tail during failures and pauses. With timeline consistency, the client sends the Get to the primary and, if no answer arrives within hbase.client.primaryCallTimeout.get (in microseconds, default 10000, which is 10 ms), also asks the secondaries and takes the first result.
try (Connection conn = ConnectionFactory.createConnection(conf);
Table table = conn.getTable(TableName.valueOf("profiles"))) {
Get get = new Get(Bytes.toBytes(userId));
get.setConsistency(Consistency.TIMELINE); // allow a secondary to answer
Result r = table.get(get);
if (r.isStale()) {
metrics.increment("profiles.stale_reads"); // served by a secondary replica
}
render(r);
}Timeline reads can go backwards in time between two requests if they hit different replicas. Use them for display data, not for read-modify-write logic, and count stale reads so you know how often it happens.
The JVM in the tail
A 400 ms garbage-collection pause becomes a 400 ms p999 on every request in flight on that server. Keep the heap moderate and move the data cache off-heap with BucketCache so the collector does not have to trace it. Use a low-pause collector suited to your Java version, and log GC pauses next to your latency histograms. If the worst latencies line up with pauses, heap tuning is the fix, not HBase configuration.
Worked example: a profile service at 180 ms p99
A profile service reads 30 rows per page from a 40-node cluster. Server Get p99 averages 25 ms, but page p99 is 180 ms. Step 1: per-server p99 shows five servers above 150 ms and the rest under 15 ms. Queue times on those five are normal, so the handlers are not the problem. Step 4 finds the cause: all five were restarted in last week's rolling upgrade and report percentFilesLocal around 0.35, and their cache misses go to remote DataNodes, one of which has a disk reporting media errors.
The fixes were ordered by speed. Replace the failing disk and enable hedged reads with a 30 ms threshold, which brought those servers under 40 ms the same day. Run major compactions on the affected regions off-peak to restore locality. Add a post-restart check to the upgrade runbook that holds the next batch until locality recovers. Then enable timeline Gets for the profile table, whose data tolerates brief staleness. Page p99 settled at 28 ms. Nothing in the table schema changed.
Failure modes and trade-offs
- Tuning the average. Cluster-wide means improve while the few bad servers stay bad. Always look per server.
- Handler inflation. Raising handlers without splitting queues moves the wait from the queue into CPU and GC.
- Hedging as a mask. A high hedge win rate hides failing hardware. Alert on it.
- Compaction for locality at peak. Major compactions restore locality but compete with reads for disk. Schedule them off-peak and throttle them.
- Stale reads in the wrong place. Timeline consistency in write paths causes lost updates.
What to do next
- Record client-side and server-side p99 per RegionServer and per table, and plot per-server p99 as a heat map.
- Compare queueCallTime with processCallTime on the slow servers to decide which layer to work on.
- If queueing, split read and write queues and give scans their own share before raising handler counts.
- If processing, check store file counts, blooms, block size and cache hit ratio against the working set.
- Alert on percentFilesLocal after restarts, enable short-circuit reads, and evaluate hedged reads with a threshold near healthy p99.
- Set explicit RPC and operation timeouts, share one Connection per process, and consider timeline replicas for read-only data.