Capacity planning for Hadoop answers one question: how many nodes of which shape do you need to run this workload for the next planning horizon without running out of anything. Disk is the obvious resource, and most first attempts stop there. A cluster can just as easily run out of NameNode heap, YARN vcores, container memory or the network bandwidth it needs to re-protect data after a node dies. Those four resources scale with different inputs, so they have to be sized separately and then reconciled.

This article builds that method from first principles, turns it into a short Python calculator, and runs a worked example in which the obvious answer (size for disk) is wrong by a factor of three. The figures are the calculator's output, not estimates typed by hand, so you can rerun it with your own inputs and get a number you can defend.

Advertisement

Why four dimensions, not one

HDFS stores bytes on DataNodes, but every file, directory and block is also an object held in the NameNode's Java heap. YARN runs work in containers whose size is set in vcores and megabytes, and the NodeManager advertises a fixed amount of each. When a node fails, HDFS re-creates the lost replicas or erasure-coded cells from the survivors, and that traffic crosses the network while production jobs are still running.

Each resource is driven by a different input. Physical storage follows bytes ingested multiplied by retention and redundancy overhead. NameNode memory follows the number of files and blocks, which depends on file size far more than on total bytes. Compute follows how much processing the jobs need inside their batch windows. Recovery time follows how much data one node holds and how much bandwidth the survivors can spare. A plan that sizes only one of them will meet the others by surprise.

Workload inputsingest/day, retention, growthFile shapeaverage file size, block sizeJob demandvcore-hours per batch windowStoragelogical x overhead / fillNameNode memoryfiles + blocksComputevcores and memory per nodeNetwork recoveryre-protect one lost nodeNode countmax(storage, compute, EC min) + 1heap, not nodescheck, not sizeReconcilere-pick the node shape if one dimension wastes the othersEach dimension is sized independently; the largest wins and the others become checks.
The sizing flow: three kinds of input feed four independent dimensions, and the node count is the largest requirement plus a spare.

Storage: from logical bytes to raw disks

Start with logical bytes per dataset: daily ingest multiplied by retention in days. Multiply by the redundancy overhead. Three-way replication, the HDFS default, stores every byte three times. The built-in Reed-Solomon policy RS-6-3 stores six data cells and three parity cells, an overhead of 1.5, and tolerates the loss of any three cells in a group. Erasure coding halves the bytes but costs CPU on write and turns recovery into a decode that reads six surviving cells for every lost one, so it suits cold or warm data that is written once and scanned, not small hot files. RS-6-3 also needs at least nine DataNodes to place a full block group, and rack-level tolerance needs the cells spread across racks; see HDFS erasure coding and rack awareness.

Then subtract what HDFS does not get. YARN local directories hold shuffle and spill data for MapReduce and Spark, logs need space, and dfs.datanode.du.reserved keeps a per-volume floor free for non-HDFS use. Reserving 15 to 25 percent of raw disk for this is a common starting point; measure your peak shuffle footprint and adjust. Finally, plan to a target fill, not to 100 percent. Around 75 percent leaves room for the balancer to work, for a node's data to land elsewhere after a failure, and for the growth you did not forecast. Well above that, individual volumes fill first, writes to those DataNodes fail, and the balancer has little free space to work with.

Advertisement

NameNode memory: count objects, not bytes

The NameNode keeps the whole namespace in heap: an inode for every file and directory and a record for every block, with the block-to-DataNode map built from block reports. A widely used rule of thumb budgets about 1 GB of heap per million blocks, with headroom on top; treat it as a starting point and measure your own NameNode with its JMX metrics. What matters for planning is that the object count is driven by file size. A petabyte stored as 512 MB files is a few million objects. The same petabyte as 1 MB files is a billion, which no single NameNode serves comfortably.

Erasure coding changes the arithmetic slightly in your favour: the NameNode tracks a block group rather than each internal cell, so with RS-6-3 and 128 MB cells one group covers up to 768 MB of data. The calculator models this with a block unit of 768 MB for erasure-coded datasets. If the forecast lands in the hundreds of millions of objects, the fix is upstream (compaction, larger files) or structural (federation), not a bigger heap; see the small files problem and NameNode internals.

Compute: vcores and memory per node

