上一篇讲了安全。本篇讲 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.x 起内置了 Prometheus 支持——通过 PrometheusMetricsSink 让 Prometheus 直接拉取 Hadoop 指标。这让 Hadoop 接入现代云原生监控栈。

工作流程:

1
2
3
4
1. Hadoop NameNode 启动,MetricsSystem 周期性采样所有 Source 的指标
2. PrometheusMetricsSink 接收指标,缓存到内存(最近 N 秒)
3. Prometheus 周期性(默认 15 秒)从 NameNode 的 /prometheus 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. 该服务把 container 的 stdout/stderr/syslog 打包成 tar.gz
4. 上传到 HDFS:/app-logs/<user>/<flow>/<applicationAttemptId>/<host>_<containerId>
5. 上传一个索引文件(哪个 container 在哪个 HDFS 文件)
6. 本地日志删除(释放 NM 节点磁盘空间)

HDFS 上的日志文件按用户组织目录,配合 Timeline Service v2 的 Flow 概念,可以按 flow 查询所有日志。

客户端查询:

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 -nodeId node-1:8041

# 仅查 stderr
yarn logs -applicationId application_1721900000000_0001 -logFiles 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 等——做安全审计和合规留档。

实验:观察可观测性数据

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 做合规留档。

分布式追踪:OpenTelemetry(Hadoop 3.4+ 实验)。跨服务的 RPC 追踪,例如 Client → NameNode → DataNode 的完整调用链。

告警:关键指标(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-06 HDFS 存储层 第一阶段(已完成)
07-12 YARN 资源管理层 第二阶段(已完成)
13-15 MapReduce 计算模型 第三阶段(已完成)
16 Hadoop RPC 协议栈
17 序列化与压缩:Writable、Avro 与 Codec
18 Hadoop 安全:Kerberos、Delegation Token、Proxy User 上一篇
19 监控与运维:Metrics V2、JMX、日志聚合 本篇(HA / 安全 / Common 篇完结)
20-22 演进、生态与对比(Hadoop 生态 / 对象存储 + K8s / 设计遗产) 第五阶段,下一篇开始

参考资料