An HBase cluster is several distributed systems running together. HBase itself has HMasters and RegionServers. Underneath it is HDFS, with NameNodes, DataNodes and JournalNodes. Beside it is a ZooKeeper ensemble that both depend on. Topology is the decision about which of these processes run on which machines, in which racks, and with which disks. It decides two things that no configuration change can fix later: how much one failure takes with it, and whether reads are served from local disk or across the network.

This article builds a topology from first principles. It starts with what each role needs, then covers co-locating RegionServers with DataNodes, rack awareness, placing the control plane, isolating workloads, and spanning sites. A worked example sizes a 15-worker cluster and walks through what happens when a whole rack fails. The roles themselves are covered in HMaster and related articles. Here the subject is where they run.

Roles and what each one needs

RoleCountWhat it needsWhat its failure costs
ZooKeeper3 or 5 (odd)low-latency disk for its transaction log, stable network, no long pauseslosing a majority stops the HBase master, RegionServer liveness tracking and NameNode failover
HMaster1 active + 1 or 2 backupsmodest CPU and memoryno failover or DDL until a backup takes over; reads and writes continue
NameNodeactive + standbylarge heap for namespace metadata, reliable diskswithout HA, all of HDFS stops
JournalNode3 (odd)low-latency disk for edit logslosing a majority stops NameNode edits
RegionServermost nodeslarge heap for memstore and block cache, many coresits regions are offline until reassigned and their WAL is replayed
DataNodesame nodes as RegionServersmany disks, network bandwidthblocks become under-replicated and are copied again

Two patterns follow. The control-plane roles (ZooKeeper, JournalNodes, NameNodes, HMasters) are small, few in number, and quorum-based or active-standby. They belong on a few reliable machines, spread across failure domains. The data-plane roles (RegionServers and DataNodes) are many, heavy, and share data. They belong together on every worker. Mixing the two, for example running ZooKeeper on a busy worker, is the most common topology mistake: the worker's load turns into control-plane latency for the whole cluster.

RegionServers and DataNodes together

A RegionServer stores its data as HFiles in HDFS. When it flushes a memstore or compacts, HDFS writes the first replica of each block to the local DataNode, if the writer runs on one. The second replica goes to a node in a different rack and the third to another node in that second rack. So a RegionServer co-located with a DataNode ends up with a local copy of almost every file it wrote, and reads can be served from local disk.

Local data only helps if the read path can use it. With short-circuit reads, the HDFS client inside the RegionServer reads block files directly, through a file descriptor passed over a Unix domain socket, instead of streaming them through the DataNode process. Configure it in the hdfs-site.xml used by both the DataNodes and the RegionServers:

<property>
  <name>dfs.client.read.shortcircuit</name>
  <value>true</value>
</property>
<property>
  <name>dfs.domain.socket.path</name>
  <value>/var/lib/hadoop-hdfs/dn_socket</value>
</property>

Locality decays whenever regions move: after a RegionServer crash, a balancer run or a rolling restart, a region is served by a node that holds none of its files, so every read crosses the network. A major compaction rewrites the region's files from the new server, which restores locality. How to measure and repair this is covered in HBase region locality. For topology, the lesson is that every worker should run both processes, with matching disk and memory, so that any region can live anywhere without becoming remote for good.

Rack awareness

HDFS only spreads replicas across racks if it knows which node is in which rack. Hadoop asks a script named by net.topology.script.file.name in core-site.xml. It is called with host names or IP addresses as arguments and must print one rack path for each. Nodes it cannot map fall into /default-rack. When every node is in the default rack, HDFS can put all three replicas behind one switch, and nothing will warn you.

#!/usr/bin/env python3
# /etc/hadoop/conf/rack_topology.py -- prints one rack per argument, in order.
import sys

RACKS = {
    "10.1.1.": "/dc1/rack1",
    "10.1.2.": "/dc1/rack2",
    "10.1.3.": "/dc1/rack3",
}
HOSTS = {"w01.example": "/dc1/rack1", "w06.example": "/dc1/rack2", "w11.example": "/dc1/rack3"}  # and so on

for arg in sys.argv[1:]:
    rack = HOSTS.get(arg) or next((r for prefix, r in RACKS.items() if arg.startswith(prefix)), None)
    print(rack or "/default-rack")

