背压不是故障,而是系统自动限速机制。它让快的地方等慢的地方,避免整条链路失控。

一、为什么要理解背压

很多人在使用 Flink 时,都遇到过这些现象:

  • 作业吞吐量突然下降
  • 数据延迟越来越高
  • Checkpoint 耗时持续增长
  • Source 明明很快,却读不动数据
  • 某个算子 CPU 很低,任务却整体hang住

表面看像是某个节点性能问题,实际往往是整条链路的速度失衡。在分布式流处理中,上游持续生产数据,下游持续消费数据。只要某一环节处理速度下降,压力就会沿着链路逐级向上传递,最终影响整个作业运行节奏。这个机制,就是 背压(Backpressure)

理解背压,本质是在理解 Flink 如何维持整个系统的稳定运行。

二、什么是背压

背压可以简单理解为:下游处理速度低于上游生产速度时,上游被迫减速甚至阻塞的现象

例如一条作业链路:

Source → Map → KeyBy → Window → Sink

如果 Sink 写入外部系统变慢:

Sink 处理慢 → Window 输出堆积 → Map 发送受阻 → Source 消费降速

这就是压力反向传导的过程。它不是异常机制,而是 Flink 的一种自我保护能力

如果没有背压控制,上游持续高速写出,最终会导致内存堆积、网络拥塞甚至任务崩溃。

三、Flink 背压产生的底层原因

很多人以为背压是线程被锁住了,实际上在 Flink 中,背压核心原因通常是:

网络缓冲区(Network Buffer)供需失衡。

数据在算子之间传输时,不是直接逐条发送,而是先写入 Buffer,再批量发送。核心组件:

组件作用
ResultPartition上游输出结果分区
InputGate下游输入消费入口
NetworkBufferPool全局网络缓冲池
LocalBufferPool单任务本地 Buffer 池
Credit-Based Flow Control下游控制上游发送速率

当下游消费变慢时:

  1. Buffer 长时间得不到释放
  2. 上游申请不到新的可写 Buffer
  3. 数据发送阻塞
  4. 背压出现

四、源码解析:背压是如何形成的

1. 上游算子输出数据

算子处理完数据后,通过 RecordWriter 写出:

recordWriter.emit(record);

进入序列化逻辑:

RecordWriter.serializeRecord(serializer, record);

最终写入分区:

ResultPartitionWriter.emitRecord(...)

2. 申请可写 Buffer

数据写出前,需要从本地缓冲池申请 Buffer。Flink 1.20 常见路径会进入:

LocalBufferPool.requestBufferBuilderBlocking()

如果当前没有可用 Buffer,线程会等待。这就是背压在代码层最直接的体现:

上游线程想发数据,但没有可用 Buffer。


3. 为什么没有 Buffer

因为 Buffer 仍被下游占用,还未归还。

下游通过:SingleInputGate.pollNext()读取数据。

若消费速度慢:

  • Buffer 积压在 InputChannel
  • Buffer 无法及时 recycle
  • 上游持续拿不到新 Buffer

最终形成完整链路:

下游慢
→ Buffer不释放
→ 上游申请Buffer阻塞
→ emit卡住
→ Source降速

五、Flink 流控机制:Credit-Based Flow Control

背压只是结果,真正维持系统稳定的是 流控机制。Flink 当前默认采用 Credit-Based Flow Control

核心思想:不是上游想发多少就发多少,而是下游决定我还有多少空余 Buffer,你最多发多少数据。即:

  • 每个可接收 Buffer 对应一个 credit
  • 下游消费一个 Buffer,归还一个 credit
  • 上游收到 credit 后才能继续发送

关键类

职责
RemoteInputChannel维护 credit 与远程拉取
PartitionRequestQueue控制网络发送队列
SingleInputGate多输入统一消费入口
NettyPartitionRequestClient网络请求客户端

注:本专栏均依据Flink 1.20版本源码

价值:这种机制避免了:

  • 单节点 Buffer 被打满
  • 网络洪峰冲击
  • 下游 OOM
  • 整条链路雪崩式积压

六、如何判断作业是否发生背压

1. Flink Web UI 指标(最直接)

Task 背压指标通常展示三种状态:

指标含义
Idle空闲等待输入
Busy正在处理数据
BackPressured输出受阻

如果某算子:BackPressured = 85%,说明大量时间花在等待下游可写能力。

2. 常见线上表现

有以下现象可以往这个方向考虑:

  • Source TPS 明显下降
  • Sink 延迟升高
  • CPU 不高但吞吐低
  • Checkpoint duration 持续增加
  • watermark 推进缓慢

七、背压治理方案(实战)

1. 找到最慢节点

多数背压根源在:

  • Sink 写库慢
  • 外部接口慢
  • 热 key 导致窗口拥塞
  • 单并行度过低

2. 提升并行度

如果最慢节点的 subtask 的busy程度都很高,说明当前并发无法处理现在的数据量,则需要提升并行度,例如:sink parallelism: 2 → 8,可直接缓解瓶颈算子压力。

如果发现增加并行度处理能力仍然变化不大,则需要确认 sink 对应服务端是否存在异常。

3. 增加网络内存

通过以下配置增加网络内存提升 Buffer 容量。:

taskmanager.memory.network.fraction: 0.15
taskmanager.memory.network.min: 256mb
taskmanager.memory.network.max: 1gb

4. 使用异步 I/O

通过异步I/O,避免 Sink 阻塞线程。例如:

  • Async JDBC
  • Async Redis
  • Async HTTP

5. 打散热点 Key

如果发现Busy subtask 只有几个,且均为输入数据量高的,则是热点 key 让单个 subtask 承担过多压力,从而导致整体吞吐下降。可通过以下方式来缓解倾斜:

  • key 加盐
  • 两阶段聚合
  • rebalance 后再聚合

6. 排查机器是否异常

如果是偶发的个别 busy,重启就恢复的情况,可以考虑是否对于 subtask 所在的底层机器存在故障,如负载高或内存、网卡硬件故障等情况。

八、背压与 Checkpoint 的关系

这是很多线上问题的根源,当背压严重时:

  • Barrier 传播变慢
  • 对齐时间变长
  • aligned checkpoint 持续超时
  • checkpoint interval 被拖垮

所以很多人看到 Checkpoint 慢,实际问题并不在 Checkpoint,而在数据链路堵塞。

九、总结:压力如何在系统中流动

我们可以这样理解 Flink 的背压机制:

环节含义
Sink 变慢出口堵塞
Buffer 堆积仓位被占满
上游阻塞生产线减速
Source 降速原料暂停进入

背压不是故障,而是系统自动限速机制。它让快的地方等慢的地方,避免整条链路失控。

Flink 正是通过 Buffer 管理与 Credit 流控,让分布式链路在高吞吐场景下依然保持稳定。当我们理解了数据如何流动、压力如何传导,下一个关键问题就是:

这些数据到底存在哪里?

Task 内存、状态内存、网络 Buffer 又如何划分?

下一篇,我们进入:《内存管理机制:从 Task 到 NetworkBuffer》

更多推荐