A RegionServer is a Java process, and Java's default home for data is the garbage-collected heap. HBase pushes huge volumes of short- and medium-lived bytes through that heap: cached HFile blocks, MemStore cells and RPC buffers. Large heaps full of such data mean long or frequent garbage-collection pauses, and a RegionServer that pauses too long loses its ZooKeeper session and is declared dead. Off-heap memory is HBase's answer: keep the bulk bytes in memory the garbage collector never scans.
This article explains how that works across the whole RegionServer, not just the block cache. It covers what direct memory is, how the off-heap read path moves data from HDFS to the client without copying it onto the heap, how ByteBuffAllocator pools RPC buffers, how the write path can put MemStore chunks off heap, and how to size HBASE_OFFHEAPSIZE so the process neither runs out of direct memory nor out of RAM. It ends with monitoring, failure modes and a checklist. BucketCache internals are covered separately in the BucketCache article.
Direct memory from first principles
Java can allocate memory outside the heap through direct ByteBuffers. The JVM tracks how much it has handed out and refuses to exceed -XX:MaxDirectMemorySize, throwing OutOfMemoryError: Direct buffer memory when the limit is reached. The garbage collector does not move or scan the contents of these buffers; it only tracks the small heap objects that point at them. Ten gigabytes of cached blocks in direct memory cost the collector almost nothing, whereas ten gigabytes of byte arrays on the heap must be traced, and copied or compacted, over their lifetime.
The price is that the application must manage the memory itself. Allocating a direct buffer is slower than allocating on the heap, so HBase pre-allocates and pools them. Freeing happens when a buffer is returned to its pool, which in turn requires knowing when nobody is using it any more, so HBase reference-counts the buffers that back cells in flight. That bookkeeping is the main source of complexity, and the main source of off-heap bugs.
In HBase's start-up scripts, HBASE_OFFHEAPSIZE in hbase-env.sh sets -XX:MaxDirectMemorySize for the RegionServer. Everything off heap that HBase allocates has to fit under that one number.
Where a RegionServer's memory goes
The diagram splits the process into its two budgets. On the heap you keep the things that are small, short-lived or need Java objects: the L1 block cache for index and bloom blocks, the MemStore's cell index, scanners and RPC bookkeeping. Off heap you put the bulk bytes: cached data blocks, RPC and HDFS read buffers and, optionally, MemStore data chunks.
HBase checks the heap side at start-up: the on-heap block cache fraction hfile.block.cache.size plus the global MemStore fraction hbase.regionserver.global.memstore.size may not exceed 0.8 of the heap. Moving data blocks off heap lets you shrink the first fraction, and with it the heap, which is the point: a smaller heap has shorter pauses. The GC tuning article covers what to do with the heap that remains.
The off-heap read path
The off-heap read path was introduced in HBase 2.0, and the reference guide's description of ByteBuffAllocator applies to 2.3 and later. A read that misses the cache follows this flow. The RegionServer reads the HFile block from HDFS. With HBase 2.3+ on a Hadoop version that supports it (the guide names Hadoop 2.10.x and 3.3.x), that read can land directly in an off-heap buffer taken from the allocator. The block is placed in the off-heap BucketCache. Cells are served from the block where it sits, and the RPC layer writes the response from the off-heap buffers to the socket. At no point are the block's bytes copied onto the heap.
A cache hit is the short version of the same path: the block is already in BucketCache, the RegionServer takes a reference on it, builds cells that point into it, writes the response and releases the reference. Because the cells are views into shared buffers, nothing may hold on to them after the RPC finishes. The reference guide is explicit for coprocessor authors: do not keep references to cells outside the scope of the hook method; clone the fields you need. A coprocessor that caches a cell reference can read bytes that have since been reused for another block.
To turn the read path on you configure BucketCache with the off-heap engine. Index and bloom blocks stay in the on-heap L1 cache, and data blocks go to BucketCache, which the block cache article describes as the combined, two-tier arrangement.
ByteBuffAllocator: pooled RPC buffers
Every RPC needs buffers for the request and the response. Allocating those on the heap per call produces garbage proportional to traffic. Since 2.3, ByteBuffAllocator keeps a pool of fixed-size direct buffers instead. The reference guide gives these settings and defaults:
| Property | Default | Meaning |
|---|---|---|
hbase.server.allocator.pool.enabled | true | Use pooled off-heap buffers |
hbase.server.allocator.buffer.size | 66560 (65 KB) | Size of each pooled buffer |
hbase.server.allocator.max.buffer.count | 2 MB x 2 x handler count / 65 KB | Maximum number of pooled buffers |
hbase.server.allocator.minimal.allocate.size | buffer.size / 6 | Requests smaller than this go to the heap instead |
In 2.x the older names hbase.ipc.server.reservoir.enabled, hbase.ipc.server.reservoir.initial.buffer.size and hbase.ipc.server.reservoir.initial.max map to the first three and are deprecated. When the pool is exhausted, the allocator falls back to heap allocation rather than failing, which is safe but brings the garbage back. The RegionServer web UI shows the allocator's statistics, including a heap allocation ratio. The guide's rule is that if heapAllocationRatio is at or above minimal.allocate.size / buffer.size as a percentage (about 16.7% with the defaults), the pool is too small and you should raise max.buffer.count.
Work the default through once. With the default 30 handlers, the pool can hold 2 MB x 2 x 30 / 65 KB, which is about 1,890 buffers of 66,560 bytes, roughly 120 MB. That is small next to a block cache, but if you raise the handler count to 200 the default pool grows to about 800 MB, and that has to be in your direct-memory budget.
The off-heap write path
Writes land in the MemStore. With MSLAB, the MemStore-Local Allocation Buffer, cell data is copied into 2 MB chunks (hbase.hregion.memstore.mslab.chunksize, default 2097152) rather than scattered small arrays, which avoids heap fragmentation. By default those chunks are on the heap. Setting hbase.regionserver.offheap.global.memstore.size to a number of megabytes moves the chunks off heap; the reference guide notes that its default of 0 means MSLAB uses on-heap chunks.
With the chunks off heap, the large cell payloads leave the heap, but the MemStore's index structure and per-cell metadata stay on it and are still governed by the on-heap global MemStore limit. Off-heap MemStore helps most on write-heavy clusters whose heap is dominated by MemStore data.
Worked example: sizing a 128 GB node
Take a RegionServer host with 128 GB of RAM that also runs a DataNode, with a read-heavy workload whose hot set is about 60 GB of blocks, moderate writes and 60 RPC handlers. Plan from the outside in:
| Budget item | Size | Reasoning |
|---|---|---|
| OS, DataNode, agents, page cache floor | about 20 GB | DataNode heap, monitoring and some page cache for HDFS short-circuit reads |
| RegionServer heap (-Xmx) | 24 GB | L1 index and bloom blocks, MemStore metadata, RPC objects |
| BucketCache | 61,440 MB (60 GB) | The hot set, data blocks only |
| Off-heap MemStore | 8,192 MB (8 GB) | Enough for the write rate between flushes |
| Allocator pool | about 250 MB | 2 MB x 2 x 60 / 65 KB = about 3,780 buffers |
| Direct-memory headroom | 2 GB | Netty, HDFS client and other direct buffers |
| Metaspace, thread stacks, code cache | about 2 GB | JVM native overhead outside both budgets |
The reference guide's sizing rule is to set the maximum direct memory a little above the allocator pool plus the off-heap cache, with an extra 1 to 2 GB that worked in its tests. The guide does not include the off-heap MemStore in that sentence; adding it is our own reasoning, because those chunks are direct buffers under the same limit. So: 61,440 + 8,192 + 250 + 2,048 MB is about 71,930 MB, rounded up to 72 GB. The whole process is then 24 + 72 + 2, about 98 GB, leaving about 30 GB for everything else on the host. If that leftover were negative, the fix is to shrink the cache, not to hope.
# hbase-env.sh
export HBASE_HEAPSIZE=24G
export HBASE_OFFHEAPSIZE=72G # becomes -XX:MaxDirectMemorySize=72G
<!-- hbase-site.xml -->
<property><name>hbase.bucketcache.ioengine</name><value>offheap</value></property>
<property><name>hbase.bucketcache.size</name><value>61440</value></property> <!-- MB -->
<property><name>hfile.block.cache.size</name><value>0.2</value></property> <!-- L1: index, bloom -->
<property><name>hbase.regionserver.global.memstore.size</name><value>0.4</value></property>
<property><name>hbase.regionserver.offheap.global.memstore.size</name><value>8192</value></property>
<property><name>hbase.regionserver.handler.count</name><value>60</value></property>
Monitoring direct memory
Watch three things. The JVM's direct buffer pool, exposed as the java.nio:type=BufferPool,name=direct MBean, reports how much direct memory is in use against the limit. The RegionServer UI and metrics report BucketCache occupancy, hit ratio and the allocator's heap allocation ratio. The operating system reports the process's resident set size, which is what the kernel and any container memory limit actually enforce.
# Poll direct-memory usage from the RegionServer's JMX servlet.
import requests
RS = "http://rs1.example.com:16030"
LIMIT = 72 * 2**30 # matches HBASE_OFFHEAPSIZE
def direct_usage():
beans = requests.get(f"{RS}/jmx", params={"qry": "java.nio:type=BufferPool,name=direct"},
timeout=5).json()["beans"]
return beans[0]["MemoryUsed"], beans[0]["Count"]
used, count = direct_usage()
print(f"direct: {used / 2**30:.1f} GiB in {count} buffers ({used / LIMIT:.0%} of limit)")
if used > 0.95 * LIMIT:
print("ALERT: direct memory near MaxDirectMemorySize")For a full native-memory breakdown, start the JVM with -XX:NativeMemoryTracking=summary in a test environment and run jcmd <pid> VM.native_memory summary. It has a small overhead, so it is a diagnostic, not a permanent setting.
Failure modes
- OutOfMemoryError: Direct buffer memory. The budget does not cover cache, pool, MemStore and headroom, often after someone raised the handler count or the cache size without touching
HBASE_OFFHEAPSIZE. Recompute the sum. - Killed by the OOM killer or a container limit. Direct memory fits the JVM limit but the process as a whole exceeds physical RAM or the cgroup limit. Budget heap plus direct plus native overhead against the host, as in the worked example.
- Reference-count leaks. A coprocessor or custom filter keeps cell references beyond the RPC, so buffers are never returned to the pool or, worse, are reused while still referenced. Symptoms are a pool that drains over days or corrupted values. Clone what you need.
- GC pauses come back. The allocator pool is too small and requests fall back to the heap. Check the heap allocation ratio rule and raise the buffer count.
- Cold cache after restart. An off-heap BucketCache is empty after a restart, so latency is high until it warms. Use rolling restarts with region moves, and plan for the warm-up in your latency budget.
- Swap. If the host swaps, off-heap memory is swapped like anything else and latency collapses. Keep swappiness low and leave real headroom.
Trade-offs
Off-heap is not free. It adds configuration, a second memory limit and a class of reference-counting bugs, and a cache hit has to follow references rather than read a plain array. It pays off when the data you want in memory is far larger than a heap you could run with acceptable pauses. For small caches, a modest heap with a modern collector may be simpler; the GC choices article compares G1, ZGC and Shenandoah for that case. If the hot set exceeds RAM, a file-backed BucketCache on fast local SSD extends the cache at lower cost, with higher latency per hit.
What to do next
- Write down your RegionServer's current heap, cache, MemStore, handler count and HBASE_OFFHEAPSIZE, and compute the direct-memory sum against the limit.
- Add the heap, direct and native totals and compare them with the host's RAM minus the DataNode and OS needs.
- Enable the off-heap BucketCache and lower hfile.block.cache.size so the on-heap L1 holds only index and bloom blocks.
- Check the allocator's heap allocation ratio under peak load and raise max.buffer.count if it breaches the guide's threshold.
- If the heap is dominated by MemStore data, trial off-heap MSLAB chunks on one RegionServer and compare GC logs.
- Alert on direct memory at 95% of the limit and on process memory close to the host or container limit.