Check the result with hdfs dfsadmin -printTopology after every hardware change, and alert if any node appears in /default-rack. Rack awareness is what lets the cluster survive a top-of-rack switch failure or a rack power loss without losing a block. It also creates cross-rack write traffic: two of the three replicas of every flush and compaction cross the spine. Size the uplinks for compaction bandwidth, not for client traffic.

Placing the control plane

A 3-rack HBase cluster: one control-plane node and five workers per rackClientsbootstrap from mastersSpine switchRack 1 top-of-rack switchControl nodeNameNode (active)HMaster (backup), ZooKeeperJournalNode, ZKFC if NameNode5 workersRegionServer + DataNode eachRack 2 top-of-rack switchControl nodeNameNode (standby)HMaster (active), ZooKeeperJournalNode, ZKFC if NameNode5 workersRegionServer + DataNode eachRack 3 top-of-rack switchControl nodeno NameNodeHMaster (backup), ZooKeeperJournalNode, ZKFC if NameNode5 workersRegionServer + DataNode eachEach HDFS block has replicas in two racks; ZooKeeper and JournalNodes keep a majority after any one rack failsWorker disks hold HDFS data; ZooKeeper's transaction log gets its own disk on each control node
Three control nodes, one per rack, each running ZooKeeper and a JournalNode; NameNodes and HMasters split across them; workers run RegionServer and DataNode together.

With three racks, put one control node in each rack. Each runs a ZooKeeper server and a JournalNode, so losing a rack leaves two of three in both quorums. Put the active and standby NameNodes in different racks, each with its ZKFC failover controller, and spread the HMasters so that the active one and at least one backup never share a rack. With five racks or a higher availability target, use five ZooKeeper servers, which tolerate two failures. More than five slows every write to ZooKeeper without adding useful tolerance.

ZooKeeper's write latency is its transaction-log fsync. Give that log its own disk on each control node, separate from the snapshot directory and from anything with heavy I/O. HBase uses ZooKeeper sessions to decide whether a RegionServer is alive. If the server misses its session timeout, set by zookeeper.session.timeout, because of a long garbage-collection pause or a ZooKeeper server stalled on disk, the master declares it dead and reassigns its regions. A slow ensemble therefore causes false failovers across the cluster.

Clients also depend on how they find the cluster. Older clients connect to ZooKeeper to locate the meta table, which ties every application to the ensemble. HBase 2.5 added RpcConnectionRegistry, which became the default in 3.0. A client bootstraps from a list of servers and refreshes it over RPC, so client connections stop counting against ZooKeeper. The older MasterRegistry is deprecated.

<!-- client hbase-site.xml (HBase 2.5+): no ZooKeeper quorum needed by applications -->
<property>
  <name>hbase.client.registry.impl</name>
  <value>org.apache.hadoop.hbase.client.RpcConnectionRegistry</value>
</property>
<property>
  <name>hbase.client.bootstrap.servers</name>
  <value>m1.example:16000,m2.example:16000,m3.example:16000</value>
</property>

Isolating workloads

One cluster often serves several workloads: a latency-sensitive serving table, a bulk-loaded analytics table, and a write-heavy event log. Topology offers three levels of isolation:

  • Separate tables and quotas on shared RegionServers. This is the cheapest option, and the weakest: a compaction storm or a long GC pause on one server affects every table hosted there.
  • RegionServer groups. Tables are pinned to a named group of servers, so the serving table's RegionServers never host the event log's regions. HDFS, ZooKeeper and the master remain shared. See HBase RSGroups.
  • Separate clusters. Nothing is shared, which also means a separate control plane to operate.

For read availability during the minutes a failed RegionServer's regions are being reassigned, region replicas keep read-only secondary copies of a region open on other servers. The balancer tries to place them on different hosts and racks, so a rack failure leaves a readable copy, at the cost of reads that may be slightly stale and extra memory per replica.

More than one site

Do not stretch one cluster across data centres. Every write waits for a WAL sync through an HDFS pipeline, and every ZooKeeper write waits for a quorum. Spread either across a WAN link and each write pays that link's latency, while a network partition can leave the minority site unable to make progress. The usual design is one complete cluster per site, connected by asynchronous replication, with applications writing to one site at a time. The second site lags by seconds, not by the length of a backup cycle, and a failover is a decision your application makes, not something the cluster does by itself.

On Kubernetes the same rules apply with different names: a rack becomes a zone or node label, ZooKeeper and JournalNodes need pod anti-affinity and local persistent volumes, and short-circuit reads need the socket shared between pods. The practical issues are in HBase on Kubernetes.

