skip to content

How would you size a new Hadoop cluster's DataNode storage, NameNode heap and worker node shape?

level: principalimportance: should knowfreq 34%

answer

  1. one number is never enough
  2. bytes and objects grow separately
  3. never plan to fill the disks
  4. the metadata cost is per file, not per gigabyte
  5. measure your own fsimage

basics

~20 s

Size three budgets separately: raw storage as logical data times the replication factor plus headroom for non-DFS use and node rebuilds; NameNode heap from the count of files, directories and blocks rather than bytes; and worker memory, cores and disks from the workload's container profile.

solid answer

~50 s

Treat storage, metadata and compute as three independent budgets. **Storage:** raw capacity is logical data × replication (3 by default, or roughly 1.5× with erasure coding on cold data), then divided by a target fill of about 70–75% — you need room for non-DFS use (`dfs.datanode.du.reserved`), YARN local and log dirs, and the re-replication burst when a node dies. **Metadata:** NameNode heap tracks the number of files, directories and blocks, not bytes; derive the object count from your own `fsimage` with `hdfs oiv` and size heap against measured growth, remembering that huge heaps mean long GC pauses. **Compute:** pick a node shape from the container profile — memory per container × concurrency, cores, and enough JBOD disks (no RAID on `dfs.datanode.data.dir`) for the scan throughput jobs need — and set `yarn.nodemanager.resource.memory-mb` to what is left after the DataNode, NodeManager and OS. Then validate against a real workload and re-plan on measured growth.

code

properties · 6 lines
properties
# key=value shorthand for the XML site properties
dfs.datanode.data.dir=/data/1/dn,/data/2/dn,/data/3/dn,/data/4/dn
dfs.datanode.du.reserved=107374182400
yarn.nodemanager.local-dirs=/data/1/nm,/data/2/nm,/data/3/nm,/data/4/nm
yarn.nodemanager.resource.memory-mb=98304
yarn.nodemanager.resource.cpu-vcores=24

go deeper

for a junior

Know the basic multiplier: with replication factor 3, a terabyte of data occupies three terabytes of raw disk, and you never plan to fill the disks completely.

for a middle

Explain why storage and metadata are separate budgets, that NameNode heap follows object count rather than bytes, and what has to be subtracted from a node before YARN gets its share.

for a senior

Show operating experience — reserving space for re-replication after a node loss, spotting when disk count rather than capacity caps job throughput, watching NameNode heap and GC, and validating a plan against a real workload.

for a principal

Own the whole model: growth and retention policy, replication versus erasure coding by data temperature, node shapes for a mixed workload, multi-tenant queue capacity, and the trigger metrics that start a procurement cycle before anything is full.

