Most HBase dashboards are wrong in small ways that add up to big wrong conclusions. A cluster-wide cache hit percentage is computed by averaging per-server percentages. A p99 latency panel shows the mean of twenty servers' p99s. A request-rate graph drops to zero and jumps back after every rolling restart. None of these is a collection problem. The numbers arrive correctly; the arithmetic applied to them is wrong because the type of each metric was ignored.
This page is about that arithmetic. It explains the three kinds of metric HBase exposes, how HBase computes its latency histograms and what window a percentile covers, which aggregations are valid for each kind, how to derive rates and ratios that stay correct across restarts, what per-region metrics cost, and how to check all of this against your own cluster. It assumes you already collect metrics; for the collection pipeline and a layer-by-layer reading of the main metrics, see HBase metrics and monitoring. Metric names below are the JMX names; Prometheus exporters rename them, so map them before copying queries.
Three kinds of metric
Every HBase metric is one of three kinds, and the kind decides what you may do with it.
A counter counts events since the process started and only goes up: readRequestCount, writeRequestCount, blockCacheHitCount, blockCacheMissCount, slowGetCount. Its absolute value is almost meaningless (it depends on uptime); what matters is how fast it grows. When the RegionServer restarts, every counter starts again from zero.
A gauge is a current value that can go up or down: memStoreSize, regionCount, storeFileCount, flushQueueLength, compactionQueueLength, and in the IPC source queueSize and numActiveHandler. A gauge is sampled at the moment it is read, so a scrape every 30 seconds sees 30-second-apart snapshots and misses everything in between. Short bursts in a queue length are invisible to a slow scrape.
A histogram summarises a distribution of measurements, usually latencies or sizes. For an operation such as Get in the RegionServer Server source, HBase emits a family of values: Get_num_ops, Get_min, Get_max, Get_mean, Get_median, and percentiles such as Get_75th_percentile, Get_90th_percentile, Get_95th_percentile, Get_99th_percentile and Get_99.9th_percentile. The IPC source has the same structure for QueueCallTime, ProcessCallTime and TotalCallTime.
Some values look like gauges but are derived from counters, such as blockCacheCountHitPercent. They are convenient for a glance at one server and dangerous for anything else, as the worked example shows.
How HBase computes its histograms
The histogram family is where most misreadings start, so it is worth knowing how it is built. In current HBase source, the metrics histogram implementation records each measurement into a bucketed structure called FastLongHistogram. Bucketing makes recording cheap enough to do on every request, at the cost of approximate percentiles: a reported p99 is accurate to within a bucket, not exact.
When the metrics system asks the histogram for its values, the implementation calls snapshotAndReset() on the underlying structure. The min, max, mean and percentiles therefore describe only the measurements recorded since the previous snapshot. The _num_ops value is different: it is reported as a counter and comes from a separate count that is not reset, so it grows for the life of the process.
Three consequences follow. First, a percentile is the percentile of one interval, not of all time and not of your dashboard's time range. Second, the length of that interval is decided by when the metrics system takes snapshots, which depends on the metrics configuration and on how reads through JMX are cached, not simply on your scrape interval; if two collectors read the same server independently, the interval each one sees can be affected by the other. Third, _num_ops is the only part of the histogram you can turn into a rate. Verify the behaviour on your own version with a quick test: read the JMX endpoint twice a few seconds apart under steady load and see whether the percentiles move and the count keeps rising.
# One RegionServer's Server source, filtered with the qry parameter (default info port 16030)
curl -s 'http://rs1.example.com:16030/jmx?qry=Hadoop:service=HBase,name=RegionServer,sub=Server' \
| python -m json.tool | grep -E '"(Get_num_ops|Get_99th_percentile|blockCacheHitCount|blockCacheMissCount)"'
Why percentiles and percentages cannot be averaged
Suppose two RegionServers each report a p99 Get latency. Server A served 9,900 Gets with a p99 of 5 ms; server B served 100 Gets with a p99 of 200 ms. The average of the two p99s is 102.5 ms, a number that describes neither server and is not the p99 of the combined 10,000 requests. The combined p99 is the boundary of the slowest 100 requests, and where it lands depends on how A's slow requests and B's requests interleave: near 5 ms if most of B's requests were fast, near B's fastest request if all of B was slow. The two p99s cannot tell you which. Averaging also gave B, with 1 percent of the traffic, half the weight.
There is no way to compute an exact cluster-wide percentile from per-server percentiles, because the percentiles have thrown away the shape of each distribution. The honest options are: show the maximum of the per-server p99s (the worst server, which is what an unlucky client sees); show each server separately as a heat map; or measure latency at the client, where a single histogram covers every request. If you need a weighted summary, weight by _num_ops rate and label the result as an approximation.
Percentages derived from counters have the same problem in a milder form. A server with a 50 percent cache hit rate on 10 requests and one with 99 percent on 100,000 requests average to 74.5 percent, while the real cluster rate is almost 99 percent. The fix is always the same: go back to the counters.
Worked example: a true cluster-wide cache hit ratio
A five-server cluster reports these block cache counter increases over a five-minute window:
| Server | Hit increase | Miss increase | Per-server hit % |
|---|---|---|---|
| rs1 | 1,200,000 | 60,000 | 95.2 |
| rs2 | 1,150,000 | 70,000 | 94.3 |
| rs3 | 1,180,000 | 65,000 | 94.8 |
| rs4 | 40,000 | 40,000 | 50.0 |
| rs5 | 1,210,000 | 55,000 | 95.7 |
Averaging the per-server percentages gives (95.2 + 94.3 + 94.8 + 50.0 + 95.7) / 5 = 86.0 percent. Summing the counters first gives hits of 4,780,000 and misses of 290,000, so the cluster ratio is 4,780,000 / 5,070,000 = 94.3 percent. The averaged figure would suggest a cache that is badly undersized; the correct figure says the cache is fine and one server is odd.
The interesting finding is rs4: few requests and a 50 percent hit rate. That pattern usually means a server that recently restarted with a cold cache, or one hosting regions with a scan-heavy access pattern. Either way, it is a per-server question, which is exactly why you should keep the per-server view next to the cluster ratio rather than replacing one with the other.
In PromQL the correct calculation is rate first, then sum, then divide. The metric names here depend on your exporter, so substitute the ones you actually have:
# Cluster-wide block cache hit ratio over 5 minutes (names depend on your exporter)
sum(rate(hbase_regionserver_server_blockcachehitcount[5m]))
/
( sum(rate(hbase_regionserver_server_blockcachehitcount[5m]))
+ sum(rate(hbase_regionserver_server_blockcachemisscount[5m])) )
# Worst-server p99 Get latency: max, never avg
max(hbase_regionserver_server_get_99th_percentile)
# Cluster Get rate from the cumulative _num_ops counter
sum(rate(hbase_regionserver_server_get_num_ops[5m]))
Counters across restarts
Because counters restart at zero, any arithmetic on raw values breaks across a rolling restart. A panel that plots the difference between consecutive raw values shows a huge negative spike; a sum of raw counters across servers drops by the restarted server's whole history. Prometheus's rate() and increase() functions detect a counter going down and treat it as a reset, which fixes most of this, as long as you apply them before summing across servers. Summing first and taking the rate of the sum hides resets inside the total and produces exactly the artefacts you were trying to avoid.
If you compute rates yourself, for example in a script that polls JMX, apply the same rule: if the new value is lower than the old one, treat the new value as the increase since the restart. The sketch below polls one server and computes the hit ratio and the Get rate for each interval, handling resets:
import json, time, urllib.request
URL = ("http://rs1.example.com:16030/jmx"
"?qry=Hadoop:service=HBase,name=RegionServer,sub=Server")
KEYS = ("blockCacheHitCount", "blockCacheMissCount", "Get_num_ops")
def read():
bean = json.load(urllib.request.urlopen(URL, timeout=5))["beans"][0]
return {k: bean[k] for k in KEYS}
def delta(new, old):
return new if new < old else new - old # counter went down: process restarted
prev = read()
while True:
time.sleep(60)
cur = read()
d = {k: delta(cur[k], prev[k]) for k in KEYS}
lookups = d["blockCacheHitCount"] + d["blockCacheMissCount"]
ratio = d["blockCacheHitCount"] / lookups if lookups else float("nan")
print(f"gets/s={d['Get_num_ops'] / 60:.1f} hit_ratio={ratio:.3f}")
prev = cur
Gauges: sampling, burstiness and what they hide
Gauges have the opposite weakness. They are easy to aggregate (sum the memStoreSize of all servers to see total unflushed data; take the maximum compactionQueueLength to find the worst server), but a gauge only shows the moment it was read. An RPC queue that fills for five seconds every minute can look empty on a 30-second scrape.
Two techniques recover the missing information. First, pair each important gauge with a related counter or histogram that does cover the whole interval: the IPC QueueCallTime histogram reveals queueing that queueSize sampling misses, and slowGetCount or slowPutCount increments show slow operations even if no scrape caught them in progress. Second, chart gauges with max-over-time rather than average-over-time when you care about saturation.
Per-region and per-table metrics: cardinality is the cost
HBase also publishes metrics per region and per table, which are invaluable for finding a hot region or a table whose read pattern changed. They also multiply. A cluster with 20 servers hosting 300 regions each, with a few dozen metrics per region, produces hundreds of thousands of time series, and each region move or split creates new series and abandons old ones. Monitoring systems that charge or slow down by series count feel this quickly.
The practical approach is to collect server-level and table-level metrics continuously and treat region-level metrics as a diagnostic tool. Use the JMX qry parameter or your exporter's allow-list to limit which beans are scraped, and either scrape region beans on demand or at a longer interval. When you do keep region metrics, put region names in labels you can drop later, not in metric names.
Failure modes in metric interpretation
- Averaged percentiles. The cluster p99 panel looks calm while one server is slow. Fix: chart max and per-server values, and measure at the client for true end-to-end percentiles.
- Averaged percentages. The cluster hit ratio looks low because a cold or idle server drags the mean. Fix: sum counter rates and divide.
- Restart artefacts. Negative spikes or step drops in request graphs after maintenance. Fix: rate before sum.
- Two collectors, one histogram. Percentile windows behave oddly after a second monitoring agent is added. Fix: have one collector read JMX and fan the data out.
- Invisible bursts. Queue gauges look healthy while clients time out. Fix: watch queue-time histograms and slow-operation counters.
- Series explosion. The monitoring system slows down after region metrics are enabled. Fix: limit region beans and use on-demand scrapes.
What to do next
- List every panel and alert on your HBase dashboard and label each input as counter, gauge or histogram.
- Rewrite any average of percentiles as a max or a per-server view, and any average of percentages as a ratio of summed counter rates.
- Check that every counter goes through rate() or increase() before it is summed across servers.
- Run the two-read JMX test on one server to confirm how percentiles reset on your HBase version and collection setup.
- Make sure exactly one collector reads each server's JMX endpoint.
- Decide which region-level metrics you need continuously and restrict the rest with qry filters or an allow-list.