Worked example: 15 workers and a lost rack

Suppose you need to store 120 TB of compressed table data. With HDFS replication of 3, that is 360 TB on disk. Compactions temporarily need room for the old and new files side by side, and a rack failure must be absorbed by the remaining racks. A worker with twelve 4 TB data disks offers 48 TB raw.

QuestionArithmeticResult
Workers for capacity alone, keeping 30% free360 / 0.7 / 48about 11
Workers so two surviving racks stay below 85% full360 / (0.85 × 48 × 2/3)about 13.2, so 15 (5 per rack)
Data per RegionServer120 TB / 158 TB
Regions per server at 20 GB regions8 TB / 20 GBabout 400
Write-active regions before global flush pressure31 GB heap × 0.4 / 128 MBabout 99

The capacity line says 11 workers, but the failure line says 15. That gap is typical: the topology is sized for the day a rack goes missing, not for the steady state. The memstore line uses the defaults hbase.regionserver.global.memstore.size = 0.4 and hbase.hregion.memstore.flush.size = 128 MB. It shows that only about 99 regions per server can each hold a full memstore before the server reaches its global memstore limit, so keep the write-heavy key ranges spread over servers rather than packed onto a few.

Now lose rack 2 at 03:00. Five RegionServers vanish, along with about a third of the regions. The master sees their ZooKeeper sessions expire, splits their WALs and reassigns the regions to the ten survivors. Rack 2 also held the standby NameNode and the active HMaster. HDFS keeps running on the active NameNode in rack 1, but with no standby until rack 2 returns, so alert on it: one more failure would stop HDFS. A backup HMaster in rack 1 or 3 takes over once the active master's ZooKeeper session expires. Two of three ZooKeeper servers and JournalNodes remain, so both quorums hold. No block is lost, because every block had replicas in two racks. About 120 TB of replicas are now missing and HDFS starts copying them: at an aggregate 5 GB/s that takes roughly seven hours, competing with compactions that are also restoring locality on the reassigned regions. Plan maintenance windows and alerts around that tail, not around the few minutes the failover takes.

Failure modes

FailureSymptomTopology fix
Rack script missing or wrongall nodes in /default-rack; a rack loss loses blocksmaintain the script and alert on the default rack
ZooKeeper on a busy workerrandom RegionServers declared deaddedicated control nodes, separate log disk
Active NameNode and HMaster on one hostone host failure triggers two failoverssplit active roles across racks
Even-sized ensemble4 servers tolerate 1 failure, like 3, but each write waits for 3 acknowledgementsalways 3 or 5
Heterogeneous workersregions on small nodes run hot; balancer churnidentical workers per group
Cluster stretched across siteshigh write latency; partition stalls a siteone cluster per site plus replication
Locality never recoversreads cross the network weeks after a restartcompact moved regions, then measure locality

Trade-offs

Dedicated control nodes cost three machines that store no data, in return for a control plane that worker load cannot disturb. Rack awareness costs cross-rack bandwidth on every flush and compaction in return for surviving a switch failure. Sizing for a lost rack leaves capacity idle most days. RSGroups isolate RegionServers but still share HDFS and ZooKeeper; separate clusters isolate everything but double the operational work. Replication between sites keeps writes fast, but the standby site always lags slightly.

What to do next

  1. Run hdfs dfsadmin -printTopology and confirm that no node is in /default-rack.
  2. List where every ZooKeeper server, JournalNode, NameNode and HMaster runs, and confirm no two of a kind share a rack.
  3. Move ZooKeeper's transaction log to its own disk, and off any worker.
  4. Enable short-circuit reads on every RegionServer and check locality after each restart.
  5. Recompute worker count for the loss of one rack, not just for capacity.
  6. Move clients to RpcConnectionRegistry if you run 2.5 or later.
  7. Use RSGroups or separate clusters for workloads whose failures must not spread.
  8. Rehearse a rack failure in staging and time each phase: reassignment, NameNode failover, re-replication.
Key takeaway: Put the quorum-based control plane on a few dedicated nodes, one per rack, with ZooKeeper's log on its own disk. Run a RegionServer and a DataNode together on every identical worker, with short-circuit reads. Make HDFS rack-aware and check it. Size the cluster for the loss of a whole rack. Isolate workloads with RSGroups or separate clusters, and connect sites with replication rather than stretching one cluster.