## Three budgets, not one number The common mistake is to size a Hadoop cluster from a single figure — "we have 300 TB". Storage, metadata and compute scale on different inputs and fail in different ways, so plan each one. ## Budget 1: raw storage Start with logical data: current volume plus projected growth over the planning horizon, after compression and after applying a retention policy (deciding what you *delete* is a capacity lever people forget). Multiply by the durability factor. Three-way replication means 3× raw for every logical byte. Hadoop 3's erasure coding is the alternative: an RS(6,3) scheme stores six data and three parity cells, giving roughly 1.5× overhead instead of 3× and tolerating three losses. It is not free — reconstruction after a failure is CPU- and network-heavy, and EC files carry functional restrictions — so it belongs on warm/cold data, not on the hot working set. Then add headroom. Do not plan to fill DataNodes: block placement, the balancer and the temporary spike of re-replication when a node fails all need slack, and jobs need space for output before old data is dropped. Planning to about 70–75% of raw is the usual discipline. Separately reserve space that HDFS must never claim with `dfs.datanode.du.reserved`, and remember that YARN's `yarn.nodemanager.local-dirs` (shuffle spill and container work dirs) and log dirs live on the same disks. So 300 TB logical at replication 3 is 900 TB of replicas, and you buy comfortably more than that — over a petabyte raw — not 900 TB exactly. ## Budget 2: NameNode metadata The NameNode holds the namespace in memory, and its heap is driven by **object count** — files, directories and blocks — not by how many bytes those files contain. A petabyte in a few thousand large files is trivial for the NameNode; the same petabyte in a hundred million tiny files may not fit in heap at all. Plan it empirically: run `hdfs oiv` over a representative `fsimage` to count objects, measure the heap your current NameNode actually uses, and project forward on the growth rate of *object count*. Rules of thumb about bytes per object circulate, but the honest engineering answer is to measure your own namespace rather than trust one. Set heap in `hadoop-env.sh` (`HADOOP_HEAPSIZE_MAX` or `HDFS_NAMENODE_OPTS`), and be aware that very large heaps bring long GC pauses — a pause that outlives the ZooKeeper session timeout will trigger a spurious HA failover. If object growth is the binding constraint, the levers are compaction of small files, HDFS federation, or moving cold data elsewhere. ## Budget 3: worker node shape Work backwards from the workload's container profile: how much memory a typical task needs, how many you want running concurrently, and whether the jobs are CPU-bound or scan-bound. - **Memory.** Node RAM minus the OS, the DataNode JVM and the NodeManager JVM is what YARN may hand out; put that figure in `yarn.nodemanager.resource.memory-mb`. It is not derived from physical RAM for you, and setting it to the full machine size is how you get OOM-killed daemons. - **Cores.** `yarn.nodemanager.resource.cpu-vcores` should reflect usable cores, allowing for the daemons and for the fact that vcores are a scheduling abstraction, not a cgroup guarantee unless you enable CPU enforcement. - **Disks.** Many independent disks, JBOD, each its own directory in `dfs.datanode.data.dir`. Do not RAID DataNode disks: HDFS already replicates across nodes, and RAID costs you both capacity and parallel spindle throughput. Sequential scan bandwidth per node is often the real ceiling on job runtime, so disk count matters as much as capacity per disk. - **Balance.** A node with huge capacity per disk but few spindles stores a lot and reads it slowly; that shows up as jobs that are slower on a "bigger" cluster. ## Master nodes and the network Master nodes are sized for reliability, not throughput: NameNode heap plus redundant metadata directories, three JournalNodes on separate machines, a ZooKeeper ensemble that is not competing with heavy I/O. The network matters too — shuffle-heavy workloads push a lot of traffic across rack uplinks, and an oversubscribed top-of-rack switch will cap the cluster long before the disks do. ## Multi-tenancy and validation If several teams share the cluster, capacity planning includes queue sizing in the scheduler and an explicit policy on what happens when the cluster is full: whose jobs wait. Decide that before the cluster is full, not during the first incident. Finally, validate rather than trust the spreadsheet. Run a representative workload on a small number of the proposed nodes, measure achieved scan throughput, container density and job runtime, and extrapolate from that. Then instrument the four early-warning signals — DataNode fill percentage, NameNode heap used and object count, queue wait time, and rack uplink utilisation — and re-plan when any of them trends toward its limit rather than when it hits it.

  • Why is NameNode heap the constraint that usually bites first on an ageing cluster?
    Because it scales with object count, which grows with ingestion frequency rather than data volume. Hourly partitions and small output files multiply files, directories and blocks while total bytes rise slowly. Disks can be added incrementally; heap cannot exceed what one JVM manages with acceptable GC pauses. That is why compaction, larger partitions and federation are planning levers, not just tidiness.
  • When does erasure coding change the storage plan, and what does it cost?
    An RS(6,3) layout stores roughly 1.5× the logical size instead of 3×, so a cold-data tier gets dramatically cheaper. The costs are reconstruction expense — rebuilding a lost cell reads from many nodes and burns CPU and network — plus functional restrictions and a different failure profile. Apply it to warm and cold data with low read concurrency, and keep the hot working set replicated.
  • Why should DataNode disks be JBOD rather than RAID?
    HDFS already provides redundancy across nodes, so RAID adds a second, redundant durability layer and gives up usable capacity for it. It also hides individual spindles behind one device, reducing the parallel I/O the DataNode can drive. List each disk separately in `dfs.datanode.data.dir` and let a failed disk be tolerated at the HDFS level. NameNode metadata directories are the exception where redundant storage is appropriate.
  • How would you plan for the capacity a node failure consumes?
    When a DataNode is lost the NameNode re-replicates every block it held, so the surviving nodes must absorb that data and the network must carry it. Reserve enough free space that losing your largest failure domain — a node, or a rack if you plan at that level — does not push the remainder past a safe fill level, and expect a period of elevated background traffic that competes with jobs.

saying these in an interview costs you the question

  • Sizes raw storage as the logical data volume with no replication factor
  • Says NameNode heap scales with terabytes stored
  • Plans to run DataNodes at 95% full
  • Puts DataNode disks behind RAID for safety
  • Sets yarn.nodemanager.resource.memory-mb to the machine's full RAM

context