Flink 源码与运行原理

流式计算 / 源码阅读 / LESSON 14

监控 & 性能 & 排障深度手册

把指标、火焰图、背压、Checkpoint 与故障现象组织为一套可执行的性能排障顺序。

阅读时间
60 分钟
学习路径
Flink 源码与运行原理
内容来源
Flink 深度笔记

预计阅读时间: 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类型配置适用
PrometheusPull (HTTP)metrics.reporter.prom.factory.class: org.apache.flink.metrics.prometheus.PrometheusReporterFactory生产标准
JMXPull (MBean)metrics.reporter.jmx.factory.class: ...JmxReporterFactory本地调试
InfluxDBPushmetrics.reporter.influxdb.factory.class: ...InfluxdbReporterFactoryGrafana
GraphitePushmetrics.reporter.graphite.factory.class: ...GraphiteReporterFactoryGrafana
StatsDPushmetrics.reporter.statsd.factory.class: ...StatsDReporterFactoryDatadog
SLF4JLogmetrics.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/sConsumer 慢
hardBackPressuredTimeMs< 10ms/s> 50ms/s整个管线慢
idleTimeMs少量过多Source 无数据 / 反压
currentOutputWatermark持续增加停滞Idle Source / 大延迟
checkpointDuration< 30s> 5min状态大 / DFS 慢 / 对齐慢
checkpointSize稳定持续递增状态泄漏
checkpointAlignmentTime< 1s> 10sChannel 间数据倾斜
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.javaregister(), unregister()Metric 注册中心
flink-runtime/.../metrics/groups/OperatorMetricGroup.java算子级 Metric
flink-runtime/.../metrics/groups/OperatorIOMetricGroup.javaIO Metric
flink-metrics/flink-metrics-prometheus/Prometheus Reporter
flink-runtime/.../checkpoint/CheckpointStatsTracker.javaCheckpoint 统计

下一步

恭喜! 你已完成 Flink 学习文档全部章节。

推荐的实际操作路径:

  1. 启动本地集群 + WordCount
  2. StreamTask.processInput()CheckpointCoordinator.triggerCheckpoint() 加断点
  3. 配置 Prometheus → Grafana → 监控本地 Job
  4. 刻意制造反压场景 (Sink 慢 / 倾斜) → 在 Web UI 中观察
  5. 从 Checkpoint/Savepoint 恢复并修改并行度