预计阅读时间: 60 分钟 前置阅读: 全部 doc (运维必备, 需要了解所有子系统)
1. Metrics 体系架构
1.1 Metric 类型
// Flink 支持的 4 种基本 Metric 类型
Counter counter = metricGroup.counter("numRecordsIn");
// 单调递增计数器 (如处理的总记录数)
Gauge gauge = metricGroup.gauge("currentOutputWatermark", () -> waterMark);
// 瞬时值 (如当前 Watermark)
Histogram histogram = metricGroup.histogram("checkpointDuration",
new DescriptiveStatisticsHistogram(1000));
// 分布统计 (P50/P90/P95/P99)
Meter meter = metricGroup.meter("numRecordsInPerSecond",
new MeterView(counter, 60));
// 速率 (如每秒处理记录数)
1.2 MetricGroup 树形结构
MetricGroup 树 (从根到叶):
TaskManagerMetricGroup
└── JobMetricGroup (job-id)
└── TaskMetricGroup (task-id)
├── OperatorMetricGroup (operator-name)
│ ├── Counter "numRecordsIn" // 接收记录数
│ ├── Counter "numRecordsOut" // 输出记录数
│ ├── Gauge "currentOutputWatermark" // 当前 Watermark
│ ├── Counter "numLateRecordsDropped" // 丢弃的延迟记录
│ └── Meter "numRecordsInPerSecond" // 接收速率
│
├── OperatorIOMetricGroup
│ ├── Counter "numBytesIn" // 接收字节数
│ └── Counter "numBytesOut" // 输出字节数
│
└── JVMMetricGroup
├── Gauge "Heap.Used" // Heap 使用
├── Gauge "Heap.Committed" // Heap 提交
├── Gauge "Heap.Max" // Heap 最大
├── Gauge "Direct.TotalCapacity" // Direct Memory 容量
├── Gauge "Direct.MemoryUsed" // Direct Memory 使用
├── Gauge "GarbageCollector.G1 Young Generation.Time"
├── Gauge "GarbageCollector.G1 Old Generation.Time"
└── Gauge "GarbageCollector.G1 Old Generation.Count"
1.3 Metrics Reporters
| Reporter | 类型 | 配置 | 适用 |
|---|---|---|---|
| Prometheus | Pull (HTTP) | metrics.reporter.prom.factory.class: org.apache.flink.metrics.prometheus.PrometheusReporterFactory | 生产标准 |
| JMX | Pull (MBean) | metrics.reporter.jmx.factory.class: ...JmxReporterFactory | 本地调试 |
| InfluxDB | Push | metrics.reporter.influxdb.factory.class: ...InfluxdbReporterFactory | Grafana |
| Graphite | Push | metrics.reporter.graphite.factory.class: ...GraphiteReporterFactory | Grafana |
| StatsD | Push | metrics.reporter.statsd.factory.class: ...StatsDReporterFactory | Datadog |
| SLF4J | Log | metrics.reporter.slf4j.factory.class: ...Slf4jReporterFactory | 开发调试 |
2. 关键监控指标
2.1 健康检查三大维度
维度 1: Checkpoint 健康
┌─────────────────────────────────────────────────────┐
│ ✓ numberOfCompletedCheckpoints (持续增长) │
│ ✓ numberOfFailedCheckpoints (= 0 或很少) │
│ ✓ lastCheckpointDuration (< interval × 0.5) │
│ ✓ lastCheckpointSize (稳定 → 健康, 持续增长 → 泄漏) │
└─────────────────────────────────────────────────────┘
维度 2: 反压健康
┌─────────────────────────────────────────────────────┐
│ ✓ softBackPressuredTimeMs (< 10ms/s → OK) │
│ ✓ hardBackPressuredTimeMs (< 10ms/s → OK) │
│ ✓ idleTimeMs (> 0 但不过大 → 正常等待) │
└─────────────────────────────────────────────────────┘
维度 3: 吞吐与延迟
┌─────────────────────────────────────────────────────┐
│ ✓ numRecordsInPerSecond (稳定) │
│ ✓ numRecordsOutPerSecond (≥ numRecordsIn × 0.9) │
│ ✓ currentOutputWatermark (持续前进, 不倒退) │
│ ✓ latencyTracking (P95 < SLA) │
└─────────────────────────────────────────────────────┘
2.2 详细指标速查
| Metric | 正常值 | 异常值 | 原因 |
|---|---|---|---|
numRecordsInPerSecond (numRecordsInRate) | 稳定 | 骤降 ↓ | 上游反压 / Source 停止 |
softBackPressuredTimeMs | < 10ms/s | > 50ms/s | Consumer 慢 |
hardBackPressuredTimeMs | < 10ms/s | > 50ms/s | 整个管线慢 |
idleTimeMs | 少量 | 过多 | Source 无数据 / 反压 |
currentOutputWatermark | 持续增加 | 停滞 | Idle Source / 大延迟 |
checkpointDuration | < 30s | > 5min | 状态大 / DFS 慢 / 对齐慢 |
checkpointSize | 稳定 | 持续递增 | 状态泄漏 |
checkpointAlignmentTime | < 1s | > 10s | Channel 间数据倾斜 |
heapUsed | 周期性 GC | 锯齿 / 持续增长 | GC 频繁 / 内存泄漏 |
gcTime (% of wall clock) | < 5% | > 10% | 大量临时对象 / Heap 不足 |
numBytesInPerSecond | 稳定 | 骤变 | 流量变化 / 网络问题 |
KafkaConsumer.records-lag-max | < 10000 | 持续增长 | 消费跟不上生产 |
3. 排障实战
3.1 Checkpoint 超时 / 失败
症状: Web UI Checkpoint History: FAILED / Timeout
排查步骤:
1. 确认超时类型:
- Duration > timeout → 超时: 看哪个 Subtask 最后 Ack
- Failed (Exception) → 看 JM/TM Log
2. 慢 Subtask 定位:
- Web UI → Job → Checkpoints → 展开最后一次超时的 Checkpoint
- 查看每个 Subtask 的 Ack 时间
- 最后一个 Ack 的 Subtask = 瓶颈
3. 常见原因:
- Barrier 对齐慢 → 启用 Unaligned Checkpoint
- RocksDB Snapshot 慢 → 启用增量 Checkpoint
- S3/HDFS 上传慢 → 增大 checkpoint 间隔 / 换存储
4. 修复验证:
- 增大 interval + timeout
- 启用 unaligned.checkpoint
- 减小状态大小 (TTL)
3.2 反压根因定位
症状: Web UI: HIGH backPressure → 性能下降
排查步骤:
1. Web UI → Job → SubTasks → 按 backPressuredTime 排序
2. 找到第一个 backPressuredTime < LOW 的算子 → 它的上游是瓶颈
3. 分析瓶颈算子:
a. numRecordsInPerSecond > numRecordsOutPerSecond (处理不过来)
→ CPU 高 → 计算瓶颈 → 增加并行度 / 优化代码
→ CPU 低 → 外部 IO 慢 → asyncIO / 批量写入
b. GC Time > 10% → 内存瓶颈 → 增大 Heap / 切换 RocksDB
c. Subtask 间 numRecordsIn 差距 > 10x → 数据倾斜 → Salt / 换 Key
4. 修复:
- 计算瓶颈: 增加并行度 / 优化热点代码
- IO 瓶颈: asyncIO / 批量写入 / 增加连接池
- 倾斜: 加 Salt / 改为 Broadcast + Filter
3.3 Watermark 不推进
症状: Window 不触发 → currentOutputWatermark 停滞
排查步骤:
1. Web UI → 检查各 Source 的 currentOutputWatermark
2. 如果某个 Source Watermark 停滞 → Idle Source → withIdleness()
3. 如果所有 Source Watermark 都在推进 → 下游算子可能有多个输入
→ min(all inputs) 停滞 (= min 运算符导致)
4. 检查下游算子的所有输入 Channel Watermark
5. 一个停滞的 Channel → 上游 Source 或中间算子的 Watermark 问题
修复:
- Source 端: withIdleness(Duration.ofMinutes(1))
- 中间算子: 检查是否有 Operators 在等待多路输入 (Union)
- 窗口端: allowedLateness + sideOutput 处理延迟数据
3.4 OOM / GC 频繁
症状: TM Log: OutOfMemoryError → Task 退出
排查步骤:
1. 确认 OOM 类型:
- Java heap space → Heap 不足 → 查看 heapUsed 趋势
- Direct buffer memory → Network/Direct Memory → 查看 directUsed
- Metaspace → 类加载过多 → 查看 Metaspace
- GC overhead limit → GC 频繁但无法回收 → Heap 太小 或 内存泄漏
2. 如果是 Heap OOM:
a. 确认是否使用 HashMapStateBackend → 切换 RocksDB → 立竿见影
b. 检查 State 是否有 TTL → 无 TTL 的状态会无限增长
c. 检查 operator 是否缓存了 data (如 ProcessFunction)
3. 如果是 Direct Memory OOM:
a. 增大 network.fraction / managed.fraction
b. 减小 parallelism (减少 channel 数)
修复:
- Heap OOM: 增大 task.heap.size / 切换 RocksDB / State TTL
- Direct OOM: 增大 network.fraction / 减小 managed.fraction
- Metaspace: 增大 metaspace.size
4. 生产监控栈
推荐监控栈:
Flink Metrics → Prometheus (拉取 :9249 URL) → Grafana (Dashboard)
关键 Grafana Panels:
Row 1: Job Overview
- Job Uptime, Restart Count, Status
Row 2: Checkpoint Health
- Duration (P50/P95/P99), # of Failures, Size
Row 3: Throughput
- numRecordsIn/s (by Operator), numRecordsOut/s
Row 4: Backpressure
- backPressuredTime (by Task)
Row 5: JVM / Memory
- Heap Used, GC Time, Direct Memory
Row 6: Source/Sink Health
- Kafka Consumer Lag, Sink Write Rate
5. 源码导航
| 文件 | 关键类/方法 | 职责 |
|---|---|---|
flink-runtime/.../metrics/MetricRegistryImpl.java | register(), unregister() | Metric 注册中心 |
flink-runtime/.../metrics/groups/OperatorMetricGroup.java | — | 算子级 Metric |
flink-runtime/.../metrics/groups/OperatorIOMetricGroup.java | — | IO Metric |
flink-metrics/flink-metrics-prometheus/ | — | Prometheus Reporter |
flink-runtime/.../checkpoint/CheckpointStatsTracker.java | — | Checkpoint 统计 |
下一步
恭喜! 你已完成 Flink 学习文档全部章节。
推荐的实际操作路径:
- 启动本地集群 + WordCount
- 在
StreamTask.processInput()和CheckpointCoordinator.triggerCheckpoint()加断点 - 配置 Prometheus → Grafana → 监控本地 Job
- 刻意制造反压场景 (Sink 慢 / 倾斜) → 在 Web UI 中观察
- 从 Checkpoint/Savepoint 恢复并修改并行度