上一篇讲了安全。本篇讲 Hadoop 的可观测性——Metrics、JMX、日志聚合。这是 HA、安全与 Common 篇的最后一篇。

Hadoop 可观测性常被介绍成"看 Web UI 和日志"。这个描述对应了基本操作但漏了系统化设计。准确的说法是:Hadoop 提供三层可观测性——Metrics V2 是结构化指标系统(Source-Sink 架构,周期性采样暴露给外部采集)、JMX 是运行时 Java 指标暴露接口(Prometheus JMX Exporter 让 Hadoop 接入现代时序数据库)、YARN Log Aggregation 是 Task 日志的集中存储和查询。这三层共同构成 Hadoop 集群的运维可观测性基座。

本篇只抓一个问题:Metrics V2 的 Source-Sink 架构怎么工作、JMX 怎么接入 Prometheus 等现代监控栈、YARN Log Aggregation 的写入和查询流程、生产环境的可观测性最佳实践。

Metrics V2:Source-Sink 架构

Hadoop 2.x 引入了 Metrics V2 框架(替代 1.x 的 Metrics V1),是所有 Hadoop 组件(NameNode、DataNode、ResourceManager、NodeManager、MRAppMaster)的指标系统。

Metrics V2 的核心抽象是三个角色:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
MetricsSource:指标来源
- 每个 Hadoop 组件注册一个或多个 Source
- Source 实现 getMetrics(MetricsCollector, boolean) 方法
- 方法被周期性调用(默认 10 秒),返回当前指标的快照

MetricsSystem:中央收集器
- 单例,每个 JVM 一个
- 周期性(默认 10 秒)调用所有已注册 Source 的 getMetrics
- 把收集到的指标 fan out 给所有已注册 Sink

MetricsSink:指标下游
- 把指标推送到外部系统
- 内置:FileSink(写本地文件)、GangliaSink(推 Ganglia)、GraphiteSink(推 Graphite)
- Hadoop 3.x:PrometheusMetricsSink(暴露给 Prometheus 拉取)

Metrics V2 的配置在 hadoop-metrics2.properties 文件:

1
2
3
4
5
6
7
8
9
10
11
12
13
# 全局配置
*.period=10 # 默认采样周期 10 秒
*.rowsPerline=10

# 各组件的 Source 配置
namenode.sink.file.class=org.apache.hadoop.metrics2.sink.FileSink
namenode.sink.file.filename=nn-metrics.out

resourcemanager.sink.ganglia.class=org.apache.hadoop.metrics2.sink.ganglia.GangliaSink31
resourcemanager.sink.ganglia.servers=ganglia-host:8649

# Prometheus(Hadoop 3.x)
namenode.sink.prometheus.class=org.apache.hadoop.metrics2.sink.PrometheusMetricsSink

每个组件可以配置多个 Sink,指标 fan out 到多个下游。例如 NameNode 同时把指标写本地文件(备份)和推 Ganglia(可视化)。

Metrics V2 的设计动机是解耦——Source 只负责产生指标,Sink 只负责转发,MetricsSystem 负责周期性调度。这种解耦让新加 Sink(例如 Prometheus)不需要改 Source 代码。

JMX:运行时 Java 指标

除了 Metrics V2,Hadoop 还通过 JMX(Java Management Extensions)暴露运行时指标。JMX 是 Java 平台标准,所有 JVM 进程都可以暴露 MBean(Managed Bean)让外部工具查询。

Hadoop 的每个组件在 JVM 启动时注册大量 MBean:

1
2
3
4
5
6
Hadoop:service=NameNode,name=NameNodeInfo       ← NameNode 状态
Hadoop:service=NameNode,name=NameNodeMetrics ← NameNode 性能指标
Hadoop:service=NameNode,name=RpcMetrics ← RPC 指标
Hadoop:service=NameNode,name=JvmMetrics ← JVM 指标(heap, gc, threads)
Hadoop:service=NameNode,name=ugi ← User Group Information 指标
Hadoop:service=NameNode,name=RetryCache ← Retry Cache 指标

