You have found the hotspot: one region, on one RegionServer, takes a disproportionate share of requests, and its server's latency drags every table it hosts. Diagnosis is covered in HBase hotspot analysis, which classifies the cause. This article is the next step: carrying out the fix on a live production table, safely, in the order that buys the most relief for the least risk.

The guiding rule is that fixes differ in cost and reversibility. Moving a region takes seconds and can be undone. Splitting is quick but cannot easily be reversed under load. Changing the row key needs a data migration and application changes. So match the fix to the cause, start with the cheapest action that addresses it, and know in advance which causes the cheap actions cannot fix.

Match the fix to the shape of the hotspot

Picking a hotspot fix: cheapest reversible action firstOne region is hotdiagnosis doneOne row?top key share > 50%Key range, spread?many keys in regionMonotonic tail?writes at last regionReads: cache, replicaswrites: shard the rowSplit at traffic midpointthen place daughtersChange the keysalt or hash, migrateMeanwhile: quota the noisy clientprotects neighbours while the real fix landsSplitting never helps a single row or a monotonic tail; it only divides a range that is already spread.
The shape of the hot traffic decides the fix. Only a range that is already spread across many keys benefits from splitting; a single row and a monotonic tail need different tools.

Three shapes cover most production hotspots. A single hot row: one key, such as a popular product or a global counter, takes most of a region's traffic. A hot range: many keys in one region are busy, for example one large tenant whose keys share a prefix. A monotonic tail: keys begin with a timestamp or sequence, so every new write lands in the last region, whichever region that currently is.

Splitting divides a range of keys between two regions. It therefore helps only the hot range. A single row is the atomic unit of an HBase region and cannot be split. A monotonic tail moves on: after a split, the upper daughter is the new tail and gets all the writes, so the hotspot just changes its region name.

Splitting a hot range

For a hot range, split at the point that divides the traffic, not the data. HBase's own split point is the midpoint of the region's largest store file, which divides bytes. If 90% of the requests go to keys in the first tenth of the region, a byte split leaves almost all the load in one daughter.

Find the traffic midpoint by sampling keys from real requests, for example from the application's request log or from a short sample of the RegionServer's request log, and taking the median key. Then split there and check where the daughters land.

# HBase shell
split 'events', 'tenant0042|'                    # explicit split point between hot keys
list_regions 'events'                            # see daughters, servers and request counts
move 'a3f9c1e2...encoded', 'rs7.example.com,16020,1727950000000'   # place one daughter elsewhere

The same in Java, useful when an automated tool picks the split point:

try (Connection conn = ConnectionFactory.createConnection(conf);
     Admin admin = conn.getAdmin()) {
    TableName t = TableName.valueOf("events");
    byte[] splitKey = Bytes.toBytes(medianKeyFromSample(sampledKeys));
    admin.split(t, splitKey);                    // asynchronous; poll regions until daughters appear
    // later, once daughters are online:
    admin.move(daughterEncodedName, ServerName.valueOf("rs7.example.com", 16020, startCode));
}

Two constraints apply. A daughter region holds reference files that point into the parent's store files until a compaction rewrites them, and a region with references cannot be split again. If you need several splits, wait for the daughters to compact, or trigger a major compaction on them first. And splits are not free to undo: merging regions back is possible with merge_region, but under write load it is another disruption. Split only when sampling shows the load really divides.

Placement and the balancer

After a split, both daughters often stay on the same server, which divides the region but not the server's load. The balancer decides placement. HBase's StochasticLoadBalancer minimises a weighted sum of costs, including region count per server, locality, and read and write request load. The request-load weights are set by hbase.master.balancer.stochastic.readRequestCost and hbase.master.balancer.stochastic.writeRequestCost, both 5 by default, small compared with the region-count weight. With default weights, the balancer may judge an even region count more valuable than spreading a hot region's traffic.

