A Hadoop cluster rarely fails all at once. It degrades: a DataNode loses a disk, the NameNode spends more time in garbage collection, a queue starts holding applications for twenty minutes instead of twenty seconds. Each of those is visible in a number the daemons already publish, minutes or hours before a user files a ticket. Monitoring Hadoop well is mostly a matter of knowing where those numbers live, which few of the several hundred matter, and what each one means when it moves.
This page builds that knowledge from the bottom. It explains how the metrics2 framework exposes every daemon's state, how to read it by hand from the /jmx servlet, which NameNode, DataNode and YARN signals deserve alerts, how to get them into a time-series store, and how to diagnose a slow cluster from them. It ends with a small poller you can run today and a checklist. The internals behind the signals are covered in the NameNode deep dive and the YARN ResourceManager; this page is about watching them.
How Hadoop exposes its health
Every Hadoop daemon, NameNode, DataNode, ResourceManager, NodeManager, JournalNode and the rest, registers metric sources with the metrics2 framework. A source is a named group of counters, gauges and rates: the NameNode has FSNamesystem for namespace and block state, RpcActivityForPort8020 for its client RPC server, JvmMetrics for the JVM, and more. Each source is also registered as a JMX MBean in the Hadoop domain, so the same numbers are reachable three ways.
- The /jmx servlet. Each daemon's embedded web server serves its MBeans as JSON. In Hadoop 3 the default HTTP ports are 9870 for the NameNode, 9864 for DataNodes, 8088 for the ResourceManager and 8042 for NodeManagers. No agent is required, which makes it the easiest place to start.
- metrics2 sinks.
hadoop-metrics2.propertiescan push every source on a period to a sink: a file, Graphite, StatsD and others ship with Hadoop. Sinks push; nothing needs to scrape. - Remote JMX. Any JMX client can read the beans, and that is how the Prometheus jmx_exporter Java agent works: it runs inside the daemon JVM and turns MBeans into a scrape page.
Two side channels complete the picture. Daemon logs carry events that metrics only count, such as the JvmPauseMonitor line Detected pause in JVM or host machine (eg GC), which reports a stall the whole process experienced. And the CLI reports, hdfs dfsadmin -report, hdfs fsck and yarn node -list -all, give per-node detail when an aggregate number tells you something is wrong but not where.
Reading /jmx by hand
Before building dashboards, read the raw beans once. It teaches you the names, and it is what you will fall back to when the monitoring stack itself is broken. The qry parameter filters to one bean, which keeps the response small:
# Namespace and block state from the NameNode
curl -s 'http://nn1:9870/jmx?qry=Hadoop:service=NameNode,name=FSNamesystem'
# Client RPC server: queue time, processing time, call queue length
curl -s 'http://nn1:9870/jmx?qry=Hadoop:service=NameNode,name=RpcActivityForPort8020'
# JVM: heap, GC count and time, blocked threads
curl -s 'http://nn1:9870/jmx?qry=Hadoop:service=NameNode,name=JvmMetrics'
# YARN cluster and root queue
curl -s 'http://rm1:8088/jmx?qry=Hadoop:service=ResourceManager,name=ClusterMetrics'
curl -s 'http://rm1:8088/jmx?qry=Hadoop:service=ResourceManager,name=QueueMetrics,q0=root'The response is a JSON object with a beans array; each bean is a flat map of attribute names to values. The RPC bean name includes the port, so a NameNode that also runs a separate service RPC port has a second bean for it. If the cluster uses Kerberos with SPNEGO on the web endpoints, use curl --negotiate -u : with a valid ticket, and expect your scraper to need the same.
Check every metric name in this page against your own /jmx output before you alert on it. Names occasionally differ between the reference table and what a version emits; for example the Hadoop Metrics reference lists the RPC processing average as RpcProcessingAvgTime, while the rate it belongs to is RpcProcessingTime. The bean is the authority.
The NameNode signals that matter
The NameNode is the single point every HDFS client touches, so most HDFS monitoring is NameNode monitoring. Group its metrics by the question they answer.
| Question | Metric | Bean | Alert when |
|---|---|---|---|
| Is data at risk? | MissingBlocks | FSNamesystem | Above 0. Every replica of a block is gone; page. |
| Is redundancy degraded? | UnderReplicatedBlocks | FSNamesystem | Rising for more than about an hour after a node loss |
| Are replicas corrupt? | CorruptBlocks | FSNamesystem | Above 0 and not falling |
| Are DataNodes alive? | NumDeadDataNodes, NumStaleDataNodes | FSNamesystem | Dead above your tolerance; stale rising |
| Is there space? | CapacityRemaining / CapacityTotal | FSNamesystem | Below 20 percent free; page below 10 |
| Is the namespace growing? | FilesTotal, BlocksTotal | FSNamesystem | Trend against heap sizing, not a threshold |
| Are clients waiting? | RpcQueueTimeAvgTime, CallQueueLength | RpcActivityForPort8020 | Queue time above your baseline for 10 minutes |
| Is the JVM healthy? | GcTimeMillis, GcNumWarnThresholdExceeded | JvmMetrics | GC time per minute rising; any warn threshold events |
| Is checkpointing working? | TransactionsSinceLastCheckpoint, LastCheckpointTime | FSNamesystem | Checkpoint older than a few multiples of the period |
| Is HA in the right state? | tag.HAState | FSNamesystem | Not exactly one active NameNode |
Three of these deserve explanation. Missing blocks is the one HDFS metric that almost always means real loss; under-replication after a node failure is normal and self-healing, missing is not. RPC queue time is the earliest warning of an overloaded NameNode: requests wait in the call queue before a handler takes them, so queue time rises before processing time does, and it rises for every client at once. Checkpoint age catches a silent failure: if the standby (or secondary) stops checkpointing, nothing breaks until the next restart, which then has to replay an enormous edit log and can take hours. MillisSinceLastLoadedEdits on the standby says the same thing from the other side.
The HA pair needs its own check: exactly one NameNode reports active. hdfs haadmin -getAllServiceState shows it from the command line, and the failover mechanics are in NameNode high availability.
DataNodes and disks
Individual DataNodes are cattle; HDFS is designed to lose them. Alert on the aggregate in the NameNode, and use per-node metrics for diagnosis. The NameNode marks a DataNode stale after 30 seconds without a heartbeat by default (dfs.namenode.stale.datanode.interval), and dead after about ten and a half minutes with default settings, when it begins re-replicating that node's blocks. A burst of dead nodes therefore shows up first as stale nodes, then as an under-replication spike and heavy network use.
Disk failures are the commonest DataNode fault. With the default dfs.datanode.failed.volumes.tolerated of 0, one failed disk takes the whole DataNode offline; many operators raise it on dense nodes. Either way, watch VolumeFailuresTotal on the NameNode and VolumeFailures on each DataNode, and treat any increase as a hardware ticket. For slow rather than dead disks, the per-DataNode FsyncNanos and TotalWriteTime counters and the DataNode heartbeat time HeartbeatsAvgTime are the places to look: one slow disk in a write pipeline slows every client writing through that node.
YARN: the ResourceManager, nodes and queues
YARN monitoring answers two questions: can the cluster run work, and is work waiting longer than it should. The first comes from ClusterMetrics on the ResourceManager; the second from QueueMetrics, which exists per queue (q0=root, q1=default and so on).
- Node health:
NumActiveNMs,NumLostNMsandNumUnhealthyNMs. An unhealthy NodeManager is usually the disk health checker refusing a full or failed local directory, which quietly removes capacity. - Waiting work:
AppsPending,PendingMB,PendingVCoresandPendingContainers. Pending is normal for seconds; pending for many minutes whileAvailableMBis high points at a placement or queue limit problem rather than a lack of capacity. - Use:
AllocatedMB,AllocatedContainersandReservedMB. Large reservations mean the scheduler is holding space for big containers that do not fit yet. - Outcomes:
AppsFailedandAppsKilledare counters; alert on their rate, not their value.
Like the NameNode, the ResourceManager runs as an HA pair; yarn rmadmin -getAllServiceState shows which is active. Remember that a standby ResourceManager does not report the live cluster metrics, so scrape both and keep only the active one's values, or your dashboards will show the cluster dropping to zero on every failover.
Getting metrics into a time-series store
Any of the three paths works; the choice is operational.
| Path | How | Strengths | Weaknesses |
|---|---|---|---|
| jmx_exporter agent | Java agent on each daemon, scraped by Prometheus | Works on any version, full control over names | A YAML rule file per daemon type, agent upgrades |
| /prom endpoint | Set hadoop.prometheus.endpoint.enabled=true (Hadoop 3.3.0 and later, default false) | No agent, native metrics2 names | Only on versions that have it; confirm on your build |
| metrics2 sink | Configure a sink in hadoop-metrics2.properties | Push, so no scrape config | Fewer sink types, restart to change |
| /jmx poller | Script reading JSON over HTTP | Zero install, good for alerts on a few metrics | You own the code and the history |
For the agent path, add it to the daemon options, for example in hadoop-env.sh:
export HDFS_NAMENODE_OPTS="$HDFS_NAMENODE_OPTS \
-javaagent:/opt/jmx/jmx_prometheus_javaagent.jar=7071:/etc/hadoop/jmx-namenode.yaml"and a sink, if you prefer to push, looks like this:
# hadoop-metrics2.properties
*.period=10
namenode.sink.file.class=org.apache.hadoop.metrics2.sink.FileSink
namenode.sink.file.filename=namenode-metrics.outWhichever you pick, keep the scrape interval at 15 to 60 seconds, label every series with cluster and role, and scrape both members of each HA pair. Exporter rule files decide the final metric names, so write your alert rules against the names your exporter produces, not the ones on this page.
A poller with alert rules you can run
When you need alerts this week and a full stack next quarter, a small poller over /jmx covers the paging-worthy signals. It reads each bean once, evaluates rules, and prints alerts you can wire into cron and email or a chat webhook.
import json, sys, urllib.request
NN = "http://nn1:9870"
RM = "http://rm1:8088"
def bean(base, qry):
with urllib.request.urlopen(f"{base}/jmx?qry={qry}", timeout=10) as r:
beans = json.load(r)["beans"]
if not beans:
raise RuntimeError(f"no bean {qry} at {base}")
return beans[0]
def check():
fs = bean(NN, "Hadoop:service=NameNode,name=FSNamesystem")
rpc = bean(NN, "Hadoop:service=NameNode,name=RpcActivityForPort8020")
cm = bean(RM, "Hadoop:service=ResourceManager,name=ClusterMetrics")
q = bean(RM, "Hadoop:service=ResourceManager,name=QueueMetrics,q0=root")
free = fs["CapacityRemaining"] / max(fs["CapacityTotal"], 1)
rules = [
("page", fs["MissingBlocks"] > 0, f"missing blocks: {fs['MissingBlocks']}"),
("page", free < 0.10, f"HDFS free {free:.1%}"),
("warn", free < 0.20, f"HDFS free {free:.1%}"),
("warn", fs["NumDeadDataNodes"] > 2, f"dead DataNodes: {fs['NumDeadDataNodes']}"),
("warn", rpc["RpcQueueTimeAvgTime"] > 50, f"RPC queue time {rpc['RpcQueueTimeAvgTime']:.0f} ms"),
("warn", cm["NumLostNMs"] + cm["NumUnhealthyNMs"] > 2, "NodeManagers lost or unhealthy"),
("warn", q["AppsPending"] > 20 and q["AvailableMB"] > q["PendingMB"],
"apps pending while memory is free: check queue limits and placement"),
]
return [(sev, msg) for sev, fired, msg in rules if fired]
if __name__ == "__main__":
try:
alerts = check()
except Exception as e: # an unreachable daemon is itself an alert
alerts = [("page", f"monitoring failed: {e}")]
for sev, msg in alerts:
print(f"{sev.upper()}: {msg}")
sys.exit(2 if any(s == "page" for s, _ in alerts) else 1 if alerts else 0)Three details matter more than they look. A daemon that does not answer is an alert, not a skipped check, because silence is the most common way monitoring fails. The thresholds are placeholders: measure a normal week first and set them from that baseline, especially RPC queue time, which varies widely between clusters. And in an HA cluster, point NN and RM at whichever node is active, or query both and use the one whose state says active.
Worked example: a slow Monday morning
At 09:10 users report that Hive queries take minutes to start. The ResourceManager looks fine: plenty of available memory, a few pending applications. On the NameNode, RpcQueueTimeAvgTime has risen from a baseline of about 2 ms to 400 ms, and CallQueueLength sits near its limit. Processing time per call has barely moved.
That shape says the handlers are not doing slow work; they are not running at all for stretches. The JvmMetrics bean confirms it: GcTimeMillis is climbing several seconds per minute, and the NameNode log shows JvmPauseMonitor lines reporting pauses of 3 to 8 seconds. FilesTotal has grown 30 percent in a month, because a new ingestion job writes one small file per event.
The immediate fix is to stop the job and compact its output, which reduces the object count and with it the old-generation pressure. The durable fix has two parts: size the heap for the namespace you actually have, which Hadoop capacity planning works through, and add FilesTotal growth to the weekly capacity review so the next small-file job is caught by a trend line rather than by users.
How monitoring itself fails
- Scraping only one HA member: after a failover every graph drops to zero or freezes, and alerts that compare against zero fire or go silent.
- Thresholds on counters:
AppsFailedonly rises; alert on its rate. - Alerting on under-replication: it spikes on every node loss and recovers on its own. Alert on missing blocks and on under-replication that does not fall.
- No absence alert: a dead exporter looks exactly like a healthy, quiet cluster.
- Kerberos surprises: enabling SPNEGO on web UIs breaks an unauthenticated scraper silently, with 401s in a log nobody reads.
- Per-node paging: paging on every DataNode disk failure at three in the morning trains people to ignore pages. Ticket it, and page on the aggregate.
What to do next
- Curl
/jmxon your active NameNode and ResourceManager and confirm every metric name you plan to alert on. - Page on missing blocks, on low HDFS free space and on an unreachable daemon; send everything else to a ticket queue.
- Record a normal week of
RpcQueueTimeAvgTime, GC time andAppsPending, then set thresholds from that baseline. - Pick one collection path, the jmx_exporter agent, the /prom endpoint on 3.3.0 or later, or a sink, and scrape both members of each HA pair.
- Alert on checkpoint age and on any JvmPauseMonitor line over a few seconds.
- Add
FilesTotal,BlocksTotaland capacity trends to a weekly review alongside your capacity plan.