JMX MBean 的访问方式:

1
2
# 通过 Web UI 的 /jmx endpoint(HTTP)
curl http://namenode-host:9870/jmx?qry=Hadoop:service=NameNode,name=NameNodeMetrics

返回 JSON 格式的 MBean 属性:

1
2
3
4
5
6
7
8
9
10
11
12
13
{
"beans": [
{
"name": "Hadoop:service=NameNode,name=NameNodeMetrics",
"modelerType": "org.apache.hadoop.hdfs.server.namenode.NameNodeMetrics",
"CreateFileOps": 12345,
"FilesCreated": 12300,
"DeleteFileOps": 100,
"GetBlockLocations": 50000,
...
}
]
}

通过 JMX 远程协议(RMI)访问需要额外配置(端口、认证),生产环境通常用 HTTP /jmx endpoint(更简单,与监控工具兼容)。

Prometheus 接入

Hadoop 3.3.0 起可开启内置 Prometheus endpoint,但默认关闭。设置 hadoop.prometheus.endpoint.enabled=true 后,daemon 才会暴露 /prom;未开启时仍可用 JMX Exporter 等方式桥接。

工作流程:

1
2
3
4
1. Hadoop NameNode 启动,MetricsSystem 周期性采样所有 Source 的指标
2. 开启 Prometheus endpoint 后,daemon 把采样结果暴露为 Prometheus 文本格式
3. Prometheus 周期性(例如 15 秒)从 NameNode 的 /prom endpoint 拉取指标
4. Prometheus 把指标存储为时间序列,配合 Grafana 可视化

Hadoop 3.x + Prometheus + Grafana 是生产集群的标准监控栈。Grafana 仪表盘可以从社区获取(Hadoop Grafana Dashboards)或者自己定制。

Hadoop 2.x 时代没有内置 Prometheus 支持,需要用 JMX Exporter(Prometheus 官方提供的 Java Agent)桥接 JMX 到 Prometheus:

1
2
# 在 Hadoop 启动参数里加 JMX Exporter agent
export HADOOP_NAMENODE_OPTS="-javaagent:/path/to/jmx_prometheus_javaagent.jar=12345:/path/to/config.yml"

这种方式 JMX Exporter 进程内嵌在 NameNode JVM 里,监听 12345 端口,把所有 JMX MBean 暴露为 Prometheus 格式。

YARN Log Aggregation:Task 日志集中化

MapReduce Task / Spark Executor 在 NodeManager 上跑,日志默认写本地磁盘(yarn.nodemanager.log-dirs)。Task 完成后日志默认保留 3 小时(yarn.nodemanager.log.retain-seconds)然后被清理。

这种本地保留模式在大型集群里有几个问题:

调试需要 SSH 到具体 NodeManager 节点看日志,运维麻烦。

节点重启后日志丢失。

3 小时保留期短,作业运行几小时后才出问题,可能已经看不到日志。

YARN Log Aggregation(Hadoop 2.x 引入)解决了这些问题——Task 完成后 NodeManager 把日志上传到 HDFS 集中存储,客户端通过 yarn logs 命令从 HDFS 拉取日志。

聚合流程:

1
2
3
4
5
6
1. Container 完成(成功或失败)
2. NodeManager 触发 LogAggregationService
3. LogAggregationFileController 按配置写成 TFile、IndexedFile 等聚合格式
4. 上传到远端根目录(默认 /tmp/logs)及 suffix(默认 logs)下
5. 以 applicationId 和 NodeManager 维度组织文件;实际路径取决于 file controller 与配置
6. 本地日志删除(释放 NM 节点磁盘空间)

远端日志通常按用户、applicationId 和节点组织;不要把某一种 file controller 的路径写成固定的 flow/applicationAttemptId 结构。

客户端查询:

1
2
3
4
5
6
7
8
9
10
11
# 查询某作业的所有日志
yarn logs -applicationId application_1721900000000_0001

# 查询某 container 的日志
yarn logs -applicationId application_1721900000000_0001 -containerId container_1721900000000_0001_01_000001

