Flink internals

STREAM PROCESSING / SOURCE READING / LESSON 14

Monitoring, performance, and troubleshooting

Turn metrics, flame graphs, backpressure, checkpoints, and symptoms into an executable troubleshooting sequence.

Reading
60 min
Track
Flink internals
Source
Chinese source notes

The source notes for this track are currently maintained in Chinese.

预计阅读时间: 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 恢复并修改并行度