Every HBase table is cut into regions, contiguous key ranges that are the unit of serving, splitting, balancing and recovery. How many regions you have, and how big each one is, quietly decides whether writes flush cleanly or in storms, how long a RegionServer takes to recover, how evenly load spreads, and how much heap is consumed before any data is stored. Yet most clusters arrive at their region count by accident: whatever the default split policy produced, or whatever someone typed into a pre-split once.
This article treats region count as a capacity decision you make deliberately. It explains what each region costs, derives the two limits that bound a sensible count, works through a sizing example with real arithmetic, shows how to pre-split and measure, and lists the symptoms of too many and too few regions. Split mechanics themselves are covered in HBase regions and splits; here the question is how many you want and why.
What a region really costs
A region looks like a key range, but on a RegionServer it is a set of live objects. Each column family in each region is a store, with its own memstore and its own set of HFiles. The costs scale with stores, not just regions, which is why column-family count appears in every formula below.
| Resource | How region count affects it |
|---|---|
| Memstore heap | Each store taking writes fills a memstore up to the flush size (hbase.hregion.memstore.flush.size, 128 MB by default). All memstores share a global limit of 0.4 of heap. |
| MSLAB chunks | With MSLAB enabled, each memstore that is receiving writes holds at least one 2 MB chunk. A thousand write-active regions with two families pin about 4 GB before storing much. |
| WAL retention | One WAL per server is shared by every region. A WAL file can be archived only after every region with edits in it has flushed, so many slowly written regions keep old WALs alive and trigger forced flushes. |
| HFiles and compaction | Each flush writes one file per store. More stores mean more, smaller files, more compaction work and more open file handles. |
| Block index and Bloom filters | Each HFile carries index and Bloom metadata that must be loaded to read it. |
| Master and meta | Every region is a row in hbase:meta and a unit the master must assign, so restart and failover time grow with count. |
Reads behave differently. The block cache is a separate budget that caches blocks rather than regions, so a cold region costs it nothing; a read-heavy workload is limited by cache size and key distribution, not by region count.
The other side is that a region is also the unit of parallelism. One region is served by one RegionServer at a time; if a table has fewer regions than servers, some servers do nothing for it. And when a server dies, its regions are reassigned and their WAL edits replayed, so fewer, larger regions mean each reassignment carries more work.
Two limits: the write bound and the data bound
Two limits bound a healthy region count per RegionServer. They measure different things, and confusing them is the root of most bad advice on this subject.
The write bound asks how many stores can be write-active before the memstore budget runs out. The HBase reference guide gives it as heap size times the global memstore fraction, divided by flush size times the number of column families. With a 16 GB heap, 0.4, 128 MB and one family that is about 50 regions. Beyond that, memstores never reach their flush size; the server hits its global limit and flushes the largest memstores early, producing small files and extra compaction. Note what it counts: regions receiving writes. Regions that are read-only or cold cost almost no memstore.
The data bound asks how many regions the data needs: the compressed size of the table's store files on this server divided by the region size you are targeting. It counts every region, hot or cold.
If the data bound is below the write bound, you are fine: every region could be written at once. If the data bound is far above it, check the write pattern. A time-series table keyed by time writes only to its newest regions, so thousands of cold regions are harmless. A table with uniformly random keys writes to every region at once, so the write bound is a hard limit and you must use bigger regions, fewer families or more servers.
Region size is governed by hbase.hregion.max.filesize, 10 GB by default: when a store grows past the threshold the region splits. In HBase 2.x the default split policy is SteppingSplitPolicy, which splits a table's first region on a server early, at twice the flush size, and then uses the full threshold. That gets a new table spread across a few servers quickly without fragmenting it. If your sizing says regions should be 20 GB, raise the threshold on that table rather than cluster-wide.
Worked example: a 20-server events cluster
Take a cluster of 20 RegionServers with 32 GB heaps and an events table with two column families: d for payload and m for small metadata. All figures are illustrative; plug in your own.
Write bound. 32 GB x 0.4 = 12.8 GB of memstore. Divided by 128 MB x 2 families = 256 MB per region, that allows 50 regions per server to be write-active at full flush size. In practice the server starts forcing flushes at the lower limit, 0.95 of the global limit, and the m family barely fills, so think of 40 to 50 as the realistic ceiling.
Data bound. The table holds 12 TB of compressed store files today and grows by 400 GB a month. HDFS replication does not count here; region size is the size of one copy. That is 600 GB per server. At 10 GB regions it needs 60 regions per server, 1,200 in total; at 20 GB, 30 per server.
Reconcile. The keys are a salted hash of device ID followed by a timestamp, with 40 salt buckets, so writes land on the newest region of every bucket: about 40 write-active regions in total, two per server. The write bound is nowhere near binding, so the data bound decides. Choosing 10 GB gives 60 regions per server, which is a manageable count, keeps compactions short, and leaves room for about 20 months of growth before passing about 100 per server.
Pre-split. Because writes are salted into 40 buckets, pre-split the new table into exactly 40 regions on bucket boundaries. Every server takes writes from day one, and natural splits subdivide each bucket as it grows.
Now change one fact: if the keys had been random UUIDs, all 1,200 regions would take writes. At 60 per server against a write bound of about 45, memstores would flush early and continuously. The fix is to use 20 GB regions (30 per server), drop to one family, or add servers, and the arithmetic tells you which costs least.
Pre-splitting and measuring
Pre-splitting is done when the table is created. The shell accepts explicit split keys or a generated layout:
# 40 regions on salt-bucket boundaries (keys start with a two-digit bucket)
create 'events', {NAME => 'd', COMPRESSION => 'ZSTD'}, {NAME => 'm'},
SPLITS => (1..39).map { |i| format('%02d', i) }
# hex-prefixed keys, 32 regions spread evenly over the hex space
create 'sessions', 'd', {NUMREGIONS => 32, SPLITALGO => 'HexStringSplit'}
# a per-table split threshold instead of the 10 GB cluster default
alter 'events', MAX_FILESIZE => '21474836480'To check what you actually have, read region counts and sizes from cluster metrics instead of the UI. This uses the HBase 2.x Admin API:
try (Connection conn = ConnectionFactory.createConnection(conf);
Admin admin = conn.getAdmin()) {
ClusterMetrics cm = admin.getClusterMetrics();
for (Map.Entry<ServerName, ServerMetrics> e : cm.getLiveServerMetrics().entrySet()) {
long regions = 0, mb = 0, memMb = 0;
for (RegionMetrics r : e.getValue().getRegionMetrics().values()) {
regions++;
mb += (long) r.getStoreFileSize().get(Size.Unit.MEGABYTE);
memMb += (long) r.getMemStoreSize().get(Size.Unit.MEGABYTE);
}
System.out.printf("%s regions=%d storeGB=%.1f avgGB=%.2f memstoreMB=%d%n",
e.getKey().getServerName(), regions, mb / 1024.0,
regions == 0 ? 0 : mb / 1024.0 / regions, memMb);
}
}Run it weekly and keep the output. Three numbers matter: regions per server against your plan, average region size against the threshold, and the spread between the busiest and quietest server.
Symptoms of getting it wrong
Too many regions shows up as a server that flushes constantly with small files, logs about forced flushes because of the global memstore limit or too many WAL files, and spends its disk bandwidth on compaction. Flushes stall when a store reaches hbase.hstore.blockingStoreFiles (16 in HBase 2.x) until compaction catches up, and writes to a region block once its memstore reaches four times flush size (hbase.hregion.memstore.block.multiplier). Restarts take noticeably longer because the master has thousands of regions to assign and open.
Too few regions shows up as uneven load: one server is hot because it holds a third of a busy table, and the balancer cannot help because it moves whole regions. Major compactions on very large regions run for hours and rewrite huge files. A server crash takes longer to recover because each region's WAL replay and reopen is bigger.
Recovery time is the hidden cost. When a RegionServer dies, its WAL is split by region (in HBase 2.x the work is handed to surviving servers), the master reassigns every region it held, and each new host replays that region's unflushed edits before opening it. Store files live on HDFS and are not copied, so a region's size on disk barely affects how fast it reopens; what matters is how many regions must be opened and how many edits each one has to replay. Sixty regions spread over nineteen survivors is about three each, opened in parallel, which is fast. Thousands of write-active regions that flush rarely mean many opens and long replays, which is one more reason to keep the write-active count within what the memstore can flush promptly.
The wrong regions, a sensible count with bad keys, looks like too few regions: all writes hit the last region of a monotonically increasing key. No count fixes that; salt or hash the key prefix.
Keeping the count right over time
Region count is not something you set once. Splits raise it automatically as data grows; the region normalizer can split oversized regions and merge undersized ones toward a target; the balancer spreads regions across servers but does not change how many there are. Decide the target size per table, let splits and the normalizer hold it there, and let the balancer handle placement.
| Choice | Gain | Cost |
|---|---|---|
| Smaller regions (1-5 GB) | Finer balancing, short compactions, fast reassignment | More stores, more files, more meta rows, write bound reached sooner |
| Default (10 GB) | Reasonable for most tables | May be small for very large, append-only tables |
| Larger regions (20-50 GB) | Fewer stores, fewer flushes, lower heap overhead | Long major compactions, coarse balancing, slower recovery per region |
| More column families | Separate tuning per family | Multiplies every per-store cost; flushes are per region |
| Pre-split | Even load from the first write | Wrong split points create empty or hot regions |
Memstore tuning interacts with all of this; memstore flushes explains the flush triggers, and the RegionServer explains where heap goes.
What to do next
- For each large table, record compressed size, growth rate, column-family count and write pattern (time-ordered, salted or random).
- Compute the write bound per server from heap, global memstore fraction, flush size and families.
- Compute the data bound from size per server and a target region size, starting at 10 GB.
- If random writes make the data bound exceed the write bound, raise region size for that table, merge families or add servers.
- Pre-split new tables on real key boundaries, such as one region per salt bucket.
- Set MAX_FILESIZE per table where the cluster default is wrong for it.
- Run the metrics script weekly; alert when regions per server drift past plan or the busiest server holds far more than the average.
- Watch flush and compaction logs for forced flushes and blocking-store-file warnings, which are the first sign of too many write-active regions.