# 查询某 node manager 的日志
yarn logs -applicationId application_1721900000000_0001 -nodeAddress node-1:8041

# 仅查 stderr
yarn logs -applicationId application_1721900000000_0001 -log_files stderr

YARN Log Aggregation 的关键配置:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
<property>
<name>yarn.log-aggregation-enable</name>
<value>true</value>
</property>
<property>
<name>yarn.nodemanager.log-aggregation.roll-monitoring-interval-seconds</name>
<value>3600</value> <!-- 长作业日志每 1 小时滚动一次上传 -->
</property>
<property>
<name>yarn.log-aggregation.retain-seconds</name>
<value>2592000</value> <!-- HDFS 上日志保留 30 天 -->
</property>
<property>
<name>yarn.log.server.url</name>
<value>http://history-server-host:19888/jobhistory/logs</value>
</property>

yarn.log.server.url 指向 MapReduce History Server(独立服务)。History Server 提供 Web UI 查看 HDFS 上聚合的日志,用户不需要 SSH 到 NodeManager。

MapReduce History Server

MapReduce History Server 是一个独立的服务(默认 19888 端口),承担两个职责:

第一,作业历史查询。每个 MapReduce 作业完成时把作业元数据(job.xml、task 列表、counter)写到 HDFS(mapreduce.jobhistory.done-dir)。History Server 加载这些元数据,提供 Web UI 让用户查看历史作业。

第二,聚合日志查询。History Server 通过 yarn.log.server.url 接受日志查询请求,从 HDFS 拉取对应 container 的聚合日志,渲染成 HTML 或返回纯文本。

History Server 是无状态服务——所有数据在 HDFS,History Server 进程崩溃后重启不丢。可以部署多个实例做负载均衡。

生产集群通常单独部署 History Server 节点(不与 NameNode / ResourceManager 同部署),避免 History Server 处理大量日志查询时影响核心服务。

Audit Log:操作审计

除了性能指标和日志,Hadoop 还有 Audit Log——记录每次文件系统操作的审计日志。

NameNode 的 Audit Log(hdfs.audit.logger)记录每个 HDFS 操作:

1
2
3
4
2026-07-26 10:00:01,234 INFO FSNamesystem.audit: allowed=true \
ugi=alice@REALM (auth:KERBEROS) \
ip=/10.0.0.5 cmd=open src=/user/alice/data.txt \
dst=null perm=null proto=rpc

每条记录包含:用户身份(ugi)、来源 IP、操作类型(cmd)、源路径(src)、目标路径(dst)、权限(perm)、协议(rpc / webhdfs)。

ResourceManager 的 Audit Log 记录每个 YARN 操作:

1
2
3
2026-07-26 10:00:01,234 INFO RMAuditLogger.audit: allowed=true \
ugi=alice ip=/10.0.0.5 cmd=submitApplication \
param=application_1721900000000_0001

Audit Log 默认输出到 log4j 配置的日志文件,可以单独配置输出到不同文件(例如 hdfs-audit.logyarn-audit.log)方便检索。

生产集群通常把 Audit Log 推送到 SIEM(Security Information and Event Management)系统——Splunk、ELK Stack、OpenSearch 等——做安全审计和合规留档。

实验:观察可观测性数据

实验状态:UNVERIFIED_RUNTIME。下面是验证步骤,本轮没有连接真实 Hadoop 集群执行。

NameNode Web UI(9870 端口)的 Metrics 标签页展示关键指标。更详细的指标通过 /jmx endpoint 访问:

1
curl http://namenode-host:9870/jmx | python3 -m json.tool | head -30

输出大量 MBean。重点关注:

NameNodeMetrics:CreateFileOps、FilesCreated、DeleteFileOps、GetBlockLocations、TransactionsAvgTime(EditLog 写入延迟)、SyncsAvgTime。

FSNamesystem:CapacityTotal、CapacityUsed、CapacityRemaining、UnderReplicatedBlocks、MissingBlocks、BlocksTotal、FilesTotal、LiveNodes、DeadNodes。

