Planning a Hadoop cluster is two jobs that people often mix up. Sizing decides how much: how many terabytes, how many cores, how much NameNode heap. Layout decides where: which service runs on which machine, how machines map to racks and switches, how disks are formatted, and how the cluster will grow without an outage. A cluster with the right totals and the wrong layout still loses data when a rack fails, still stalls when a single master host dies, and still needs a weekend outage to add nodes.

The arithmetic is covered in our Hadoop capacity planning guide, so this article starts where it ends. Given a node count, it explains node classes and role placement, the high-availability quorums, rack topology and block placement, the network, worker disk and operating system layout, YARN resources per node, edge access, a worked three-rack plan, and how to grow and upgrade. Property names are from Apache Hadoop 3; packaged distributions may manage them through their own tooling and ship different defaults.

Advertisement

Node classes and role placement

Every Hadoop cluster has the same small set of node classes, and most planning mistakes come from blurring them.

ClassRunsHardware emphasisCount
MasterNameNode, ResourceManager, JournalNode, ZooKeeper, history serversRAM, reliable disks (RAID for metadata), redundant power and NICs3 in an HA cluster, 5 at large scale
WorkerDataNode and NodeManager togetherMany data disks, balanced cores and memoryMost of the cluster
EdgeClient libraries, Hive and Spark gateways, ingestion tools, user shellsModest; network access to users1 to a few, behind a load balancer
UtilityMetastore database, monitoring, Kerberos KDC, configuration managementSmall, reliableOften virtual machines

Run the DataNode and NodeManager on the same worker, because the whole point of Hadoop's scheduling is to send computation to the node holding the block. Keep masters free of worker roles; a NodeManager running a memory-hungry container next to an active NameNode is how a busy afternoon turns into a long garbage collection pause and a failover. Keep users off masters and workers entirely: everything a person or external system does should go through an edge node, which gives you one place to install clients, enforce access and absorb mistakes such as a large hdfs dfs -get to local disk.

Rack 1Rack 2Rack 3Master ANameNode 1, JN, ZKMaster BNameNode 2, JN, ZK, RM 1Master CJN, ZK, RM 2, historyWorkers 1 to 9DataNode + NodeManagerWorkers 10 to 18DataNode + NodeManagerWorkers 19 to 27DataNode + NodeManagerTop-of-rack switchTop-of-rack switchTop-of-rack switchSpine switchesuplinks sized to the oversubscription targetEdge nodesclients, gatewaysOne master per rack: losing a rack keeps a NameNode, a ResourceManager and a quorum of JN and ZK
A three-rack layout: each rack holds one master node and a third of the workers, so the HA quorums of JournalNodes and ZooKeeper survive the loss of any single rack, and HDFS spreads replicas across racks through the topology script.

The high-availability control plane

The HDFS namespace lives in the NameNode, so it must not have a single point of failure. In the standard design, an active NameNode writes every edit to a quorum of JournalNodes, standby NameNodes tail those edits, and a ZooKeeper Failover Controller next to each NameNode uses ZooKeeper to elect the active one. Hadoop 3 allows more than two NameNodes; the documentation recommends three and advises against going above five because of the overhead of communication.

<!-- hdfs-site.xml (excerpt) -->
<property><name>dfs.nameservices</name><value>prod</value></property>
<property><name>dfs.ha.namenodes.prod</name><value>nn1,nn2</value></property>
<property><name>dfs.namenode.rpc-address.prod.nn1</name><value>master-a:8020</value></property>
<property><name>dfs.namenode.rpc-address.prod.nn2</name><value>master-b:8020</value></property>
<property><name>dfs.namenode.shared.edits.dir</name>
  <value>qjournal://master-a:8485;master-b:8485;master-c:8485/prod</value></property>
<property><name>dfs.ha.automatic-failover.enabled</name><value>true</value></property>
<!-- core-site.xml -->
<property><name>ha.zookeeper.quorum</name><value>master-a:2181,master-b:2181,master-c:2181</value></property>

JournalNodes and ZooKeeper are both majority quorums, so use an odd number, three or five, and place them on different racks with different power feeds. Three tolerates one failure, five tolerates two. Give JournalNodes and ZooKeeper their own disks: both fsync on every write, and sharing a spindle with NameNode checkpoints or logs adds latency to every namespace change. Make the ResourceManager highly available too, with a ZooKeeper-backed state store, so a failover does not lose the list of running applications. The NameNode HA article covers failover and fencing in detail.

Advertisement

Racks, the topology script and replica placement

HDFS does not know what a rack is until you tell it. The property net.topology.script.file.name points to a script that receives host names or IP addresses and prints a rack path for each. Without it, every node is in /default-rack and HDFS has no way to keep replicas apart.

