Flink internals

STREAM PROCESSING / SOURCE READING / LESSON 08

Backpressure

Understand backpressure through symptoms, propagation paths, and network flow control.

Reading
60 min
Track
Flink internals
Source
Chinese source notes

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

预计阅读时间: 60 分钟 前置阅读: doc-02(流数据链路), doc-09(网络栈) 下一次阅读: doc-09(网络栈详解)


反压机制深度解析 图 01

1. Credit-based Flow Control 协议

1.1 协议设计目标

Flink 的 Credit-based 流控是接收端驱动的流控机制,设计目标:

  1. 无阻塞反压:Consumer 慢时自动抑制 Producer,不需要中心调度
  2. 精确反压:反压只影响特定的 Subpartition-Channel 对,不影响其他 Channel
  3. 无死锁:保证独占缓冲 (exclusive credits) + 浮动缓冲 (floating buffers) 的分配

1.2 Credit 协议五阶段

Phase 1: 连接建立
  Consumer → Producer: PartitionRequest(initialCredit)
    initialCredit = networkBuffersPerChannel (默认 2)
  Producer: credit_counter = initialCredit

Phase 2: 数据传输
  Producer: 有 credit → 发送 BufferResponse(buffer, sequenceNumber, backlog)
  Consumer: 接收 buffer → 消费 → 释放 → unannouncedCredit++

Phase 3: Credit 归还
  Consumer: 当 unannouncedCredit > 0 → AddCredit(n) → unannouncedCredit = 0
  Producer: credit += n → 可继续发送

Phase 4: 积压通知
  Producer: 当 backlog > 0(有排队数据) → BacklogAnnouncement(backlog)
  Consumer: 根据 backlog 请求浮动缓冲 → AddCredit(新缓冲数)

Phase 5: 反压
  Consumer 慢 → buffer 消费慢 → unannouncedCredit 不增长
  Producer credit 用尽 → 无法获取新 buffer → isAvailable=false → 反压

1.3 缓冲类型

类型来源数量特性
Exclusive (独占)initialCredit固定每个 Channel 保证的最低缓冲,防止死锁
Floating (浮动)共享池按需根据 backlog 动态分配,无背压时释放
Overdraft (透支)共享池限制超出 currentPoolSize 的临时缓冲,处理突发尖峰

2. 反压检测与传播

2.1 反压检测点

// Task.java — 反压判断
public boolean isBackPressured() {
    if (invokable == null || partitionWriters.length == 0
        || (executionState != ExecutionState.INITIALIZING
            && executionState != ExecutionState.RUNNING)) {
        return false;
    }
    // 检查所有输出写入器
    for (ResultPartitionWriter writer : partitionWriters) {
        if (!writer.isAvailable()) {  // 任何一个无法获取 buffer
            return true;              // → 反压!
        }
    }
    return false;
}

// ResultPartition.getAvailableFuture()
// → LocalBufferPool.getAvailableFuture()
// → availabilityHelper (当无法分配 NetworkBuffer 时变为不可用)

2.2 反压的逐级传播

Sink 写入外部 DB 慢 (100ms/条)
  → Sink 消费速度慢 → NetworkBuffer 积累在 Sink 的 InputChannel 中
  → Sink 的 InputGate 不再归还 buffer → 不发送 AddCredit
  → 上游 Window 算子的 Output Subpartition 没有新 credit
  → Window 的 LocalBufferPool 缓冲池耗尽
  → Window 的 ResultPartition.isAvailable = false
  → Window 算子反压! → processInput() 暂停
  → 继续传播到 Source...
  → Source 停止 KafkaConsumer.poll()

反压传播的关键机制:不是通过主动通知,而是通过能力不足的被动传播

  • 下游慢 → 不归还 buffer → 无 credit → 上游无法发送 → 上游缓冲区满 → 上游停止消费 → 更上游反压...

2.3 Web UI 反压显示

OK (绿色):  backPressuredTimeMs < 10ms/s    — 无压力
LOW (黄色): backPressuredTimeMs 10-50ms/s  — 轻度
HIGH (红色): backPressuredTimeMs > 50ms/s   — 重度