If request skew is your recurring problem, raise those two weights and observe several balancer runs. Change one weight at a time, because higher request weights make the balancer react to short traffic spikes and move regions more often, and every move costs a brief unavailability and some cache warmth. For a one-off incident, a manual move is simpler and more predictable. Disable the balancer with balance_switch false while you move regions by hand, or it may move them back, and remember to turn it on again.

Containment with quotas

While a real fix is underway, protect neighbouring tables with RPC quotas. With hbase.quota.enabled set to true, you can throttle a user, a namespace or a table.

set_quota TYPE => THROTTLE, TABLE => 'events', LIMIT => '5000req/sec'
set_quota TYPE => THROTTLE, USER => 'batch_loader', THROTTLE_TYPE => WRITE, LIMIT => '20M/sec'
list_quotas

Throttled clients receive a throttling exception and back off. That is the intent, but check client retry settings first: clients that retry immediately and in large numbers turn a throttle into a retry storm. Quotas are a containment tool, not a fix; remove them once the hotspot is gone.

Fixing a single hot row

A single hot row is fixed differently depending on whether it is read or written.

For reads, the cheapest fix is a cache. A hot row is already in the block cache, so the bottleneck is the RegionServer's handler threads and network, not disk. An application-level cache with a short TTL, or request coalescing so a thousand concurrent requests for the same key become one, takes the load off HBase entirely. Where clients can accept slightly stale data, region replicas help too: with REGION_REPLICATION set to 2 or 3, secondary replicas serve reads issued with timeline consistency, which spreads one region's reads over several servers. Reads from a secondary can be behind the primary, by milliseconds when asynchronous WAL replication to secondaries is enabled and by much longer when secondaries only pick up flushed files, and clients can check whether a result is stale.

alter 'catalog', REGION_REPLICATION => 2

// Java client: allow a secondary replica to answer
Get g = new Get(Bytes.toBytes("product#12345"));
g.setConsistency(Consistency.TIMELINE);
Result r = table.get(g);
boolean maybeStale = r.isStale();

For writes to a single row, typically a counter, spread the row over N shard rows and sum on read. Writes pick a random shard; reads fetch all N, which is cheap because N is small.

static final int SHARDS = 16;

void increment(Table t, String counter, long delta) throws IOException {
    int shard = ThreadLocalRandom.current().nextInt(SHARDS);
    byte[] row = Bytes.toBytes(String.format("%02d|%s", shard, counter));
    t.incrementColumnValue(row, CF, Q_COUNT, delta);
}

long read(Table t, String counter) throws IOException {
    List<Get> gets = new ArrayList<>();
    for (int s = 0; s < SHARDS; s++) gets.add(new Get(Bytes.toBytes(String.format("%02d|%s", s, counter))));
    long sum = 0;
    for (Result r : t.get(gets)) if (!r.isEmpty()) sum += Bytes.toLong(r.getValue(CF, Q_COUNT));
    return sum;
}

The shard prefix comes first, so the 16 rows of one counter fall into different regions when the table is pre-split on shard boundaries. The trade-off is that a read is no longer atomic across shards: it sums values that may be mid-update. That is fine for view counts and wrong for balances that must never be read inconsistently.

Changing the key on a live table

A monotonic tail is a key-design problem, and the only real fix is a new key: a salt bucket or hash prefix in front of the timestamp, as covered in HBase salting. On a live table, that means a migration, because the old and new keys put the same row in different places. Do it in five steps.

  1. Create the new table pre-split on the salt boundaries, with the same column families and settings.
  2. Dual-write. Change writers to write both tables, the old one first, with a feature flag. From this moment, the new table has every new row.
  3. Backfill. Take a snapshot of the old table and run a Spark or MapReduce job over the snapshot files that rewrites each row with the new key, copying every cell with its original timestamp. That way a row updated after the snapshot keeps its newer dual-written value, because HBase returns the cell with the higher timestamp; a backfill stamped with the current time would overwrite it with stale data. Snapshot reads do not load the RegionServers, which matters because the old table is the one in trouble.
  4. Verify and switch reads. Compare row counts and sampled rows per time range, then move readers to the new table, behind a flag so you can switch back.
  5. Retire. Stop dual-writing, keep the old table for an agreed period, then drop it.