#!/usr/bin/env bash
# /etc/hadoop/conf/topology.sh  -- maps each argument to a rack path
# topology.map lines look like: 10.20.1.15 /dc1/rack1
MAP=/etc/hadoop/conf/topology.map
for host in "$@"; do
  rack=$(awk -v h="$host" '$1 == h {print $2}' "$MAP")
  echo "${rack:-/dc1/default-rack}"
done

With topology known, the default placement policy for three replicas puts the first replica on the writer's node (or a random node if the writer is outside the cluster), the second on a node in a different rack, and the third on a different node in the same rack as the second. One copy therefore survives the loss of any rack, while two of the three replicas share a rack to save cross-rack bandwidth on the write path. YARN uses the same topology to prefer rack-local tasks when a node-local slot is not free.

Practical rules follow. Use at least three racks, so a rack failure leaves both a quorum and a second rack to re-replicate into. Keep racks roughly equal in capacity; a small rack fills first because one replica of everything must land outside the writer's rack. Erasure coded directories spread the cells of each block group across as many racks as they can, which is another reason to count racks before choosing a policy. Keep the map file in configuration management and treat an unmapped host as a deployment failure. The rack awareness article goes deeper on the policy.

Network design

Hadoop moves a lot of data east-west: shuffle between tasks, replication pipelines on write, and re-replication after a failure. A two-tier leaf-spine fabric with a top-of-rack switch per rack and uplinks to every spine switch is the common design. The number to plan is the oversubscription ratio: the bandwidth of all server ports in a rack divided by the rack's uplink bandwidth.

  • For example, 20 workers at 25 Gb/s is 500 Gb/s of server bandwidth. Four 100 Gb/s uplinks give 400 Gb/s, a ratio of 1.25:1, which is comfortable for shuffle-heavy workloads. Two uplinks give 2.5:1, which is acceptable for mostly scan-and-filter work but slows recovery after a rack loss.
  • Bond two NICs per server to two switches or use dual-homed racks if a single top-of-rack switch failure must not take the whole rack offline.
  • Use a separate management network for out-of-band access, so you can still reach a node when its data network is saturated.
  • Keep forward and reverse DNS consistent for every host. Kerberos and the topology script both depend on names resolving the same way in both directions.

Worker disks and operating system

Workers use plain disks, not RAID. HDFS already keeps replicas on other machines, so RAID on a DataNode spends capacity and write bandwidth to protect against a failure HDFS handles better. Mount each data disk separately and list them all in dfs.datanode.data.dir; the DataNode spreads blocks across volumes. Masters are the opposite: protect the NameNode metadata directories with RAID 1, and list two directories in dfs.namenode.name.dir so the image is written twice.

# worker: one filesystem per disk, no access-time updates
/dev/sdb1  /data/1  xfs  defaults,noatime  0 0
/dev/sdc1  /data/2  xfs  defaults,noatime  0 0
# ... one line per data disk

# kernel and limits for Hadoop daemons
sysctl -w vm.swappiness=1
echo never > /sys/kernel/mm/transparent_hugepage/enabled
echo never > /sys/kernel/mm/transparent_hugepage/defrag
# /etc/security/limits.d/hadoop.conf : nofile and nproc limits for hdfs and yarn users

Set dfs.datanode.failed.volumes.tolerated deliberately. Its default of 0 shuts the DataNode down when one disk fails, which removes every other disk on the node from service too. On a twelve-disk worker, tolerating one or two failed volumes keeps the node serving while you schedule a replacement. Keep the operating system and logs on a separate pair of disks so a full log partition never competes with data. Put YARN's yarn.nodemanager.local-dirs, where shuffle and spill data land, on the same data disks so intermediate I/O spreads across all spindles.

YARN resources per node

YARN does not discover how much it may use; you declare it per node with yarn.nodemanager.resource.memory-mb and yarn.nodemanager.resource.cpu-vcores. Start from physical memory and subtract the operating system, the DataNode heap, the NodeManager heap and anything else on the box, then give YARN what remains. On a 256 GB worker that might be 20 GB for the system and daemons and about 230 GB for containers. Leaving too little headroom invites the kernel's out-of-memory killer, which tends to choose the DataNode.

If the cluster mixes generations of hardware, the declared resources differ per node, which configuration management handles with host groups. If some nodes have GPUs or much more memory, use node labels so only the queues that need them can schedule there, instead of letting any job land on expensive hardware.

Security and edge access

