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
| Role | Count | What it needs | What its failure costs |
|---|---|---|---|
| ZooKeeper | 3 or 5 (odd) | low-latency disk for its transaction log, stable network, no long pauses | losing a majority stops the HBase master, RegionServer liveness tracking and NameNode failover |
| HMaster | 1 active + 1 or 2 backups | modest CPU and memory | no failover or DDL until a backup takes over; reads and writes continue |
| NameNode | active + standby | large heap for namespace metadata, reliable disks | without HA, all of HDFS stops |
| JournalNode | 3 (odd) | low-latency disk for edit logs | losing a majority stops NameNode edits |
| RegionServer | most nodes | large heap for memstore and block cache, many cores | its regions are offline until reassigned and their WAL is replayed |
| DataNode | same nodes as RegionServers | many disks, network bandwidth | blocks 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
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.
| Question | Arithmetic | Result |
|---|---|---|
| Workers for capacity alone, keeping 30% free | 360 / 0.7 / 48 | about 11 |
| Workers so two surviving racks stay below 85% full | 360 / (0.85 × 48 × 2/3) | about 13.2, so 15 (5 per rack) |
| Data per RegionServer | 120 TB / 15 | 8 TB |
| Regions per server at 20 GB regions | 8 TB / 20 GB | about 400 |
| Write-active regions before global flush pressure | 31 GB heap × 0.4 / 128 MB | about 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
| Failure | Symptom | Topology fix |
|---|---|---|
| Rack script missing or wrong | all nodes in /default-rack; a rack loss loses blocks | maintain the script and alert on the default rack |
| ZooKeeper on a busy worker | random RegionServers declared dead | dedicated control nodes, separate log disk |
| Active NameNode and HMaster on one host | one host failure triggers two failovers | split active roles across racks |
| Even-sized ensemble | 4 servers tolerate 1 failure, like 3, but each write waits for 3 acknowledgements | always 3 or 5 |
| Heterogeneous workers | regions on small nodes run hot; balancer churn | identical workers per group |
| Cluster stretched across sites | high write latency; partition stalls a site | one cluster per site plus replication |
| Locality never recovers | reads cross the network weeks after a restart | compact 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
- Run
hdfs dfsadmin -printTopologyand confirm that no node is in/default-rack. - List where every ZooKeeper server, JournalNode, NameNode and HMaster runs, and confirm no two of a kind share a rack.
- Move ZooKeeper's transaction log to its own disk, and off any worker.
- Enable short-circuit reads on every RegionServer and check locality after each restart.
- Recompute worker count for the loss of one rack, not just for capacity.
- Move clients to
RpcConnectionRegistryif you run 2.5 or later. - Use RSGroups or separate clusters for workloads whose failures must not spread.
- Rehearse a rack failure in staging and time each phase: reassignment, NameNode failover, re-replication.