HBase splits a region into two when it grows too large. That sentence is true, and running a busy cluster on it alone is how teams end up with split storms after a bulk load, a brand-new table pinned to one RegionServer during its first hour, and regions that keep splitting but never take load off the server that is hot. Split management is the practice of deciding deliberately when regions split, where the boundaries fall and which component is allowed to trigger a split.
This article does not re-explain the split procedure itself (reference files, daughter compaction, parent cleanup) or catalogue every policy. Separate articles cover both and are linked below. It covers the decisions an operator makes: computing boundaries from real keys, when to take splitting away from the RegionServers, how to run a budgeted splitter, how to keep it from fighting the normalizer and the balancer, and what to do when splits go wrong.
What split management decides
Every split answers three questions. When: what size or load triggers it. Where: which row key becomes the boundary. Who: whether the RegionServer's policy, an operator, the normalizer or an external tool starts it.
Out of the box, the RegionServer answers all three. The policy decides when. The split point is the middle key of the region's largest store, which is the middle of the data by size, not by traffic. The trigger comes from the server after a flush or compaction. That is a good default for tables whose keys are well spread and whose growth is steady. It is a poor fit for three common cases: a new table that will receive heavy writes from minute one, a bulk load that lands many gigabytes at once, and a table whose hot keys are a narrow range that a size-based midpoint does not isolate.
The default policy, in numbers
Two numbers drive the default. The memstore flush size defaults to 128 MB, and the split policies derive an initial size of twice that, 256 MB. The maximum region size, hbase.hregion.max.filesize, defaults to 10 GB.
| Policy | Split size for a store | Effect |
|---|---|---|
SteppingSplitPolicy (default) | 256 MB if this is the only region of the table on this server, otherwise the max file size | A new table spreads quickly, then regions grow to 10 GB |
IncreasingToUpperBoundRegionSplitPolicy | min(max file size, 256 MB x n3), n = regions of this table on this server | 256 MB, 2 GB, 6.75 GB, then capped at 10 GB |
ConstantSizeRegionSplitPolicy | The max file size | Predictable, slow to spread a new table |
DisabledRegionSplitPolicy | Never | Splits only when someone asks |
The stepping rule has a property worth knowing. A new table splits at 256 MB only while each server holds exactly one of its regions. Once the daughters land on the same server, the next split waits for 10 GB. A table that starts on one server can therefore stall at two or three regions, all on that server, and absorb its first day of writes there. That is the most common reason to pre-split.
There is also a per-server ceiling, hbase.regionserver.regionSplitLimit, default 1000. Past it, the server stops requesting splits. That is a safety limit, not a target, and a cluster near it has a region-count problem to fix rather than a setting to raise.
The management loop
Figure 1 shows the shape of a managed setup. Day-one boundaries come from data. Later splits come from one place, whether that is the RegionServer policy, the normalizer or your own tool, and that place has a budget. The balancer moves the daughters afterwards. Every component in the figure exists in a default cluster too. The difference is that each has a known job and nothing triggers splits behind the others' backs.
Pre-split points from real keys
Generic pre-splitting with HexStringSplit or UniformSplit (via RegionSplitter -c 32 -f cf or the shell's NUMREGIONS/SPLITALGO) assumes keys are uniform over a known alphabet. That holds for hashed or salted keys. It is wrong for natural keys such as tenant IDs, URLs or account numbers. Their distribution is lumpy, and evenly spaced boundaries put most data into a few regions.
The fix is to compute boundaries from a sample of the keys you will actually write. Take a uniform random sample (from a source-system export, a log, or an existing table scanned with a key-only filter), sort it, and take quantiles. If the sample is weighted by expected write volume rather than row count, the boundaries balance writes rather than bytes.
import random
def split_points(sample_keys, regions, weights=None):
"""Return regions-1 boundaries so each region gets ~equal weight of the sample."""
pairs = sorted(zip(sample_keys, weights or [1] * len(sample_keys)))
total = sum(w for _, w in pairs)
step, acc, out = total / regions, 0.0, []
for key, w in pairs:
acc += w
if acc >= step * (len(out) + 1) and len(out) < regions - 1:
if not out or key > out[-1]: # boundaries must be strictly increasing
out.append(key)
return out
keys = [line.strip() for line in open("tenant_keys_sample.txt")]
sample = random.sample(keys, min(200_000, len(keys)))
with open("splits.txt", "w") as f:
for k in split_points(sample, regions=48):
f.write(k + "\n")hbase> create 'events', {NAME => 'd', COMPRESSION => 'ZSTD'}, SPLITS_FILE => 'splits.txt'Size the region count so the table spreads across every RegionServer at least twice, and so that each region's expected size after a few months stays well under the split threshold. A table pre-split into 48 regions on 12 servers starts with four regions per server and lets the balancer place them before the first write.
Taking control of automatic splits
Taking automatic splitting away from the RegionServers is reasonable when you will replace it with something better. It is dangerous when you simply switch it off and forget. There are two levers.
- The cluster switch.
splitormerge_switch 'SPLIT', falsein the shell, oradmin.splitSwitch(false, true)in Java, stops all splits on all tables. Use it for short windows, such as a bulk load, an upgrade or an incident, and record that you did. A cluster left in this state grows regions without bound. - The table setting. The table descriptor's
SPLIT_ENABLEDflag (setSplitEnabled(false)onTableDescriptorBuilder), or theDisabledRegionSplitPolicy, stops splits for one table. This is the right lever for a table whose splits a tool will manage permanently.
TableDescriptor td = admin.getDescriptor(TableName.valueOf("events"));
admin.modifyTable(TableDescriptorBuilder.newBuilder(td)
.setSplitEnabled(false) // the managed splitter owns this table's splits
.build());With automatic splits off, a region that keeps growing still slows down. Its compactions take longer, and recovery after a server crash takes longer because more data must be reopened. The owner of the splits must therefore keep up, and an alert on maximum region size is the safety net if it stops.
A managed splitter under a budget
A managed splitter is a small scheduled job. It reads region metrics, picks candidates and asks the master to split a bounded number of them per run. The budget is the point. Each split leaves two daughters that must compact away their reference files before they can split again, and that compaction is real I/O. Ten splits at once on one server is a self-inflicted incident.
long thresholdMb = 6 * 1024; // split before the 10 GB ceiling, on our schedule
int budget = 4; // splits per run, cluster-wide
List<RegionMetrics> candidates = new ArrayList<>();
for (ServerName sn : admin.getClusterMetrics().getLiveServerMetrics().keySet()) {
for (RegionMetrics rm : admin.getRegionMetrics(sn, table)) {
if (rm.getStoreFileSize().get(Size.Unit.MEGABYTE) > thresholdMb) candidates.add(rm);
}
}
candidates.sort(Comparator.comparingDouble(
(RegionMetrics rm) -> rm.getStoreFileSize().get(Size.Unit.MEGABYTE)).reversed());
for (RegionMetrics rm : candidates.subList(0, Math.min(budget, candidates.size()))) {
try {
admin.splitRegionAsync(rm.getRegionName()).get(10, TimeUnit.MINUTES); // master picks midkey
} catch (Exception e) {
log.warn("split skipped for {}: {}", rm.getNameAsString(), e.toString()); // e.g. references remain
}
}Run it every 15 to 30 minutes, outside peak hours if your traffic has peaks. To split on a traffic boundary rather than a size midpoint, pass an explicit split key to splitRegionAsync(regionName, splitPoint). The key might come from request sampling that shows where the hot range starts. The RegionSplitter -r rolling split, with -o limiting outstanding splits, is a ready-made version of the same idea for doubling a whole table.
Coordinating with the normalizer, balancer and compactions
Three other components act on the same regions, and each can undo the others' work.
- The normalizer plans splits of oversized regions and merges of undersized ones, aiming at an average or at
NORMALIZER_TARGET_REGION_SIZE_MBorNORMALIZER_TARGET_REGION_COUNTon the table. If your splitter uses a 6 GB threshold and the normalizer's target makes 3 GB regions look "small", they will split and merge the same ranges in turn. Either give the normalizer a target consistent with your threshold or setNORMALIZATION_ENABLEDfalse on tables the splitter owns. Note that the normalizer is also gated by the split and merge switches. - The balancer moves daughters to even out load. Splitting many regions on one server and then letting the balancer move half of them is normal. Running the balancer during a burst of splits adds region moves on top of compactions. Pause it for planned split windows and turn it back on afterwards.
- Compactions are what make a split finish. A daughter cannot split again until its references are compacted away. Throttled or queued compactions therefore slow every split. Check the compaction queue before you raise the splitter's budget.
Signals to watch
| Signal | Where | What it tells you |
|---|---|---|
splitQueueLength | RegionServer metrics | Splits requested but not started; a persistent queue means splits are blocked or too frequent |
splitRequestCount | RegionServer metrics | Split rate; a spike after a bulk load or schema change is a storm |
| Largest region size per table | Region metrics | Whether the split owner is keeping up; alert well before the max file size |
| Regions per server | Master UI or cluster metrics | Distance from the 1000 split limit and from your heap budget |
| Compaction queue | RegionServer metrics | Whether daughters can shed references and split again |
| Procedures in progress | list_procedures in the shell | Split procedures that are stuck rather than slow |
Failure modes and runbooks
Split storm after a bulk load. Many regions cross the threshold together, split together and compact together, and latency spikes for an hour. Prevention: pre-split the target table to the bulk load's key distribution, or switch splits off for the load and let the managed splitter catch up under its budget.
Splits that do not relieve a hotspot. A monotonically increasing key (a timestamp or sequence prefix) sends every write to the last region. Splitting it creates a new last region that receives every write. No split policy fixes this. The fix is in the key: salting or hashing a prefix, then pre-splitting on the salt buckets.
Region that will not split. The usual cause is reference files still present from an earlier split, because compaction has not caught up. Major-compact the region with admin.majorCompactRegion(regionName), wait, then retry. Also check that splits are not disabled at the cluster or table level by an earlier change.
Stuck split procedure. A split that shows as running in list_procedures for far longer than its peers, often with a region in transition, needs the master's logs first. Most resolve once the underlying problem, such as an unavailable server or HDFS trouble, is fixed. Bypassing a procedure with HBCK2 is a last resort, done with the HBCK2 documentation open, because forcing state can leave the parent and daughters inconsistent in hbase:meta.
Too many small regions. Aggressive pre-splitting or a low threshold leaves thousands of tiny regions, each with memstore overhead and open file handles. Merge neighbours with mergeRegionsAsync or let the normalizer do it, and raise the threshold.
Trade-offs
Leaving splits to the default policy costs nothing to run and suits tables with spread-out keys and steady growth. Pre-splitting from a sample costs one script and removes the first-day hotspot, but the boundaries are only as good as the sample. A managed splitter gives predictable timing and traffic-aware boundaries, at the cost of a job you must keep running and monitor. A forgotten splitter is worse than none. The usual answer is all three in sequence: pre-split every heavy table, keep the default policy for ordinary tables, and move a table to managed splitting only when its load justifies the extra job.
For more depth, see the split procedure internals, regions and split policies, the region normalizer, region count sizing and salting row keys.
What to do next
- List your top tables by write volume and record each one's region count, largest region and regions per server.
- For any table that will take heavy writes from day one, sample its keys and create it with quantile split points instead of evenly spaced ones.
- Check whether the cluster split switch is on; if anyone turned it off, find out why and when.
- Add alerts on largest region size per table and on
splitQueueLength. - Make the normalizer's target consistent with your split threshold, or disable normalization on tables you split manually.
- Before the next bulk load, pre-split the target or plan a split-off window with a budgeted catch-up.
- If a hot range survives splitting, fix the row key rather than the split policy.