Compute is sized from demand, not from disks. Take the vcore-hours your critical jobs consume and divide by the window they must finish in: 6,000 vcore-hours in a six-hour nightly window means 1,000 vcores busy at once. Divide by what one NodeManager advertises. On a 64-core, 512 GB worker, leave cores and memory for the OS, the DataNode, the NodeManager and page cache, and give YARN the rest. Keep the memory-to-vcore ratio of the node close to the ratio your containers ask for, or one resource strands the other: 8 GB per vcore suits Spark executors around that ratio, while a fleet of 2 GB containers would leave memory idle.

<!-- yarn-site.xml for a 64-core, 512 GB worker -->
<property><name>yarn.nodemanager.resource.memory-mb</name><value>458752</value></property>  <!-- 448 GB -->
<property><name>yarn.nodemanager.resource.cpu-vcores</name><value>56</value></property>
<property><name>yarn.scheduler.minimum-allocation-mb</name><value>1024</value></property>
<property><name>yarn.scheduler.maximum-allocation-mb</name><value>65536</value></property>

<!-- hdfs-site.xml: keep non-HDFS space off the table -->
<property><name>dfs.datanode.du.reserved</name><value>107374182400</value></property>  <!-- 100 GiB per volume -->
<property><name>dfs.datanode.failed.volumes.tolerated</name><value>1</value></property>

Container sizing and the scheduler's minimum and maximum allocation are covered in YARN containers. For capacity planning, the output is one number: usable vcores per node, and a check that memory per node covers those vcores at your typical container ratio.

Network: how long until data is safe again

When a DataNode dies, every block it held is under-protected until the NameNode schedules new copies. The time to recover is roughly the data on the failed node divided by the re-replication bandwidth the cluster will actually spend, which is throttled on purpose by settings such as dfs.namenode.replication.max-streams so recovery does not starve production. More nodes mean less data per node and more sources to copy from, so recovery gets faster as the cluster widens. Dense storage nodes with 200 TB or more each are where this bites: a single failure can leave data exposed for many hours. Recovery time is a check, not a sizing input. If it is too long, use more, smaller nodes or a faster network.

A calculator you can run

The function below takes datasets, growth and a node shape and returns all four dimensions. Growth compounds to the end of the horizon and is applied to both bytes and compute, because more data usually means more processing. Re-replication bandwidth is assumed to be 10 percent of each surviving node's NIC, a deliberately cautious figure. EC reconstruction reads several cells per lost one, so for erasure-coded data the real network cost is higher than this model shows.

import math

TB = 1e12

def plan(datasets, growth_per_year, horizon_months, node, job, fill=0.75, nic_gbps=25):
    g = (1 + growth_per_year) ** (horizon_months / 12)
    logical = physical = objects = 0
    for d in datasets:
        size = d["tb_per_day"] * d["retention_days"] * g * TB
        files = size / d["avg_file_bytes"]
        blocks_per_file = math.ceil(d["avg_file_bytes"] / d["block_unit_bytes"])
        logical += size
        physical += size * d["overhead"]
        objects += files * (1 + blocks_per_file)
    hdfs_per_node = node["disks"] * node["disk_tb"] * TB * (1 - node["non_hdfs_frac"])
    storage_nodes = math.ceil(physical / fill / hdfs_per_node)
    vcores_needed = job["vcore_hours"] * g / job["window_hours"]
    compute_nodes = math.ceil(vcores_needed / node["yarn_vcores"])
    nodes = max(storage_nodes, compute_nodes, node["min_nodes"]) + 1  # N+1 spare
    per_node = physical / nodes
    rebuild_s = per_node / ((nodes - 1) * nic_gbps / 8 * 1e9 * 0.1)  # 10% of each NIC
    return dict(growth=g, logical_tb=logical / TB, physical_tb=physical / TB,
                storage_nodes=storage_nodes, compute_nodes=compute_nodes, nodes=nodes,
                fill_pct=100 * physical / (nodes * hdfs_per_node),
                nn_objects_m=objects / 1e6, rebuild_h=rebuild_s / 3600)

datasets = [
    dict(name="raw", tb_per_day=2.0, retention_days=30, overhead=3.0,
         avg_file_bytes=16e6, block_unit_bytes=128e6),
    dict(name="curated", tb_per_day=0.8, retention_days=365, overhead=1.5,
         avg_file_bytes=512e6, block_unit_bytes=768e6),  # RS-6-3: one block group per 6 x 128 MB
]
node = dict(disks=12, disk_tb=16, non_hdfs_frac=0.20, yarn_vcores=56, min_nodes=9)
job = dict(vcore_hours=6000, window_hours=6)
print(plan(datasets, growth_per_year=0.40, horizon_months=18, node=node, job=job))

