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.

Advertisement

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

Before and after a region move: the files stay, the region leavesHost A: RegionServer + DataNodeserves region R, holds replica 1 of every blockHost B: DataNodereplica 2 (other rack)Host C: DataNodereplica 3 (same rack as B)locality of R on A = 1.0reads use short-circuit: local disk, no socketbalancer, restart or crash moves RHost D: RegionServer + DataNodenow serves R, holds none of its blocksHosts A, B, Cstill hold all three replicasremote readlocality of R on D = 0.0every cache miss crosses the network to a DataNodeHealing on Dnew flushes and compactions write replica 1 locally; major compaction or block moves restore the restLocality is a property of the pair (region files, hosting server), not of the data alone
A region moves from host A to host D. Its HFiles keep all three replicas on A, B and C, so D reads everything remotely until data is rewritten or moved.

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.

Advertisement

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 -n

2. 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

SituationWhat goes wrongResponse
Compaction storm after a mass restartmany regions major compact at once, disk and network saturate, latency worsens before it improvesrank by remote bytes, cap compactions per run, use off-peak windows
Balancer undoes locality workregions compacted locally are moved away hours laterraise locality weight moderately, or balance less often during recovery
Locality-forced compaction too aggressivea high threshold rewrites cold data every period after each movestart around 0.7 and watch compaction bytes
Erasure-coded HFilesstriped blocks spread over many DataNodes, so no single host holds a fileaccept remote reads; see HBase on erasure-coded HDFS
Object storage instead of HDFSno locality concept existsrely on block cache and bucket cache sizing
Region replicassecondary replicas report low localityjudge 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

  1. Record current locality per RegionServer from the master UI or percentFilesLocalPrimaryRegions, alongside block cache hit ratio.
  2. Confirm short-circuit reads are enabled on every DataNode and RegionServer, with matching socket paths.
  3. Replace every plain restart in your runbooks and automation with graceful_stop.sh or RegionMover unload and load.
  4. Run the report tool in report-only mode and look at the top twenty regions by remote megabytes.
  5. If cold single-file regions dominate, set the locality-forced major compaction threshold to about 0.7 and watch compaction volume for a week.
  6. Queue a small nightly budget of targeted major compactions for the rest.
  7. Review how and when the HDFS balancer runs on the HBase cluster.
  8. Re-measure p99 read latency and localBytesRead after each change to confirm it paid off.
Key takeaway: Locality is the share of a region's bytes stored on the server that hosts it. HDFS creates it automatically through local first replicas, and moves, restarts, crashes, bulk loads and the HDFS balancer destroy it. Measure it byte-weighted and alongside cache hit ratio, protect it with region movers and a quiet balancer, and restore it selectively, ranking by remote bytes, with forced major compaction of cold stores or block movement.