JvmMetrics:MemHeapUsedM、MemHeapCommittedM、MemHeapMaxM、GcCount、GcTimeMillis、ThreadsNew、ThreadsRunnable。

ResourceManager Web UI(8088 端口)类似,可以查看 RM 的 JMX 指标。

MapReduce 作业的可观测性通过 ApplicationMaster Web UI。每个作业的 AM 都有 Web UI(端口随机),点进 Tasks 标签页看每个 Task 的状态、Counter、日志。

聚合日志查询:

1
2
# 假设你刚跑完一个作业 application_1721900000000_0001
yarn logs -applicationId application_1721900000000_0001 | head -50

输出该作业所有 container 的 stdout/stderr,按 container 分段。

如果想看历史作业,访问 MapReduce History Server Web UI(默认 19888 端口),按用户、队列、时间范围搜索。

生产环境的可观测性最佳实践

生产集群的可观测性栈通常是:

指标:Hadoop Metrics V2 + JMX Exporter + Prometheus + Grafana。Prometheus 周期性从每个 Hadoop 节点拉取 JMX 指标,Grafana 渲染仪表盘。关键告警(NameNode heap 使用率 > 80%、UnderReplicatedBlocks > 1000、DeadNodes > 0)通过 Prometheus Alertmanager 触发。

日志:YARN Log Aggregation + HDFS 集中存储 + ELK / Splunk 检索。容器日志自动聚合到 HDFS,应用层日志(Hive SQL、Spark Job)通过 Flume / Filebeat 推到 ELK。

审计:Audit Log + SIEM(Splunk / OpenSearch)。安全相关操作(Kerberos 认证、ACL 检查、Proxy User 访问)记录到 SIEM 做合规留档。

分布式追踪:Hadoop 3.4 并没有发布 OpenTelemetry 支持;HTrace 已被替换成 No-Op tracer,OpenTelemetry 支持仍停留在未发布的 JIRA 补丁状态。跨服务 RPC 追踪需要外部探针或自行集成。

告警:关键指标(NameNode heap、UnderReplicated Blocks、DeadNodes、ResourceManager RPC queue length)配置阈值告警,通过 Slack / 钉钉 / PagerDuty 通知。

这套栈的可观测性比 Hadoop 早期(手工 tail 日志、SSH 看进程)强大得多。Hadoop 3.x 在可观测性上的最大改进是把 Metrics V2 与现代云原生监控栈(Prometheus / Grafana)无缝集成。

模式提炼

Hadoop 可观测性体现的设计模式:

1
2
3
4
5
6
7
模式:三层可观测性 + Source-Sink 解耦 + 集中存储 + 多协议暴露

- 指标层(Metrics V2):Source-Sink 架构,新加 sink 不改 source
- 运行时层(JMX):标准 Java 指标接口,外部工具可插拔
- 日志层(YARN Aggregation):本地落盘 → 集中存储 → 命令行 / Web 查询
- 审计层(Audit Log):每次操作可追溯,合规留档
- 多协议暴露(HTTP /jmx + Prometheus pull + JMX RMI + REST API)

这个模式不只是 Hadoop。Kubernetes 的可观测性栈(Prometheus + Grafana + Fluentd + Jaeger)几乎对应。Kafka 的可观测性(JMX + Prometheus JMX Exporter + Burrow + Cruise Control)也是同样思路。

数据库领域的对应物是"工作负载监控"。Oracle Enterprise Manager、PostgreSQL pg_stat_statements、MySQL Performance Schema 都提供类似的指标 / 查询 / 审计三层。

工程迁移表

Hadoop 概念 Kubernetes Kafka Cassandra Database
Metrics V2 Source kubelet metrics kafka metrics metrics-core performance schema
Metrics Sink - - - -
JMX Bean - JMX MBean JMX MBean pg_stat views
/jmx endpoint /metrics endpoint - - pg_stat views
Prometheus native JMX Exporter JMX Exporter postgres_exporter
Log Aggregation Loki / EFK Log Compaction Audit Log binlog / WAL
Audit Log Kubernetes Audit Kafka Audit Cassandra Audit DB audit