Worked example: the cluster is compute-bound

The inputs describe a typical event platform. Raw landing data arrives at 2 TB a day in 16 MB files, is kept for 30 days and is replicated three ways because it is hot and rewritten. Curated Parquet arrives at 0.8 TB a day in 512 MB files, is kept for a year and uses RS-6-3. Ingest grows 40 percent a year, and the plan covers 18 months. Nodes have twelve 16 TB disks with 20 percent reserved for non-HDFS use, and expose 56 vcores to YARN. The nightly pipeline needs 6,000 vcore-hours in a six-hour window.

QuantityCalculator output
Growth factor at 18 months1.66
Logical data583 TB
Physical data after replication and EC1,024 TB
Nodes needed for storage at 75% fill9
Nodes needed for compute30
Nodes to buy (largest plus one spare)31
Actual HDFS fill on 31 nodes21.5%
NameNode objects14.3 million
Re-protect one node at 10% of a 25 Gbit/s NICabout 1 hour

Sizing for disk alone would have bought nine or ten nodes, and the nightly window would have overrun by a factor of three. The real answer is 31 nodes, at which point the disks are about one-fifth full. That is the reconcile step: the node shape is wrong for this workload. Halving the disks to six per node puts fill at about 43 percent on the same 31 nodes and removes a large share of the hardware cost; alternatively, move curated data to denser storage-only nodes and run compute on diskless or lightly disked workers.

The NameNode result is just as instructive. The raw zone holds 17 percent of the logical bytes but about 87 percent of the 14.3 million objects, because its files are 16 MB. Using the rule of thumb, that suggests a heap in the mid-teens of gigabytes, so a 32 GB heap has room to grow. Compacting raw files to 128 MB before they age into the curated zone would cut the raw zone's object count by roughly a factor of eight and the total to about a quarter, which is cheaper than any heap tuning.

Failure modes

  • Planning to 100 percent. The cluster hits full volumes during the first node failure, because the failed node's data has nowhere to go. Keep a fill target and a spare node.
  • Forgetting non-HDFS space. Shuffle-heavy jobs fill YARN local directories on the same disks and fail with disk-full errors while HDFS still looks half empty.
  • Ignoring file shape. Byte forecasts look fine while the NameNode slides into long garbage collection pauses from small files.
  • Mismatched memory and vcores. Containers ask for 16 GB each on nodes with 4 GB per vcore, so half the cores sit idle while the scheduler reports the cluster as full.
  • Erasure coding hot data. Small or frequently rewritten files under RS-6-3 add CPU and reconstruction traffic and gain little storage. Keep EC for large, cold files.
  • Growth applied to bytes only. Storage keeps up, and the batch window quietly stretches past its deadline.

Operating the plan

A capacity plan is a forecast, so compare it against reality monthly. Track HDFS used and remaining, DataNode volume fill skew, NameNode files and blocks, heap after garbage collection, YARN pending containers and the finish time of critical jobs. Re-run the calculator with measured ingest and growth each quarter, and order hardware when the forecast crosses the fill target within your procurement lead time, not when the cluster is already full. Adding nodes is also a data movement event: run the balancer with a bandwidth cap after each expansion so new disks take a share of the existing data.

What to do next

  1. List every dataset with ingest per day, retention, redundancy policy and average file size; measure file size from the namespace rather than guessing.
  2. Measure vcore-hours for the jobs that have deadlines and write down each window.
  3. Run the calculator with your node shape and an 18 to 24 month horizon.
  4. Identify which dimension sets the node count, and test a second node shape that wastes less of the others.
  5. Check recovery time for one failed node and the fill after that failure; both should stay within target.
  6. Schedule a quarterly re-run against measured growth, and set alerts on fill, NameNode object count and job finish times.
Key takeaway: Hadoop capacity is four separate budgets: physical storage after redundancy and reserved space, NameNode objects driven by file size, YARN vcores and memory driven by job demand, and the bandwidth needed to re-protect a failed node. Size each from its own inputs, take the largest node count plus a spare, then use the others as checks and change the node shape if one resource is mostly wasted. Re-run the calculation against measured growth every quarter.