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

1. Credit-based Flow Control 协议
1.1 协议设计目标
Flink 的 Credit-based 流控是接收端驱动的流控机制,设计目标:
- 无阻塞反压:Consumer 慢时自动抑制 Producer,不需要中心调度
- 精确反压:反压只影响特定的 Subpartition-Channel 对,不影响其他 Channel
- 无死锁:保证独占缓冲 (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 差距 > 10x | SubTasks 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.java | isBackPressured() | 反压判断 |
flink-runtime/.../partition/ResultPartition.java | getAvailableFuture() | Producer 端可用性 |
flink-runtime/.../buffer/LocalBufferPool.java | requestMemorySegment(), recycle() | 缓冲池管理 |
flink-runtime/.../netty/CreditBasedPartitionRequestClientHandler.java | channelRead(), userEventTriggered() | Credit 协议 |
flink-runtime/.../partition/consumer/RemoteInputChannel.java | onBuffer(), notifyBufferAvailable() | Consumer 端 Credit |
flink-runtime/.../buffer/BufferDebloater.java | recalculate() | 缓冲自适应 |
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) 延迟了真正的反压信号 → 可能导致更长的恢复时间。
下一步
- 网络栈详解: doc-09-network-stack.md
- TM 内存: doc-10-taskmanager-memory.md