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.

Advertisement

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.properties can 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.
Where Hadoop health comes from, and where it should goNameNodeHTTP 9870, RPC 8020DataNodesHTTP 9864ResourceManagerHTTP 8088NodeManagersHTTP 8042metrics2 sourcesper-daemon beans/jmx servletJSON over HTTPsame beansmetrics2 sinksfile, Graphite, StatsDScraperjmx_exporter or /promPollercurl /jmx, cronpushpullpullTSDB and alertingrules, dashboards, pagesSide channels: daemon logs, CLI reportsdfsadmin -report, fsck, yarn node -list, JvmPauseMonitor lineslog alerts
Every daemon publishes the same metrics2 sources as JMX beans. You either push them through a sink or pull them over HTTP; logs and CLI reports fill the gaps the metrics leave.

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.

Advertisement

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.

QuestionMetricBeanAlert when
Is data at risk?MissingBlocksFSNamesystemAbove 0. Every replica of a block is gone; page.
Is redundancy degraded?UnderReplicatedBlocksFSNamesystemRising for more than about an hour after a node loss
Are replicas corrupt?CorruptBlocksFSNamesystemAbove 0 and not falling
Are DataNodes alive?NumDeadDataNodes, NumStaleDataNodesFSNamesystemDead above your tolerance; stale rising
Is there space?CapacityRemaining / CapacityTotalFSNamesystemBelow 20 percent free; page below 10
Is the namespace growing?FilesTotal, BlocksTotalFSNamesystemTrend against heap sizing, not a threshold
Are clients waiting?RpcQueueTimeAvgTime, CallQueueLengthRpcActivityForPort8020Queue time above your baseline for 10 minutes
Is the JVM healthy?GcTimeMillis, GcNumWarnThresholdExceededJvmMetricsGC time per minute rising; any warn threshold events
Is checkpointing working?TransactionsSinceLastCheckpoint, LastCheckpointTimeFSNamesystemCheckpoint older than a few multiples of the period
Is HA in the right state?tag.HAStateFSNamesystemNot 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, NumLostNMs and NumUnhealthyNMs. An unhealthy NodeManager is usually the disk health checker refusing a full or failed local directory, which quietly removes capacity.
  • Waiting work: AppsPending, PendingMB, PendingVCores and PendingContainers. Pending is normal for seconds; pending for many minutes while AvailableMB is high points at a placement or queue limit problem rather than a lack of capacity.
  • Use: AllocatedMB, AllocatedContainers and ReservedMB. Large reservations mean the scheduler is holding space for big containers that do not fit yet.
  • Outcomes: AppsFailed and AppsKilled are 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.

PathHowStrengthsWeaknesses
jmx_exporter agentJava agent on each daemon, scraped by PrometheusWorks on any version, full control over namesA YAML rule file per daemon type, agent upgrades
/prom endpointSet hadoop.prometheus.endpoint.enabled=true (Hadoop 3.3.0 and later, default false)No agent, native metrics2 namesOnly on versions that have it; confirm on your build
metrics2 sinkConfigure a sink in hadoop-metrics2.propertiesPush, so no scrape configFewer sink types, restart to change
/jmx pollerScript reading JSON over HTTPZero install, good for alerts on a few metricsYou 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.out

Whichever 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: AppsFailed only 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

  1. Curl /jmx on your active NameNode and ResourceManager and confirm every metric name you plan to alert on.
  2. Page on missing blocks, on low HDFS free space and on an unreachable daemon; send everything else to a ticket queue.
  3. Record a normal week of RpcQueueTimeAvgTime, GC time and AppsPending, then set thresholds from that baseline.
  4. 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.
  5. Alert on checkpoint age and on any JvmPauseMonitor line over a few seconds.
  6. Add FilesTotal, BlocksTotal and capacity trends to a weekly review alongside your capacity plan.
Key takeaway: Every Hadoop daemon publishes its state as metrics2 sources, readable as JSON from the /jmx servlet or exported through a sink or the jmx_exporter agent. Read the beans by hand first so you know the names. Page on the few signals that mean real harm: missing blocks, low HDFS space, an unreachable or not-active NameNode or ResourceManager. Watch RPC queue time, GC time and checkpoint age as early warnings, use queue metrics to tell lack of capacity from placement limits, scrape both HA members, and alert on the monitoring going silent.