实现:每个 Task 抽样 backPressuredTime,取 3 次抽样的比例
- softBackPressuredTime: 输出反压但输入仍有数据
- hardBackPressuredTime: 输入和输出都反压

3. 反压常见根因与排查

3.1 根因定位方法论

排查步骤:
1. Web UI → Job → SubTasks → 找出 backPressuredTime 最高的算子
2. 确定反压是产生于此算子还是传播自此算子:
   - 如果 Source 就反压 → 反压由下游产生,传播到 Source
   - 如果只有 Sink 反压 → Sink 本身是瓶颈
3. 找到第一个不反压的算子 → 它的上一个算子就是瓶颈
4. 分析瓶颈算子的工作:
   - CPU 高 → 计算瓶颈 → 优化代码/增加并行度
   - CPU 低 + Sink 慢 → 外部 IO 瓶颈 → 批量写入/异步 IO
   - GC 频繁 → 内存瓶颈 → 调整 Heap/切换 RocksDB
   - 数据倾斜 → 某个 Subtask 处理量远大于其他 → 重新选 Key/Salt

3.2 常见原因表

原因症状排查解决方案
Sink 写入慢外部 DB/API 延迟高检查外部系统延迟asyncIO, 批量写入, 增加 Sink 并行度
数据倾斜Subtask 间 numRecordsIn 差距 > 10xSubTasks Metrics加 Salt, 重新设计 Key
GC 频繁TM GC Time > 5%TM Metrics → GC Time增大 Heap/Managed, 使用 G1GC
外部 API 调用asyncIO 容量饱和asyncIO Metrics增大 asyncIO 容量
State 操作慢RocksDB 慢 (compaction)RocksDB Metrics调整 RocksDB 配置
Checkpoint 占用对齐期间处理暂停Checkpoint Metrics启用 Unaligned Checkpoint

4. BufferDebloater — 自适应缓冲大小

// BufferDebloater.java — 基于吞吐量动态调整缓冲
public int recalculate(int currentBufferSize, long throughput) {
    // throughput: bytes/ms
    // targetDuration: 1ms (默认)
    // 最优缓冲大小 = throughput × 1ms = 保证 1ms 的数据量
    // 高吞吐场景 → 缓冲自动增大 → 减少 Netty 调用
    // 低吞吐场景 → 缓冲自动减小 → 减少延迟
    return Math.clamp(throughput * targetDuration, minSize, maxSize);
}

5. 源码导航

文件关键类/方法职责
flink-runtime/.../taskmanager/Task.javaisBackPressured()反压判断
flink-runtime/.../partition/ResultPartition.javagetAvailableFuture()Producer 端可用性
flink-runtime/.../buffer/LocalBufferPool.javarequestMemorySegment(), recycle()缓冲池管理
flink-runtime/.../netty/CreditBasedPartitionRequestClientHandler.javachannelRead(), userEventTriggered()Credit 协议
flink-runtime/.../partition/consumer/RemoteInputChannel.javaonBuffer(), notifyBufferAvailable()Consumer 端 Credit
flink-runtime/.../buffer/BufferDebloater.javarecalculate()缓冲自适应

6. 常见问题 / 面试题

Q1: Credit-based 流控和 TCP 流控的区别?

A: TCP 流控是字节级(sliding window = bytes),Flink Credit 是缓冲级(credit = buffer count)。关键区别:Flink 的缓冲维度直接对应内存管理——每个 buffer 是独立的 32KB 内存段,消费完成后可以回收给缓冲池复用,同时控制内存使用量(而不只是网络流量)。

Q2: 反压会最终导致 Source 停止消费吗?

A: 是的。反压是一级一级传播的——Sink 慢 → Agg 的输出 buffer 满 → Agg 处理慢 → Map 的输入 buffer 满 → Map 处理慢 → Source buffer 满 → Source 停止 poll()。这是精确的 backpressure chain,保证了不会 OOM。

Q3: 透支缓冲在什么场景下有用?过度增大会如何?

A: 透支缓冲允许临时超出 currentPoolSize,处理数据量的短时尖峰。过度增大的代价:(1) 占用更多全局缓冲池中的内存 → 可能实际导致其他 Channel 更早反压;(2) 延迟了真正的反压信号 → 可能导致更长的恢复时间。


下一步