Plan security before the first byte lands, because enabling it later means touching every service and every client. A production cluster needs Kerberos for authentication, which needs a KDC that is itself highly available and reachable from every node, service principals per host, and keytabs distributed by automation. Wire encryption, HDFS encryption zones for sensitive directories, and an authorization layer for tables follow from the same plan. The Hadoop Kerberos guide covers the setup. From a layout point of view, the key decision is that edge nodes are the only hosts users can log in to, and that firewalls allow cluster ports only between cluster nodes and from edge nodes.

Worked example: a three-rack cluster

Suppose capacity planning produced 27 workers with twelve 16 TB disks, 64 cores and 256 GB each, three masters, and two edge nodes. The figures are illustrative; the reasoning is what transfers.

  1. Racks: three racks of nine workers plus one master each, as in the diagram. Each rack holds about a third of raw capacity, so losing one rack leaves two thirds of the disks to re-replicate into, which capacity planning confirmed would fit below 80 percent utilisation.
  2. Quorums: JournalNode and ZooKeeper on all three masters, one per rack, so any single rack loss leaves two of three. NameNodes on masters A and B, ResourceManagers on B and C, so neither active role is pinned to one host and each can fail over independently.
  3. Network: nine workers at 25 Gb/s per rack is 225 Gb/s; two 100 Gb/s uplinks give about 1.1:1. Masters are dual-homed.
  4. Disks: each worker mounts twelve XFS data disks with noatime and tolerates two failed volumes; operating system on a separate mirrored pair. Masters use RAID 1 for NameNode, JournalNode and ZooKeeper directories, each on its own pair.
  5. YARN: about 230 GB and 60 vcores declared per worker, leaving the rest for the DataNode, NodeManager and system.
  6. Growth: the spine has free ports for three more racks, and the rack plan reserves IP ranges and topology map entries for them.

Growth, decommissioning and upgrades

Add capacity in whole racks when you can, so racks stay balanced and the placement policy keeps working well. A new node joins by being added to the topology map and the include files, starting its daemons and refreshing nodes; it then receives only new writes, so run the HDFS balancer with a bandwidth limit to move existing blocks gradually. Remove nodes by decommissioning, which re-replicates their blocks before you switch them off, never by powering them down.

Upgrades are easiest when the cluster was planned for them. With HA NameNodes and ResourceManagers, a rolling upgrade can update one master at a time and workers in batches, as the rolling upgrade guide explains. Keep the batch size below the number of nodes whose loss the cluster tolerates, which in practice means no more than part of one rack at a time.

Failure modes and trade-offs

Planning mistakeWhat happensPrevention
No topology scriptAll replicas of a block can land in one rack; a rack failure loses dataMap every host; alert on hosts in the default rack
Quorum on two racksLosing the rack with two members stops HDFS writesOne member per rack, odd count, three or more racks
RAID on workersLess capacity, slower writes, longer rebuildsJBOD with failed volume tolerance
Failed volumes tolerated set to 0One bad disk takes a whole node outTolerate one or two on dense nodes
Masters running workersGC pauses and failovers under loadDedicated masters
High oversubscriptionSlow shuffle, very slow recovery after a rack lossSize uplinks to workload; measure
Users on cluster nodesLocal disks fill, ad hoc load on mastersEdge nodes only

The trade-offs are mostly money against blast radius. More racks and more masters cost more but shrink the share of the cluster any single failure affects. Denser workers lower cost per terabyte but make each node failure and each decommission move more data. Converged nodes keep data locality, while separating storage and compute allows independent scaling at the price of network traffic for every read. Decide these explicitly and write them down next to the capacity numbers.

What to do next

  1. Finish the sizing with the capacity planning guide, then draw the rack diagram for your node count before ordering hardware.
  2. Assign every role to a host class and confirm no master runs worker roles and no user logs in outside the edge nodes.
  3. Place JournalNodes and ZooKeeper one per rack with odd counts, and give each its own disk.
  4. Write the topology script and map file, and add a check that fails deployment if any host resolves to the default rack.
  5. Set the oversubscription target with your network team and verify it with a test transfer between racks.
  6. Standardise the worker disk layout, mount options, kernel settings and failed volume tolerance in configuration management.
  7. Plan Kerberos and the edge access model now, and reserve spine ports, address ranges and rack positions for the next expansion.
Key takeaway: Sizing gives the totals; layout decides whether the cluster survives failures and grows without outages. Keep masters, workers and edge nodes separate, put one JournalNode and ZooKeeper member per rack in odd quorums, and tell HDFS the topology so replicas span racks. Use plain disks on workers with failed volume tolerance, size uplinks to the workload, declare YARN resources with headroom, plan security first, and grow in balanced racks with decommissioning and rolling upgrades.