Proving the fix worked

A fix is done when the numbers say so, not when the command returns. Before you act, record a baseline over a fixed interval, for example 15 minutes at the same time of day: each region's read and write request rate (the difference between two samples of the cumulative per-region counters, never the raw totals), each RegionServer's p99 latency, its handler queue length and the client-side error rate.

After the change, take the same measurements over the same interval. Three checks matter. The hot region's share of requests should fall towards its fair share, which is one divided by the number of regions taking that table's traffic. The worst server's p99 should approach its peers'. And the load must not simply have moved: if another region or server is now the top outlier, the fix relocated the hotspot rather than removing it. Keep the before and after figures with the change record; they are what you compare against when the hotspot returns.

Worked example: one tenant melts a RegionServer

A worked example. A multi-tenant table, events, is keyed tenant|timestamp|id, with 120 regions over 12 RegionServers. One RegionServer shows p99 latency of 900 ms against 40 ms elsewhere. Diagnosis finds one region taking 35% of all writes; its keys all belong to tenant 0042, which onboarded a large customer last week. Inside tenant 0042 the keys are monotonic by time.

The first action is containment: a write throttle on the bulk-loading user, which brings that server's p99 down to 300 ms within minutes. Splitting the region at the traffic midpoint would not help, because within the tenant the writes all hit the newest timestamps, so the upper daughter would get everything. The real fix is a key change for that tenant's traffic: the team adds a 16-bucket salt after the tenant prefix, tenant|salt|timestamp|id, so the large tenant's writes spread over 16 regions on different servers, while small tenants keep their locality. They migrate with dual writes and a snapshot backfill over two days. Afterwards, the busiest region takes 4% of writes, every server's p99 is under 60 ms, and the quota is removed.

Failure modes

  • Splitting a monotonic tail. The hotspot moves to the new last region. Check the key shape before splitting.
  • Balancer undoes manual moves. A manually placed daughter moves back on the next balancer run. Switch the balancer off during the operation, or tune its weights so it agrees with you.
  • Split blocked by references. The daughter cannot be split again until a compaction rewrites its reference files.
  • Throttle causes a retry storm. Aggressive client retries multiply load under a quota. Check client backoff before you throttle.
  • Stale replica reads. A client that writes and then reads its own data through a TIMELINE get may not see the write. Use replica reads only where staleness is acceptable.
  • Automatic splitting adds to the noise. Size-based split policies, SteppingSplitPolicy by default in HBase 2.x, split by bytes, not traffic. They are no substitute for a deliberate split point.

Trade-offs

FixTime to reliefReversible?Fixes
Move regionSecondsYesHot server hosting several busy regions
QuotaSecondsYesContainment only
Split at traffic midpointMinutesPartly (merge)Hot range
Cache, coalescing, replicasHours to daysYesHot row reads
Sharded rowsDaysNeeds migrationHot row writes
New key and migrationDays to weeksNeeds migrationMonotonic tail, chronic skew

Related reading: split management for automating splits, and the stochastic load balancer for the cost functions behind placement.

What to do next

  1. Classify the hotspot as single row, hot range or monotonic tail before you touch anything.
  2. Contain first: throttle the noisy user or table with a quota, after checking client backoff.
  3. For a hot range, sample request keys, split at the traffic median and place the daughters on different servers.
  4. For a hot row, add caching or coalescing for reads, region replicas where staleness is acceptable, and sharded rows for writes.
  5. For a monotonic tail, plan a salted-key migration with dual writes and a snapshot backfill.
  6. Afterwards, remove quotas, re-enable the balancer and check the busiest region's share of requests against its fair share.
Key takeaway: Fix an HBase hotspot according to its shape. Split only a hot range, at the traffic midpoint, and place the daughters on different servers. Use caching, replicas or sharded rows for a single hot row, and a salted key with an online migration for a monotonic tail. Quotas buy time while the real fix lands.