A Hadoop cluster is a stack of layers: disks and network at the bottom, the HDFS data path on top of them, the NameNode's metadata service beside it, then YARN scheduling and the MapReduce or Spark shuffle on top. When a job is slow, any of those layers can be at fault, and a benchmark is only useful if you know which layer it exercises. Running TeraSort and reporting one number tells you that the whole stack was slow or fast that afternoon. It does not tell you why.
This article treats benchmarking as measurement, not ritual. It maps the tools that ship with Hadoop to the layers they stress, shows how to run and read TestDFSIO, the TeraGen, TeraSort and TeraValidate trio and the NameNode throughput benchmark, works through an acceptance test from first principles so you know what number to expect before you run anything, and ends with the mistakes that make benchmark results lie. For tuning the knobs a benchmark points at, see HDFS performance tuning.
The benchmark stack
Which layer does each benchmark measure?
Every benchmark answers one question about one layer, and its result is capped by every layer underneath. TestDFSIO writes and reads large files through HDFS, so it measures the DataNode data path, replication pipeline and disks, but also the network, because a replicated write crosses the wire twice. TeraSort adds map compute, sorting, the shuffle and reducer writes on top of all of that. The NameNode benchmarks barely touch disks at all: they measure how many metadata operations per second the NameNode can handle, which is a CPU, lock and RPC problem.
That gives the rule for diagnosing a slow cluster: measure bottom up. First establish raw disk throughput per node with a tool such as fio and raw network throughput between nodes with iperf3, outside Hadoop. Then run TestDFSIO and compare it against what the hardware should deliver. Only when HDFS is close to its ceiling does a TeraSort result tell you anything about YARN or shuffle settings. If you skip the bottom layers, a bad NIC, a disk in a degraded RAID controller or a misconfigured MTU shows up as a mysterious slow sort.
TestDFSIO: the HDFS data path
TestDFSIO ships in the MapReduce job-client tests jar. It runs one map task per file; each map writes or reads its file sequentially and records bytes and time, and a single reduce aggregates the per-task figures. Control files go under io_control in the base directory, which defaults to /benchmarks/TestDFSIO and can be changed with -Dtest.build.data=.... The usage, as printed by the tool, is:
TestDFSIO [genericOptions] -read [-random | -backward | -skip [-skipSize Size]]
| -write | -append | -truncate | -clean
[-compression codecClassName] [-nrFiles N] [-size Size[B|KB|MB|GB|TB]]
[-resFile resultFileName] [-bufferSize Bytes]
[-storagePolicy storagePolicyName] [-erasureCodePolicy erasureCodePolicyName]A typical write-then-read run with 64 files of 4 GB each looks like this. The size unit defaults to MB when no suffix is given; older guides use -fileSize, which current versions replace with -size.
JAR=$(ls $HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-client-jobclient-*-tests.jar)
hadoop jar $JAR TestDFSIO -write -nrFiles 64 -size 4GB -resFile /tmp/dfsio_write.log
hadoop jar $JAR TestDFSIO -read -nrFiles 64 -size 4GB -resFile /tmp/dfsio_read.log
hadoop jar $JAR TestDFSIO -cleanEach run appends a block like this to the result file (values here are placeholders, not measurements):
----- TestDFSIO ----- : write
Date & time: ...
Number of files: 64
Total MBytes processed: 262144
Throughput mb/sec: ...
Average IO rate mb/sec: ...
IO rate std deviation: ...
Test exec time sec: ...Reading these lines correctly is most of the skill. Throughput mb/sec is total megabytes divided by the sum of every task's I/O time, so it is a per-task figure weighted by bytes, not cluster bandwidth. Average IO rate is the plain mean of each task's own rate, so a few fast tasks pull it above throughput. IO rate std deviation is the spread across tasks: a large value means some nodes or disks are much slower than others, which is often the most useful finding in the whole run. Test exec time is wall-clock time for the job, including scheduling. A rough cluster-wide figure is total MB divided by exec time, or throughput multiplied by the number of maps that actually ran concurrently; quote which one you used.
TeraGen, TeraSort and TeraValidate
The Tera trio lives in the MapReduce examples jar. TeraGen writes synthetic 100-byte rows with 10-byte keys, so the row count is the data size divided by 100. TeraSort samples the input to build a total-order partitioner, so each reducer receives one contiguous key range and the concatenated outputs are globally sorted. TeraValidate reads the sorted output and checks that keys are in order within and across files, writing a report directory.
EX=$(ls $HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-examples-*.jar)
# 1 TB = 10,000,000,000 rows of 100 bytes
hadoop jar $EX teragen -Dmapreduce.job.maps=400 10000000000 /bench/tera-in
hadoop jar $EX terasort -Dmapreduce.job.reduces=200 /bench/tera-in /bench/tera-out
hadoop jar $EX teravalidate /bench/tera-out /bench/tera-reportTeraGen is a map-only write benchmark and behaves much like a TestDFSIO write. TeraSort is the interesting one: it reads the full dataset, sorts it in map-side buffers, spills, shuffles every byte across the network once and writes the result. That makes it sensitive to mapreduce.task.io.sort.mb, the reducer count, compression of intermediate map output, and the shuffle service, which is covered in the shuffle and sort article. Record the replication factor of the output you used, because writing three replicas instead of one changes the write phase dramatically and published TeraSort figures are not always explicit about it.
NameNode benchmarks and workload replay
Large clusters often fail on metadata long before bandwidth. Millions of small files, aggressive listing from query engines and bursts of creates from many jobs all land on one NameNode, described in the NameNode article. NNThroughputBenchmark measures that service directly. Its mode follows fs.defaultFS, which -fs overrides. If the scheme is unset or file, it starts a NameNode in the same process (standalone mode) and its threads call NameNode methods directly, which isolates NameNode CPU and locking from RPC and network effects. Otherwise it runs against the remote NameNode over client RPCs. On a cluster gateway, where fs.defaultFS points at production, that means the default is remote mode, so pass -fs file:/// explicitly when you want standalone.
hadoop org.apache.hadoop.hdfs.server.namenode.NNThroughputBenchmark \
-fs file:/// -op create -threads 16 -files 200000 -filesPerDir 1000 -close
hadoop org.apache.hadoop.hdfs.server.namenode.NNThroughputBenchmark \
-fs hdfs://nn1.example.com:8020 -op fileStatus -threads 32 -files 200000The supported operations are all, create, mkdirs, open, delete, fileStatus, rename, blockReport, replication and clean; -op must be the first option. It reports elapsed time, operations per second and average operation time. Only run it against a shared NameNode in a maintenance window, and use -keepResults deliberately, since without it the tool cleans up the namespace it created.
Other tools cover neighbouring questions. nnbench drives NameNode operations from a MapReduce job, so the load arrives over RPC from many hosts. mrbench runs a tiny job many times to measure per-job overhead, which matters for workloads made of short queries. SliveTest generates a mixed metadata workload. Dynamometer replays a real NameNode's audit log against an emulated cluster, which is the closest you can get to testing a NameNode upgrade on production traffic. YARN's Scheduler Load Simulator replays application traces against the scheduler without real containers. Check each tool's own help output before relying on a flag.
Worked example: accepting a 20-node cluster
Suppose you are accepting a new 20-node cluster. Each worker has 12 data disks and a 25 Gbit/s NIC. Before running anything, derive the expected ceiling. All numbers below are arithmetic from assumed hardware figures, not measurements.
- Disk ceiling. If fio shows about 180 MB/s of large sequential writes per disk, each node can write about 12 x 180 = 2,160 MB/s, and the cluster about 43,200 MB/s of physical writes.
- Replication. With replication 3, every logical megabyte becomes three physical megabytes. The logical write ceiling from disks is therefore about 14,400 MB/s for the cluster, or 720 MB/s per node.
- Network. The write pipeline sends each block to two remote DataNodes, so each node both sends and receives about two logical megabytes for every one it writes. At 25 Gbit/s, roughly 3,000 MB/s usable per direction, the network cap is about 1,500 MB/s per node, above the disk cap. Disks bind.
- Expectation. Run one file per data disk,
-nrFiles 240, with enough containers for all 240 writers to run at once. Each block replica streams to one disk, so 240 writers create 720 replica streams, about three per disk, and perfect balance gives about 14,400 / 240 = 60 MB/s per task. Expect the TestDFSIO throughput figure somewhat below that, and a total-MB-over-exec-time figure lower still because of task startup and stragglers.
Now the measurement has meaning. If write throughput per task comes back near 55 MB/s with a small standard deviation, HDFS is healthy. If it comes back at 25 MB/s with a huge deviation, look for slow nodes: sort the per-task records by rate, map tasks to hosts in the job history, and you will usually find one or two nodes with a bad disk, a half-duplex link or a missing jumbo-frame setting. Only after TestDFSIO is close to its ceiling should you move on to TeraSort and NameNode tests. The same arithmetic, run in reverse, is the core of capacity planning.
Collect results in a machine-readable form so runs can be compared over time. This short parser turns a TestDFSIO result log into one record per run:
import re, json, sys
FIELDS = {
"Number of files": "files",
"Total MBytes processed": "total_mb",
"Throughput mb/sec": "throughput_mb_s",
"Average IO rate mb/sec": "avg_io_rate_mb_s",
"IO rate std deviation": "io_rate_stddev",
"Test exec time sec": "exec_time_s",
}
def parse(path):
runs, cur = [], None
for line in open(path, encoding="utf-8"):
m = re.match(r"----- TestDFSIO ----- : (\w+)", line)
if m:
cur = {"op": m.group(1)}
runs.append(cur)
continue
key, _, value = line.partition(":")
if cur is not None and key.strip() in FIELDS:
cur[FIELDS[key.strip()]] = float(value.replace(",", "").strip())
for r in runs:
r["cluster_mb_s"] = r["total_mb"] / r["exec_time_s"]
return runs
print(json.dumps(parse(sys.argv[1]), indent=2))
Running a benchmark campaign
A benchmark campaign is an experiment, so run it like one. Quiesce the cluster or reserve a dedicated queue, because a concurrent ETL job silently halves your numbers. Snapshot the configuration: Hadoop version, hdfs-site.xml, mapred-site.xml, yarn-site.xml, JVM options, kernel and disk mount options. Do one warm-up run and discard it. Repeat each measurement at least three times and report the median and the range, not the best run. Change one variable per experiment, and keep the result logs next to the configuration snapshot in version control so that a regression three months later can be compared against a known baseline. Feed the steady-state numbers into your monitoring as reference lines, so production dashboards show how far normal load is from the measured ceiling.
Failure modes that make results lie
- Page cache reads. Reading files right after writing them often serves data from the operating system cache, producing read rates above what the disks can do. Use a dataset larger than cluster memory or drop caches on every node before the read run.
- Too few files. If
-nrFilesis smaller than the number of disks, most spindles sit idle and you measure one task's speed, not the cluster's. Use at least as many files as data disks, and confirm the tasks actually ran concurrently. - Comparing different shapes. Throughput from 16 files of 16 GB and from 256 files of 1 GB are not comparable; the concurrency is different. Keep the file count and size fixed across runs you intend to compare.
- Speculative execution. A speculative duplicate of a slow writer adds I/O the result does not account for. Disable speculation for benchmark jobs.
- Hidden replication or encoding changes. A directory with an erasure coding policy or a non-default replication factor changes the physical write volume. Record both, or set them explicitly with
-erasureCodePolicyand-D dfs.replication=3. - Forgetting cleanup. A 1 TB TeraGen left behind fills disks and skews the next balancer run. Clean up with
-cleanand by deleting the Tera directories with-skipTrash. - Treating TeraSort as the workload. It is a uniform-key, CPU-light sort. A real Hive or Spark workload with skewed joins and small files can behave entirely differently; prove it with workload replay before sizing on it.
Trade-offs
Synthetic benchmarks are cheap, repeatable and comparable across clusters, which makes them ideal for acceptance tests and regression checks after an upgrade or a configuration change. They are poor predictors of real job latency, because real workloads have skew, small files, mixed concurrency and metadata storms that no single synthetic tool reproduces. Replay tools such as Gridmix and Dynamometer close that gap at the cost of setup effort and the need to capture production traces. A sensible split is synthetic tests for every change to hardware or configuration, and replay before major version upgrades or when sizing a migration.
What to do next
- Measure raw disk and network throughput on every worker with fio and iperf3 and record them.
- Write down the expected HDFS write and read ceilings from those numbers and your replication or erasure coding policy.
- Run TestDFSIO write and read with at least one file per data disk, three times, on a quiet cluster, and parse the logs into a table.
- Investigate any run with a large IO rate standard deviation before moving up the stack.
- Run TeraGen, TeraSort and TeraValidate at a size larger than cluster memory, recording replication and reducer count.
- Run NNThroughputBenchmark for create, open and fileStatus in a maintenance window and keep the operations-per-second baseline.
- Store the configuration snapshot and results in version control and rerun the suite after every upgrade.