HBase stores its data in HDFS, and HDFS is a separate system with its own ideas about where bytes live. When a RegionServer writes a file, HDFS puts the first replica of each block on that same machine, so a region that has stayed on one server reads its data from local disk. The moment the region moves to another server, that link is broken: the files have not moved, only the process that serves them. Every block cache miss now crosses the network to a DataNode on another host. That property, the fraction of a region's bytes that sit on the server hosting it, is data locality.
This page explains locality from the HDFS placement rule upward: how it is created, the routine operations that destroy it, how to measure it correctly, what it costs when it is low, and the options for restoring it, from the cheapest (do not lose it in the first place) to the most expensive (rewrite the data). It assumes the RegionServer read path described in the RegionServer deep dive.
How locality is created
The HDFS default placement policy works per block. If the writer is itself running on a DataNode, the first replica goes to that DataNode. The second goes to a node on a different rack, and the third to a different node on the same rack as the second. HBase writes all of its long-lived data through this path: MemStore flushes create new HFiles, and compactions rewrite existing HFiles into new ones. Both happen inside the RegionServer that hosts the region, so every file a region produces starts with one full replica on its own host.
A region that has been served by one server for a while therefore reaches locality close to 1.0 without anyone doing anything. Reads that miss the block cache can then use HDFS short-circuit reads, where the DataNode passes an open file descriptor to the RegionServer over a Unix domain socket and the RegionServer reads the block file directly from local disk, skipping the DataNode's TCP data path. The configuration and its pitfalls are in HDFS short-circuit reads. Without locality, short-circuit reads have nothing to short-circuit.
How locality is lost
Anything that changes which server hosts a region, or which hosts hold its blocks, lowers locality. The common causes, roughly in order of how often they bite:
- Balancer moves. The master's balancer moves regions to even out load. Each move sends the region to a server that typically holds at most one replica of its blocks, and often none.
- Restarts without a region mover. A plain rolling restart closes every region on a server; the master reassigns them across the cluster, and when the server returns it receives a different set. One careless rolling restart can drop cluster locality from near 1.0 to a fraction of that.
- Crashes. When a RegionServer dies, its regions are reopened elsewhere after WAL splitting. This loss is unavoidable; recovery speed is what matters.
- The HDFS balancer and disk replacement. The HDFS balancer moves blocks between DataNodes to even out disk usage and knows nothing about which RegionServer reads them. See the HDFS balancer for how it chooses blocks; on HBase clusters, run it rarely, off-peak, with a low bandwidth cap, and expect a locality drop afterwards.
- Bulk loads. HFiles produced by a MapReduce or Spark job are written by whatever hosts ran the tasks, then adopted by the region. Their first replica is wherever the task ran.
- Node replacement. A new host joins with empty disks. Regions assigned to it start at zero.
Locality heals slowly on its own: new flushes are local, and every minor compaction rewrites some files locally. A region with heavy writes recovers in hours; a region of cold historical data with one large HFile may never recover without help.
Measuring it correctly
HBase computes locality from the HDFS block distribution of each store file: for every host, the number of bytes of the file's blocks that have a replica there. A region's locality on its server is the byte-weighted share of its store file bytes with a replica on that server, from 0.0 to 1.0. Because it is weighted by bytes, one large cold HFile dominates many small fresh ones.
There are three places to read it. The master web UI shows a locality column per RegionServer and per region. The RegionServer exports metrics, including percentFilesLocal, percentFilesLocalPrimaryRegions and percentFilesLocalSecondaryRegions, plus localBytesRead for actual traffic. Programmatically, Admin#getClusterMetrics returns per-region RegionMetrics with getDataLocality(). If you use region replicas, watch the primary-regions metric: secondary replicas read the primary's files from another host, so the HBase metric description itself says their locality is not likely to reach 100 percent.
Two cautions. Locality says where bytes are, not how often they are read: a region at 0.2 locality that is entirely served from block cache costs nothing. Pair locality with block cache hit ratio and localBytesRead before deciding a low number is a problem. And the figures are refreshed periodically, not on every move, so a number taken minutes after a large rebalance can still be stale.
What low locality costs
A local cache miss with short-circuit reads is a disk read inside the same machine. A remote miss is a request to another DataNode, a disk read there, and the block streamed back over TCP through that DataNode process. Each remote read adds network round trips and DataNode thread work, consumes cross-rack bandwidth when the replica is on another rack, and couples your latency to another host's disk queue and garbage collector.
Consider an illustrative RegionServer serving 2 TB of store files with an 85 percent block cache hit ratio. Fifteen percent of reads reach HDFS. At locality 0.95 almost all of them are local. After an unmanaged rolling restart the server's locality is 0.35, so roughly two thirds of the misses are now remote. Median latency often barely moves because the cache absorbs most reads; the tail moves, because the slowest reads are now the ones waiting on a remote DataNode. Measure your own p99 before and after rather than trusting any rule of thumb, including this one.
Restoring locality: the options, cheapest first
1. Do not lose it. Use graceful_stop.sh or the RegionMover tool for every planned restart. They unload a server's regions to others, restart it and load the same regions back, so the regions return to the host that holds their blocks. Disable the balancer during the operation so it does not shuffle regions mid-roll.
# Rolling restart that keeps locality: unload a server, restart it, load the SAME regions back.
bin/graceful_stop.sh --restart --reload rs-host-07.example.com
# The same with the RegionMover tool directly, saving the region list to a file.
hbase org.apache.hadoop.hbase.util.RegionMover -r rs-host-07.example.com -o unload -f /tmp/rs07.regions
# ... restart the RegionServer process ...
hbase org.apache.hadoop.hbase.util.RegionMover -r rs-host-07.example.com -o load -f /tmp/rs07.regions
# Keep the balancer from undoing your work while you operate, then restore it.
echo "balance_switch false" | hbase shell -n
echo "balance_switch true" | hbase shell -n2. Let the balancer value it. The StochasticLoadBalancer scores candidate plans with weighted cost functions; hbase.master.balancer.stochastic.localityCost (default 25) and hbase.master.balancer.stochastic.rackLocalityCost (default 15) are two of them. Raising the locality weight makes the balancer prefer moving a region to a server that already holds its blocks, at the price of a less even region count. The balancer can only choose among existing replicas; it cannot create locality. How the cost functions combine is covered in the StochasticLoadBalancer article.
3. Rewrite the data with major compaction. A major compaction rewrites every HFile of a store into one new file, written locally, so locality returns to near 1.0. It is also the most expensive option: the whole region is read and rewritten, with three-way replication of the output. The setting hbase.hstore.min.locality.to.skip.major.compact defaults to 0 (off). Despite the name, a non-zero value forces work rather than skipping it: when the periodic major-compaction check finds an old store with a single file, it compacts it anyway if that file's locality is below the threshold. That targets exactly the cold regions that never heal by themselves. How compactions are scheduled is in HBase compaction.
4. Move the blocks, not the data. Instead of rewriting a region, move one replica of each non-local block to the hosting server through the HDFS block-move mechanism the HDFS balancer uses. That moves bytes without decompressing, re-encoding or re-replicating them. HubSpot described a production daemon built this way and the upstream work it needed, tracked as HBASE-26250 and HBASE-26304 on the HBase side, so the balancer and metrics see the new locations, and HDFS-15199 for refreshing block locations in open readers. Check your HBase and Hadoop versions for which pieces you have before planning on this route.
The report tool below lists regions below a floor, ranked by remote megabytes rather than by raw locality, so a 50 GB region at 0.5 comes before a 200 MB region at 0.1. With a non-zero budget it queues a few major compactions, which is the safe way to use it from a nightly job.
import java.util.*;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.hbase.*;
import org.apache.hadoop.hbase.client.*;
import org.apache.hadoop.hbase.util.Bytes;
/** Report regions below a locality floor, worst-and-biggest first, and optionally compact a few. */
public class LocalityReport {
public static void main(String[] args) throws Exception {
float floor = Float.parseFloat(args.length > 0 ? args[0] : "0.8");
int compactBudget = Integer.parseInt(args.length > 1 ? args[1] : "0"); // 0 = report only
Configuration conf = HBaseConfiguration.create();
try (Connection conn = ConnectionFactory.createConnection(conf); Admin admin = conn.getAdmin()) {
ClusterMetrics cm = admin.getClusterMetrics(EnumSet.of(ClusterMetrics.Option.LIVE_SERVERS));
List<Object[]> low = new ArrayList<>();
for (Map.Entry<ServerName, ServerMetrics> e : cm.getLiveServerMetrics().entrySet()) {
for (RegionMetrics rm : e.getValue().getRegionMetrics().values()) {
float loc = rm.getDataLocality(); // 0.0 .. 1.0, byte weighted
double mb = rm.getStoreFileSize().get(Size.Unit.MEGABYTE);
if (loc < floor && mb > 0) {
// remote megabytes: what a full cold read of this region would pull over the network
low.add(new Object[] { (1 - loc) * mb, loc, mb, e.getKey(), rm.getRegionName() });
}
}
}
low.sort((a, b) -> Double.compare((double) b[0], (double) a[0]));
for (Object[] r : low) {
System.out.printf("%8.0f MB remote locality=%.2f size=%.0f MB %s %s%n",
r[0], r[1], r[2], r[3], Bytes.toStringBinary((byte[]) r[4]));
}
// Rewriting a region is expensive: cap how many are queued per run, per cluster.
for (int i = 0; i < Math.min(compactBudget, low.size()); i++) {
admin.majorCompactRegion((byte[]) low.get(i)[4]);
}
}
}
}
Configuration in one place
Short-circuit reads must be enabled on both the DataNodes and the RegionServer's HDFS client, with the same domain socket path. The HBase settings below are the locality-related knobs this page discusses; leave the balancer weights at their defaults until you have measured a problem.
<!-- hdfs-site.xml on every DataNode and in the RegionServer's HDFS client config -->
<property>
<name>dfs.client.read.shortcircuit</name>
<value>true</value>
</property>
<property>
<!-- Unix domain socket shared by the DataNode and local clients; the directory must not be world-writable -->
<name>dfs.domain.socket.path</name>
<value>/var/lib/hadoop-hdfs/dn_socket</value>
</property>
<!-- hbase-site.xml -->
<property>
<!-- 0 (default) disables it. With 0.7, a periodic major-compaction check on an old store
holding a single file compacts it anyway if that file's block locality is below 0.7 -->
<name>hbase.hstore.min.locality.to.skip.major.compact</name>
<value>0.7</value>
</property>
<property>
<!-- StochasticLoadBalancer weights; source defaults are 25 and 15 -->
<name>hbase.master.balancer.stochastic.localityCost</name>
<value>25</value>
</property>
<property>
<name>hbase.master.balancer.stochastic.rackLocalityCost</name>
<value>15</value>
</property>
Failure modes and trade-offs
| Situation | What goes wrong | Response |
|---|---|---|
| Compaction storm after a mass restart | many regions major compact at once, disk and network saturate, latency worsens before it improves | rank by remote bytes, cap compactions per run, use off-peak windows |
| Balancer undoes locality work | regions compacted locally are moved away hours later | raise locality weight moderately, or balance less often during recovery |
| Locality-forced compaction too aggressive | a high threshold rewrites cold data every period after each move | start around 0.7 and watch compaction bytes |
| Erasure-coded HFiles | striped blocks spread over many DataNodes, so no single host holds a file | accept remote reads; see HBase on erasure-coded HDFS |
| Object storage instead of HDFS | no locality concept exists | rely on block cache and bucket cache sizing |
| Region replicas | secondary replicas report low locality | judge the primary-regions metric only |
The deeper trade-off is between balance and locality. A perfectly even cluster keeps moving regions and never stays local; a perfectly local cluster never moves regions and develops hot spots. Most clusters do best with infrequent balancing, region movers for all planned work, and targeted compaction for the regions that matter most.
What to do next
- Record current locality per RegionServer from the master UI or
percentFilesLocalPrimaryRegions, alongside block cache hit ratio. - Confirm short-circuit reads are enabled on every DataNode and RegionServer, with matching socket paths.
- Replace every plain restart in your runbooks and automation with
graceful_stop.shorRegionMoverunload and load. - Run the report tool in report-only mode and look at the top twenty regions by remote megabytes.
- If cold single-file regions dominate, set the locality-forced major compaction threshold to about 0.7 and watch compaction volume for a week.
- Queue a small nightly budget of targeted major compactions for the rest.
- Review how and when the HDFS balancer runs on the HBase cluster.
- Re-measure p99 read latency and
localBytesReadafter each change to confirm it paid off.