An HBase table is a sorted map from row key to cells, cut into contiguous key ranges called regions. A RegionServer is a process that hosts some of those regions, and almost everything that goes right or wrong in an HBase cluster is a story about regions: where they live, how big they are, how they move, and what happens to their unflushed edits when the process hosting them dies.
This article looks at the RegionServer architecture through that lens. The companion piece on RegionServer internals covers handlers, heap budgets and the read path inside the process. Here the subject is the region itself: its layout on HDFS, the sequence ids tying it to the shared write-ahead log, the state machine the master drives it through, splits, crash recovery, and how many regions a server should carry.
What a region is on disk
A region is described by a table name, a start key (inclusive), an end key (exclusive), a region id (usually its creation timestamp) and an encoded name, which is a hash of the rest. The full region name looks like profile,user-5000,1727900000000.3f9c2a..., and the encoded part is what you see in logs and directory names. The first region of a table has an empty start key and the last has an empty end key, so the regions of a table always tile the whole key space with no gaps.
On HDFS each region is a directory under the table: /hbase/data/<namespace>/<table>/<encoded-name>/. Inside it there is a .regioninfo file describing the region, one subdirectory per column family, and, after a crash, a recovered.edits directory. A column family inside a region is called a store. Each store has its own MemStore in memory and its own set of immutable HFiles on disk, which is why adding column families multiplies flushes and files: a region with three families is three stores that flush and compact somewhat independently.
Nothing about the region lives only in the RegionServer: placement and state are in hbase:meta and the data is on HDFS, so any other RegionServer can rebuild the open region from there.
What a region is on disk
A region is described by a table name, a start key (inclusive), an end key (exclusive), a region id (usually its creation timestamp) and an encoded name, which is a hash of the rest. The full region name looks like profile,user-5000,1727900000000.3f9c2a..., and the encoded part is what you see in logs and directory names. The first region of a table has an empty start key and the last has an empty end key, so the regions of a table always tile the whole key space with no gaps.
On HDFS each region is a directory under the table: /hbase/data/<namespace>/<table>/<encoded-name>/. Inside it there is a .regioninfo file describing the region, one subdirectory per column family, and, after a crash, a recovered.edits directory. A column family inside a region is called a store. Each store has its own MemStore in memory and its own set of immutable HFiles on disk, which is why adding column families multiplies flushes and files: a region with three families is three stores that flush and compact somewhat independently.
Nothing about the region lives only in the RegionServer: placement and state are in hbase:meta and the data is on HDFS, so any other RegionServer can rebuild the open region from there.
How a RegionServer hosts many regions
A RegionServer keeps a map from encoded region name to open region. Regions on the same server do not share MemStores or HFiles, but they do share three things, and those shared resources explain most of the behaviour you see under load.
- The write-ahead log. By default every region on a server appends to the same WAL, so a single sync makes edits from many regions durable at once. Group commit is efficient, but it couples regions: a slow HDFS pipeline delays every write on the server, and the WAL can only be archived when every region that wrote into it has flushed past those edits.
- The global MemStore budget.
hbase.regionserver.global.memstore.size(0.4 of heap by default) caps the sum of all MemStores. When the server nears it, flushes are forced on the largest regions and, past the limit, writes block. Many small regions therefore each get a small share of a fixed budget and flush small files, which turns into extra compaction work. - The BlockCache. Read caching is server-wide, so a scan-heavy region can evict the blocks a latency-sensitive neighbour depends on.
A region is the unit of placement and recovery; the RegionServer is the unit of resource contention.
Sequence ids tie the WAL to the HFiles
Every edit applied to a region receives a monotonically increasing sequence id, written into the WAL entry alongside the edit. When a store flushes, the resulting HFile records the highest sequence id it contains. Two consequences follow, and both are central to the architecture.
First, WAL retention. The server tracks, for every region, the oldest sequence id still sitting unflushed in a MemStore. A WAL file can be archived only when every edit in it is older than all of those markers. If one rarely written region never fills its MemStore, it pins every WAL file since its oldest edit. hbase.regionserver.maxlogs bounds the count: when it is exceeded the server forces flushes on the regions holding the oldest edits, which is why a cluster with thousands of cold regions sees bursts of tiny flushes.
Second, recovery. When a region is replayed after a crash, any edit whose sequence id is at or below the maximum already recorded in that store's HFiles is skipped. Replay is therefore idempotent, and only the genuinely unflushed tail is re-applied.
The region state machine
The master's AssignmentManager owns region placement. Since HBase 2, assignment runs on the procedure framework: each move is a durable state machine whose progress is persisted, so a master that fails over mid-move resumes rather than leaving the region half-open. In current releases the procedure that drives open, close and move is TransitRegionStateProcedure. The states a region passes through are visible in the master UI and in hbase:meta:
| State | Meaning | What to do if it sticks |
|---|---|---|
| OFFLINE / CLOSED | Not served anywhere | Normal for disabled tables; otherwise assign it |
| OPENING | A server has been told to open it | Check that server log for store open errors |
| OPEN | Serving reads and writes | Nothing |
| CLOSING | Flushing and releasing on the old server | Look for a slow flush or HDFS trouble |
| SPLITTING / MERGING | Parent being replaced by daughters, or the reverse | Inspect the split or merge procedure |
| FAILED_OPEN | Open retried and gave up | Fix the cause (corrupt file, missing class), then assign |
A region in any state other than OPEN or CLOSED is in transition, and a region in transition for minutes is the single most common symptom of a sick cluster. The shell gives you the levers, and HBCK2 exists for the cases where the procedure state itself is wrong:
# shell: inspect and move regions
list_regions 'profile' # regions, servers, sizes, request counts
move '3f9c2a...', 'rs3.example.com,16020,1727900000000'
unassign '3f9c2a...'
assign '3f9c2a...'
balance_switch false # freeze the balancer while you investigate
# HBCK2 (separate jar) for stuck procedures, used with care
hbase hbck -j hbase-hbck2.jar assigns 3f9c2a...
hbase hbck -j hbase-hbck2.jar bypass -o <pid> # last resort
Opening and closing a region
To open a region a RegionServer reads .regioninfo, opens every store, reads the trailer, index and Bloom filter metadata of every HFile, replays any recovered.edits into the MemStore, flushes them if needed, and only then reports OPEN to the master, which records the new location in meta. A region with forty HFiles in one store takes visibly longer to open than one with four, so a compaction backlog does not only hurt reads: it lengthens every move, every rolling restart and every crash recovery.
Closing is the mirror image: stop writes, flush so the next owner needs no WAL, release handles. A close that cannot flush is what leaves a region stuck in CLOSING.
Splits and merges
A split cuts one region into two daughters at a midpoint key, usually chosen from the largest store. It is cheap because no data is rewritten: each daughter gets reference files that point at the top or bottom half of each parent HFile. The parent is marked offline in meta, the daughters open, and later compactions rewrite the references into real files; a chore, the CatalogJanitor, then removes the parent.
When to split is a policy. The default in HBase 2 is SteppingSplitPolicy: the first region of a table on a server splits at twice the flush size, and after that regions split at hbase.hregion.max.filesize (10 GB by default). Merges are the inverse, combining adjacent regions, and the region normalizer can split and merge automatically to keep sizes even. For predictable workloads, pre-split tables at creation so writes spread from the first minute. The details of policies and split points are in regions and splits.
Recovery when a RegionServer dies
A RegionServer holds a ZooKeeper session, and its ephemeral node disappears when the session expires (zookeeper.session.timeout). A long garbage-collection pause looks exactly like death. The master then runs a crash procedure for that server, in this order:
- If the dead server hosted
hbase:meta, recover and reassign meta first, because nothing else can be located without it. - Split the dead server's WAL files. Each WAL interleaves edits from many regions, so splitting reads it and writes per-region
recovered.editsfiles. The work is distributed to surviving RegionServers, not done by the master; it is coordinated either through ZooKeeper or, withhbase.split.wal.zk.coordinated=false, through the procedure framework (oneSplitWALProcedureper WAL), which is the default in HBase 3. - Assign every region of the dead server elsewhere. Each new host replays that region's recovered edits during open, skipping anything already in HFiles by sequence id, then reports OPEN.
Mean time to recovery is therefore detection time plus WAL split time plus open time. You shorten detection by keeping GC pauses small enough to allow a shorter session timeout, WAL split time by flushing regularly so WALs stay small, and open time by keeping HFile counts low. The splitting mechanics are covered in WAL splitting.
Worked example: four regions and one crash
Take a profile table pre-split into four regions on a cluster of two RegionServers with 16 GB heaps and one column family. rs1 hosts regions 1 and 2, rs2 hosts 3 and 4. A client writes user-7712; its cached location says region 3 on rs2. rs2 assigns sequence id 1408 for region 3, appends the edit to the shared WAL, syncs, inserts into region 3's MemStore and acknowledges.
Region 3 is hot and flushes at sequence id 9,000, writing an HFile with that max sequence id. Region 4 receives one write an hour, so its oldest unflushed edit is sequence id 96 from yesterday, and it pins every WAL file rs2 has written since. When the WAL count passes hbase.regionserver.maxlogs, rs2 flushes region 4 to release them: one tiny file.
Now rs2 stops responding for 40 seconds because of a full GC, and its session expires. The master starts the crash procedure. rs1 splits rs2's WALs into recovered.edits for regions 3 and 4, then opens both. Region 3's replay skips everything at or below 9,000 and re-applies only the 600 edits written after the flush. Clients see NotServingRegionException or connection errors, refresh their location cache from meta, and retry against rs1. When rs2 wakes up it finds its session gone and aborts.
Watching regions from code
The Admin API exposes per-region metrics, which is how you find the regions behind most problems: the one with too many store files, the one with an enormous MemStore, or the one taking most of the requests.
import org.apache.hadoop.hbase.*;
import org.apache.hadoop.hbase.client.*;
import java.util.Map;
try (Connection conn = ConnectionFactory.createConnection(HBaseConfiguration.create());
Admin admin = conn.getAdmin()) {
ClusterMetrics cluster = admin.getClusterMetrics();
for (Map.Entry<ServerName, ServerMetrics> e : cluster.getLiveServerMetrics().entrySet()) {
int regions = e.getValue().getRegionMetrics().size();
System.out.printf("%s hosts %d regions%n", e.getKey().getServerName(), regions);
for (RegionMetrics r : e.getValue().getRegionMetrics().values()) {
if (r.getStoreFileCount() > 20) {
System.out.printf(" %s files=%d memstore=%s reads=%d writes=%d%n",
r.getNameAsString(), r.getStoreFileCount(), r.getMemStoreSize(),
r.getReadRequestCount(), r.getWriteRequestCount());
}
}
}
}Clients find meta through ZooKeeper by default in HBase 2.x; HBase 3.0 defaults to an RPC-based registry that bootstraps from masters or RegionServers. Either way locations are cached until a miss or error. The structure of that catalog is described in the hbase:meta table.
How many regions per server
The HBase reference guide gives a rule of thumb for the number of actively written regions a server can carry: heap times the MemStore fraction, divided by flush size times the number of column families. With a 16 GB heap, 0.4 and a 128 MB flush size on one family, that is 16384 x 0.4 / 128, about 51 regions. Beyond that, each region's share of the budget drops below a full flush and files get smaller. Cold, read-mostly regions are cheaper, so real servers often host more, but the formula marks where write-heavy tables start paying.
Where regions go is decided by the balancer; see the stochastic load balancer.
Failure modes
- Regions stuck in transition. Usually a store that cannot open (corrupt HFile, a missing coprocessor class) or a close that cannot flush. Read the target server log before touching HBCK2.
- Hot region. Monotonic row keys send all writes to the last region; splitting only moves the hot spot. Fix the key design with salting or hashing.
- Too many regions per server. Tiny flushes, compaction storms, long WAL retention and slow restarts.
- Slow recovery. Large WALs from infrequent flushes, plus high HFile counts, stretch the crash procedure from seconds to many minutes.
- Meta unavailable. Every region lookup fails until meta is reassigned.
Trade-offs
| Choice | Gain | Cost |
|---|---|---|
| Large regions (10 to 20 GB) | Fewer files and meta rows, fewer WAL pins | Slow moves, coarse balance |
| Small regions | Even load, parallel recovery | More flushes, more compaction, more overhead |
| Pre-splitting | Writes spread from day one | Needs a key distribution you can predict |
| Short ZooKeeper session timeout | Faster failure detection | GC pauses become false deaths |
| Region replicas | Reads continue through a crash | Extra memory, possibly stale reads |
What to do next
- Run
list_regionson your busiest table and note size, server and request counts per region. - Apply the region count formula to your heap and flush size and compare it with what each server hosts.
- Find regions with more than twenty store files and check whether compaction is keeping up.
- Measure your WAL count per server against
hbase.regionserver.maxlogsand look for cold regions pinning logs. - Rehearse a RegionServer kill in staging and time detection, WAL split and region open separately.
- Pre-split any table you are about to bulk load, and decide whether the normalizer should run on it.
- Alert on regions in transition for longer than a few minutes, and keep HBCK2 ready before you need it.