注意 Kubernetes 这一列。K8s 的 Pod 指标(kubelet 内置)通过 /metrics endpoint 暴露 Prometheus 格式,不需要额外桥接——K8s 设计时就考虑了 Prometheus。Hadoop 是后集成 Prometheus(3.x 才有原生支持),所以需要 JMX Exporter 桥接。

常见误解

误解一:“Metrics V2 完全取代了 Metrics V1”。Metrics V2 框架取代了 V1,但很多 Hadoop 组件内部仍然用 V1 API(兼容代码),V2 框架做了适配。从用户视角看不出来差别,但开发自定义 Source 时需要按 V2 API。

误解二:“JMX 必须用 RMI 协议”。Hadoop 的 /jmx HTTP endpoint 是 JMX MBean 的 HTTP 桥接,不是 RMI。这让 Web 浏览器和 curl 都能访问,比 RMI 简单得多。RMI 仍然可用(配置 com.sun.management.jmxremote),但生产环境通常不用。

误解三:“YARN Log Aggregation 让本地日志立即删除”。本地日志在 Container 完成后才上传到 HDFS,上传完成后才删除。如果上传失败,本地日志保留(重试上传)。上传是异步的,Container 完成到日志出现在 HDFS 可能有几秒到几分钟延迟。

误解四:“History Server 处理日志查询很慢”。History Server 本身是 I/O 密集——从 HDFS 拉取日志、解压、返回。在大集群上查询高 QPS 时 History Server CPU 可能瓶颈。生产环境通常部署多个 History Server 实例 + 负载均衡。

误解五:“Audit Log 记录所有 HDFS 操作”。Audit Log 默认记录通过 RPC 的操作,不记录通过 WebHDFS / HttpFS 的 HTTP 操作(需要单独配置)。本地 DataNode 操作(block 复制、删除)也不在 Audit Log 里。完整审计需要配合其他日志。

练习

  1. 在 Hadoop 集群上访问 http://namenode-host:9870/jmx,找 FSNamesystem 和 JvmMetrics MBean。观察 CapacityTotal / CapacityUsed / CapacityRemaining、MemHeapUsedM / GcTimeMillis 在不同时刻的值。

  2. 提交一个 MapReduce 作业,作业完成后用 yarn logs -applicationId <appId> 查询聚合日志。如果集群没启用 Log Aggregation,对比本地日志访问(SSH 到 NM 节点)的体验差异。

  3. apache/hadoop 源码里找到 MetricsSystem.javaMetricsSource.javaPrometheusMetricsSink.javaLogAggregationService.java(hadoop-common-project 和 hadoop-yarn-server-nodemanager 模块),观察指标系统和日志聚合的核心实现。

  4. 思考题:如果让 Hadoop 把所有 Audit Log 推送到 Kafka 而不是本地文件,需要改造什么?这种改造对实时安全监控(异常行为检测)有什么价值?

系列导航

序号 主题 状态
00 导读:节点总会失败
01 HDFS 架构与三层切分
02 文件写入路径与流水线
03 文件读取路径与副本选择
04 NameNode 内存模型与启动恢复
05 HDFS HA 与脑裂防御
06 HDFS 3.x 演进与纠删码
07 YARN 架构与三方契约
08 YARN 资源模型、Container 与 NodeLabel
09 YARN 调度器对比:FIFO、Capacity、Fair
10 YARN 应用程序生命周期
11 YARN HA 与 Federation
12 YARN Timeline Service v2
13 MapReduce 编程模型与分而治之
14 MapReduce Shuffle 全流程
15 MRv2 on YARN ApplicationMaster 与 Task Attempt
16 Hadoop RPC 协议栈
17 序列化与压缩
18 Hadoop 安全 Kerberos Token ProxyUser 上一篇
19 监控与运维 Metrics JMX 日志聚合 本篇
20 Hadoop 生态 Hive HBase Pig 下一篇
21 Hadoop 与对象存储 Kubernetes 演进对比
22 Hadoop 设计遗产从 GFS MapReduce 到云原生

参考资料