A hotspot is a small part of an HBase table that receives a disproportionate share of the traffic, so one RegionServer works flat out while the rest of the cluster idles. Latency rises for every client whose requests touch that server, the server's call queue grows, and adding nodes does nothing. The fixes are well known: salting, hashing, pre-splitting, caching, rebalancing. The hotspotting and row-key design article covers them. But each fix only works for one cause, and applying the wrong one wastes days or makes things worse.
This article is about the part before the fix: proving there is a hotspot, locating it down to the server, the region and, where possible, the key, and classifying it. You will get a taxonomy, the data sources HBase actually exposes, a script that computes per-region request rates from JMX, a sampling technique for finding hot keys, a worked example, and a decision table. The commands and metric names below were checked against the HBase source and documentation; verify flags against the version you run.
A taxonomy of hotspots
Start by naming the possibilities, because each one leaves a different fingerprint.
| Type | Fingerprint | Split helps? |
|---|---|---|
| Monotonic write hotspot | Writes concentrate on the region holding the end of the key space; after a split the new tail region is hot | No: the hotspot follows the tail |
| Hot key range | One region has high read or write rate across many distinct keys | Yes, if the split point divides the traffic |
| Single hot row | One region is hot and almost all requests hit one row key | No: a row cannot be split |
| Scan hotspot | Read rate or bytes scanned high on a region; few requests but expensive ones | Sometimes; usually fix the scan |
| False hotspot | One server is slow but request rates per region are even | No: the problem is the server |
The last row matters more than it looks. A RegionServer with long garbage-collection pauses, a failing disk, poor data locality or simply more regions than its peers shows the same symptom, high latency on one server, without any skew in the data. Rebalancing or splitting will not help, and salting keys would be a costly rewrite for nothing.
Where the evidence comes from
Per-region counters. Each RegionServer publishes a JMX bean named Hadoop:service=HBase,name=RegionServer,sub=Regions. It holds one attribute per region and metric, named Namespace_<ns>_table_<table>_region_<encoded>_metric_<name>, where the region part is the encoded region name. The useful metrics include readRequestCount, writeRequestCount, filteredReadRequestCount, memStoreSize and storeFileSize. How this bean gets exported to Prometheus or other systems is covered in the metrics and monitoring article.
These counters are cumulative since the region opened on its current server, and they start again when the region moves or splits. That makes a naive ranking wrong: a region open for a month outranks a region that is hot right now. Always take two snapshots and rank by the difference over the interval, and discard any region whose counter went down.
hbtop. Since HBase 2.1.7, 2.2.2 and 2.3.0, hbase hbtop gives a top-like live view. Its default Region mode shows #REQ/S, #READ/S, #FREAD/S, #WRITE/S, store file size and count, memstore size, locality and the region's start key (SKEY), sorted by request rate. Use -m s for RegionServer mode and -m t for table mode; press i to drill into a record, s to change the sort field, and o to filter. Recent versions add -b for batch output you can log.
The online slow log. From HBase 2.3, with hbase.regionserver.slowlog.buffer.enabled set to true (it defaults to false), each RegionServer keeps a ring buffer of slow and large RPCs. In the shell, get_slowlog_responses '*', {'TABLE_NAME' => 'events'} reads all servers, filtered by table; region, client address and user filters also exist. Each record carries the region name, method, client address, processing and queue times, response size and block bytes scanned, which tells you which clients and which regions send the expensive calls.
The client side. None of the server sources reliably tells you which row keys are hot. For that you need sampling in the client, described below.
The workflow
Work from coarse to fine. Confirm the symptom first: p99 latency by server, call-queue length and handler saturation. If one server stands out, compare its per-region request rates. If one or two regions carry most of that server's traffic, you have a data hotspot and can classify it. If its regions look like everyone else's, it is a false hotspot: look at GC logs, disk latency, locality and region count, and at whether the balancer has recently moved regions onto it.
import json, re, time, urllib.request
BEAN = "Hadoop:service=HBase,name=RegionServer,sub=Regions"
KEY = re.compile(r"^Namespace_(.+)_table_(.+)_region_([0-9a-f]+)_metric_"
r"(readRequestCount|writeRequestCount|filteredReadRequestCount)$")
def snapshot(rs_hosts, port=16030):
out = {}
for host in rs_hosts:
url = f"http://{host}:{port}/jmx?qry={BEAN}"
bean = json.load(urllib.request.urlopen(url, timeout=10))["beans"][0]
for k, v in bean.items():
m = KEY.match(k)
if m:
ns, table, region, metric = m.groups()
out[(host, ns, table, region, metric)] = v
return out
def hot_regions(rs_hosts, interval_s=60, top=10):
a = snapshot(rs_hosts); time.sleep(interval_s); b = snapshot(rs_hosts)
rates = {}
for key, v in b.items():
if key in a and v >= a[key]: # region moved or split: counter restarted
rates[key] = (v - a[key]) / interval_s
by_metric = {}
for (host, ns, table, region, metric), r in rates.items():
by_metric.setdefault(metric, []).append((r, host, f"{ns}:{table}", region))
for metric, rows in by_metric.items():
rows.sort(reverse=True)
total = sum(r for r, *_ in rows) or 1
print(metric)
for r, host, table, region in rows[:top]:
print(f" {r:10.1f}/s {100*r/total:5.1f}% {table} {region} on {host}")The script polls each RegionServer's /jmx servlet with a qry filter (16030 is the default RegionServer info port), takes two snapshots, and prints the top regions by rate with their share of the total. Run it with a 60-second interval during the problem period. A healthy table shows shares roughly proportional to region count; a hotspot shows one region with a share many times its fair share.
Finding the hot key
Once you have the region, the next question is whether many keys or one key are hot. Instrument the client: wrap the table access path, and for a sampled fraction of calls, feed the table, operation and a row-key prefix into a heavy-hitter sketch. The Space-Saving algorithm keeps a fixed number of counters and is guaranteed to retain every key whose true frequency exceeds the total divided by its capacity, which is exactly the property you need to catch a celebrity row among millions of cold ones.
class SpaceSaving:
"""Heavy-hitter sketch: finds keys above 1/capacity of traffic in fixed memory."""
def __init__(self, capacity=1000):
self.cap, self.counts = capacity, {}
def add(self, key):
if key in self.counts or len(self.counts) < self.cap:
self.counts[key] = self.counts.get(key, 0) + 1
return
victim = min(self.counts, key=self.counts.get) # use a heap in production
self.counts[key] = self.counts.pop(victim) + 1 # overestimates by at most victim's count
def top(self, n=20):
return sorted(self.counts.items(), key=lambda kv: -kv[1])[:n]
# In the client wrapper, for a sampled fraction of calls:
# sketch.add((table, op, row_key[:16])) # prefix groups rows sharing a key designReport the sketch's top entries per client process to your metrics system and merge them centrally. If the top key accounts for most of the hot region's traffic, you have a single hot row. If the top hundred keys share a prefix such as today's date or one tenant ID, you have a hot range caused by key design. If the keys are spread evenly but all fall in one region's range, the region is simply too large a slice of a popular range.
The split test. When sampling is not possible, a controlled split answers the same question. Split the hot region at a key near the middle of its traffic, for example split 'ENCODED_REGION_NAME', 'splitkey' in the shell, let the balancer or a manual move put the daughters on different servers, and re-run the rate script. If the load divides, it was a range; if one daughter keeps nearly all of it, the heat is concentrated on a few keys. If the new tail region becomes hot, the keys are monotonic.
Worked example
An events table with 40 regions on 8 RegionServers sees write p99 latency climb every day from 09:00. hbtop in RegionServer mode shows one server at about six times the write rate of the others. The rate script over 60 seconds shows a single region taking 71 percent of all writes to the table, and its start key is a timestamp from earlier that morning. The row key design is yyyyMMddHHmmss|deviceId.
The classification follows directly. Many distinct keys hit the region, so it is not a single row. Every key is newer than the last, so the tail region absorbs all writes. The team splits the region as a test, and within minutes the new tail region is hot, confirming a monotonic hotspot. No amount of splitting or balancing can fix it. The fix is a key redesign, such as a salt or hash prefix derived from the device ID, with a migration plan; the hotspotting article compares the options. The team adds a dashboard alert on the top region's share of writes so a regression is caught the day it ships.
A second case on the same cluster looks similar but is not. A profiles table has one region at 40 percent of reads. The client sketch shows that a single row, a shared configuration record read on every request, accounts for most of it. Splitting would only move that row. The fix is to cache it in the application for a few seconds, and for reads that tolerate slightly stale data, to consider region replicas with timeline consistency.
From classification to fix
| Cause | Fix | What to verify afterwards |
|---|---|---|
| Monotonic keys | Salt, hash or reverse the leading key component; bucket time | Write share of the top region falls toward its fair share |
| Hot key range | Split at the traffic midpoint; pre-split new tables; let the balancer move daughters | Load divides between daughters on different servers |
| Single hot row | Application cache, read replicas, or spread the row's data across several rows (for example sharded counters) | Top key's share of the region drops |
| Expensive scans | Tighten start and stop rows, add filters that run server-side, reduce caching batch size | Block bytes scanned per call falls in the slow log |
| False hotspot | Fix GC, disk or locality; even out region counts | Server latency matches peers with no data change |
For the balancing step, know what the balancer optimises. The StochasticLoadBalancer weighs region counts, request load, locality and other costs, so it may leave a hot region where it is if moving it would worsen another cost. After splitting, check where the daughters actually landed, and move them yourself if needed.
Mistakes that mislead the analysis
- Ranking cumulative counters. Old regions look hot. Always use deltas over an interval.
- Measuring after the balancer has acted. A region moved mid-interval restarts its counters and disappears from the ranking. Pause the balancer during diagnosis, or discard intervals with moves.
- Counting requests, not cost. A small get and a scan over large rows each look like requests, but their cost differs by orders of magnitude. Cross-check with bytes scanned and response size from the slow log.
- Sampling interval too long. A five-minute burst at 09:00 averages away over an hour. Match the interval to the symptom.
- Ignoring hbase:meta. Clients that do not cache region locations hammer the meta table; a hot meta region looks like cluster-wide slowness.
- Fixing before classifying. Salting a table whose problem is one hot row costs a migration and changes nothing.
What to do next
- Add the per-region delta script, or an equivalent query over your exported metrics, to your runbook, and record the top region's share per table as a standing metric.
- Enable the online slow log ring buffer on RegionServers and practise reading it with get_slowlog_responses.
- Install hbtop on an operations host and learn Region and RegionServer modes before an incident.
- Add sampled heavy-hitter tracking of row-key prefixes to your client wrapper.
- For every hotspot, write down its class from the taxonomy before choosing a fix, and verify the fix with the same rate measurement.
- Read the row-key design guidance for the fixes, and alert when any region's share of traffic exceeds several times its fair share.