前面几篇把 Flink 的运行时架构、时间与 Watermark、端到端 Exactly-Once、DataStream 与状态容错 一路拆到了底。这些讲的都是「正确性」——作业算得对、故障能恢复、结果不重不丢。但真正把一个流作业推上生产、扛住广告晚高峰几百万 QPS 的洪峰时,你面对的第一个敌人往往不是正确性,而是性能:作业跑着跑着吞吐塌了、端到端延迟从秒级涨到分钟级、Kafka 消费 lag 一路飙红。绝大多数这类事故,根因都指向同一个词——反压(backpressure)。
本篇聚焦「容错运维」里的性能这一环:把反压是什么、怎么定位、怎么调优(尤其是数据倾斜这个 AdTech 里的头号杀手)讲透。
一句话定位:反压不是故障,而是 Flink 用「让上游慢下来」来保护自己不被压垮的自愈信号。调优的目标从来不是「消灭反压」,而是找到那个拖慢全链路的最慢算子并把它救回来——在广告实时链路里,这个最慢算子十有八九是被热 key 打爆的聚合算子。
TL;DR
- 反压 = 下游处理不过来、逐级向上游传导的「减速信号」:某个算子消费速度跟不上,它的输入缓冲被占满,会反向逼停上游,一级级传到 source,最终让 source 放慢/暂停从 Kafka 拉取。
- 反压本身是好事:它是 Flink 的天然反压/限流机制,防止内存被无限堆积撑爆(对比无反压系统直接 OOM)。它是一个症状指示器,告诉你「链路里有个瓶颈」,而不是故障本身。
- Credit-based flow control 是现代 Flink(1.5+)的流控内核:接收端按自己空闲的 buffer 数给发送端发放 credit,发送端只在有 credit 时才发数据——替代了早期 TCP-based 反压「一堵堵一片」的粗暴做法,反压更精准、不会误伤同一条 TCP 连接上的其它 channel。
- 定位反压看三样:Web UI 的 Backpressure 面板(High/OK)、算子指标
busyTimeMsPerSecond/backPressuredTimeMsPerSecond/idleTimeMsPerSecond,以及火焰图(Flame Graph)。 - 找「起点」而非「表象」:反压会从瓶颈算子一路往上游染红,真正的元凶是**第一个「busy 高但 backpressured 低」**的算子——它自己很忙、却不被下游拖着,说明瓶颈就在它这。
- 四大常见成因:① 数据倾斜(keyBy 后某 key 数据量畸大,打满单个 subtask);② 外部系统慢(sink 写库、维表同步 join 阻塞);③ 状态访问 / GC 慢(RocksDB 磁盘 IO、大状态、Full GC);④ 序列化开销(低效 TypeInfo、Kryo 兜底)。
- 数据倾斜是 AdTech 头号问题:按
campaignId/advertiserIdkeyBy 时,爆量大广告主会形成超级热 key,数据全压到一个 subtask,其余 subtask 空闲——整个作业吞吐由这一个最慢 subtask 决定。 - 倾斜解法三板斧:两阶段聚合(local pre-aggregate + global merge)、加盐打散(热 key 拼随机后缀分散再去盐汇总)、必要时
rebalance()/rescale()重分布——本质就是缓存系列里「热 key 打散」思想在流处理里的翻版。 - 外部访问必须异步:维表 join / 外部服务调用用
AsyncFunction+AsyncDataStream做异步 IO,把「同步阻塞等 RTT」变成「并发在途请求」,吞吐能涨一个数量级;再叠一层本地缓存维表命中率更高。 - 其它调优旋钮:并行度(对齐 Kafka 分区数、给瓶颈算子单独提并行度)、network buffer、RocksDB(块缓存 / bloom filter / 托管内存)、算子链与 slot 共享、以及 unaligned checkpoint 化解「反压时 checkpoint 做不完」的死循环。
- AdTech 实践:竞价实时统计作业曾因某爆量广告主的 campaign 变成超级热 key,单 subtask 反压传导让全作业吞吐塌陷、大盘卡死;靠两阶段聚合 + 加盐 + 维表 join 改异步救回,思路与缓存的热 key 处理完全一致。
Table of contents
Open Table of contents
1. 反压是什么、为什么会发生
流处理是一条算子流水线:source → map → keyBy 聚合 → sink,数据像水一样自左向右流。每两个相邻算子之间都有有限大小的缓冲区(network buffer / 队列)。只要下游某个算子的处理速度跟不上上游的产出速度,它们之间的缓冲就会被慢慢填满。
缓冲满了之后会发生什么?这正是反压的核心机制:
- 下游算子处理不过来 → 它的输入缓冲被占满 → 它没法再接收新数据;
- 上游算子想往下游发数据,但发不出去(下游缓冲满)→ 上游的输出缓冲也开始堆积 → 进而它自己的输入缓冲也满;
- 这个「满」的状态一级一级向上游传导,最终传到 source;
- source 拿不到「可以继续发」的许可,只能放慢甚至暂停从 Kafka 拉取数据。
反压逐级传导:Sink 变慢 → 输入 buffer 满 → 逐段回传 credit=0 → 最终逼停 Source。整条链路的吞吐,由那个最慢的算子决定。
关键认知:反压是一件好事。 设想一个没有反压机制的系统:上游只管猛发、下游处理不过来的数据全堆在内存里,迟早 OOM 崩溃。Flink 的反压机制让整条流水线自动降速到最慢算子的节奏——用「暂时慢下来」换「不崩溃、不丢数据」。所以你在 Web UI 看到红色的反压,不该慌着「消灭它」,而应该顺着它去找那个真正的瓶颈算子。
反压真正的坏处在于它的连带效应:一旦发生,端到端延迟会显著上涨(数据在缓冲里排队)、Kafka 消费 lag 累积、Checkpoint 因为 barrier 对齐被卡住而超时(§7 会讲 unaligned checkpoint 如何破解)。这些才是需要治的,而治它们的前提是先定位瓶颈。
2. Credit-based flow control:现代 Flink 的流控内核
反压能「逐级精准传导」,靠的是 Flink 1.5 之后引入的基于信用的流量控制(credit-based flow control)。要理解它的价值,先看它替代掉的旧方案。
2.1 早期 TCP-based 反压的问题
Flink 1.5 之前,反压是靠 TCP 的滑动窗口天然实现的:下游 buffer 满了不再读 socket,TCP 窗口收缩,上游 send() 阻塞,反压就自然产生了。听起来很优雅,但有个致命缺陷:一个 TaskManager 到另一个 TaskManager 之间通常复用一条 TCP 连接,上面承载了多个 subtask channel 的数据。只要其中一个 channel 的下游变慢把 TCP 窗口堵死,同一条连接上所有其它 channel 都被连累——哪怕它们的下游一点都不忙。这叫队头阻塞(head-of-line blocking),反压「误伤一片」。
而且 TCP-based 方案里,反压信号要等 TCP 缓冲、socket 缓冲层层填满才生效,传导慢;barrier 也可能被堵在 socket 缓冲里,影响 checkpoint。
2.2 Credit-based:接收端发「许可」,发送端凭「许可」发
Credit-based flow control 把流控从 TCP 层上移到 Flink 自己的应用层,且做到 channel 粒度:
- 每个下游 input channel 都有自己的一组独占 buffer(exclusive buffers)和一个共享的浮动 buffer 池(floating buffers);
- 下游持续把「我现在还有多少个空闲 buffer」作为 credit(信用额度) 汇报给对应的上游;
- 上游的 subpartition 只有在对应 channel 还有 credit 时才发送数据,每发一个 buffer,credit 减一;
- 下游消费掉数据、腾出 buffer,就补发 credit 给上游。
这样一来:哪个 channel 的下游慢,就只有那个 channel 的 credit 归零、只有那条逻辑流被反压,同一条 TCP 连接上的其它 channel 照常跑——彻底消除了队头阻塞。而且上游在发数据前就能通过 credit 知道下游还能不能收,反压传导更快、更精准,backlog 信息还能帮下游提前申请 floating buffer。
上一节那张图里「credit=0 逐段回传」画的正是这个机制:Sink 慢 → 它给聚合算子的 credit 归零 → 聚合算子发不出、自己也满 → 给 Map 的 credit 归零……一路到 source。这是一条应用层的、按 channel 精确控制的反压链,而不再是 TCP 窗口的粗暴收缩。
相关参数(一般不用动,理解即可):
taskmanager.network.memory.buffers-per-channel(每 channel 独占 buffer 数,默认 2)、taskmanager.network.memory.floating-buffers-per-gate(每个 input gate 的浮动 buffer 数,默认 8)。buffer 太少会限制吞吐、太多会拖慢反压传导速度和 checkpoint(§7)。
3. 定位反压:Web UI、指标与火焰图
治反压第一步永远是定位瓶颈算子。Flink 给了三层由粗到细的工具。
3.1 Web UI 的 Backpressure 面板
作业的 Job Graph 里,每个算子(vertex)会显示一个反压状态:OK(绿)/ LOW / HIGH(红)。Flink 通过周期性采样各 task 线程「是否卡在申请输出 buffer」来估算反压比例。看图的要点是:红色会从瓶颈算子一路往上游蔓延——如果 source → A → B → C 里 C 是瓶颈,你会看到 A、B、C 全红,唯独 C 的下游(或 C 本身)不再往下红。
所以别被「一片红」骗了。真正的元凶是红色链条最下游的那个算子——它是反压的起点。
3.2 三个关键指标:busy / backPressured / idle
Web UI 的红绿灯太粗,Flink 1.13+ 提供了每个 subtask 的时间占比指标(单位 ms/s,即每秒有多少毫秒处于该状态),这是定位反压的主力:
busyTimeMsPerSecond:算子真正在干活(处理记录)的时间占比。接近 1000 说明它 CPU/处理已经打满。backPressuredTimeMsPerSecond:算子被下游反压、卡在等输出 buffer 的时间占比。高 = 它在等下游。idleTimeMsPerSecond:算子没数据可处理、在等上游输入的时间占比。高 = 上游喂不饱它。
三者理论上加起来约等于 1000。用它们判断瓶颈的黄金法则:
瓶颈算子 = 第一个「
busy高、但backPressured低」的算子。 它自己很忙、却没被下游拖着,说明「慢」就出在它身上。它上游的那些算子会是「backPressured高」(被它拖着),下游算子则是「idle高」(它喂不饱下游)。
举例:source(idle高) → map(backPressured高) → agg(busy高、backPressured低) → sink(idle高),一眼就能锁定 agg 是瓶颈——反压从它这里发源。
3.3 火焰图:瓶颈算子「慢在哪一行」
锁定瓶颈算子后,还要知道它慢在哪里——是卡在状态访问?序列化?还是某个正则/JSON 解析?Flink 1.13+ 内置了 Flame Graph(rest.flamegraph.enabled: true 开启),对选中算子的线程做栈采样,画成火焰图。宽度越大的栈帧 = 占用 CPU 越多的调用。常见「宽帧」:
RocksDBMapState.get/put(状态访问慢,见 §4.3、§7);- Kryo 序列化相关栈(类型没注册,走了低效兜底,见 §4.4);
- 用户函数里的重逻辑(正则、JSON、大对象拷贝);
- 同步的外部调用(
jdbc、http阻塞——该改异步了,见 §6)。
三层工具连起来用:Web UI 看哪条链红 → busy/backPressured 指标锁定起点算子 → 火焰图看它慢在哪行代码。定位准了,调优才不会瞎试。
4. 常见成因:反压从哪来
瓶颈算子的「慢」通常逃不出这四类。
4.1 数据倾斜(AdTech 头号元凶)
keyBy 会按 key 的哈希把数据分发到下游各个 subtask。理想情况下 key 分布均匀、每个 subtask 分到差不多的量。但现实里 key 的分布经常是长尾的——少数 key 占了绝大多数流量。
数据倾斜:一个超级热 key 把单个 subtask 打满、其余 subtask 空转。作业吞吐 = 最慢单点吞吐,横向扩容对单 key 无效。
在广告场景里这几乎是必然:按 campaignId / advertiserId 做实时统计时,某个爆量大广告主(比如大促期间猛砸预算的头部客户)的 campaign 会瞬间成为超级热 key,它的曝光/点击事件全被哈希到同一个 subtask。结果:那个 subtask busy=1000、疯狂反压,其余 subtask idle 很高、闲得发慌。这时候加并行度、加机器都没用——因为一个 key 的数据不会被拆到多个 subtask,瓶颈永远是那一个点。解法见 §5。
4.2 外部系统慢:sink 与同步维表 join
第二类元凶是算子里同步调用了外部系统,被外部的 RTT 死死拖住:
- Sink 写外部库慢:写 MySQL/HBase/ES,如果是逐条同步写,每条都要等一个网络往返 + 磁盘写入,吞吐被外部系统的单点能力锁死。批量写(buffer 攒批 + flush)能缓解,但根子上同步就是慢。
- 同步维表 join:流上每来一条数据,就同步去查一次维表(Redis/MySQL 里的用户画像、广告主信息)。假设一次查询 5ms,单个 subtask 串行处理时每秒最多 200 条——这是灾难性的低。
同步外部调用的火焰图特征很明显:大量时间卡在 socketRead / jdbc 栈上。解法是异步 IO(§6)。
4.3 状态访问慢与 GC
有状态算子(窗口、聚合、去重)每处理一条都要读写状态。当状态后端是 RocksDB(大状态场景的默认选择,见 状态与 Checkpoint 篇)时:
- 每次
state.get()可能触发磁盘 IO(RocksDB 是 LSM-Tree,数据在磁盘 + block cache); - 状态越大、block cache 越小,缓存命中率越低,磁盘 IO 越频繁,算子越慢;
- 序列化/反序列化状态本身也有 CPU 开销(RocksDB 存的是字节,每次读写都要 ser/deser)。
此外,JVM GC 也会造成周期性反压:堆内对象过多、大窗口累积、Full GC 停顿期间算子完全不干活,busy 掉零、上游瞬间被反压。这和本地缓存篇里「缓存配太大触发 Full GC」是同一类问题——内存即预算。
4.4 序列化开销
Flink 在算子之间传递数据、读写状态时都要序列化。如果你的数据类型没有被 Flink 的 TypeInformation 高效识别(比如用了泛型擦除后认不出的类型、或没注册的 POJO),就会回退到 Kryo 这个通用但慢得多的序列化器。火焰图里出现大片 Kryo 栈就是信号。解法:用 Flink 原生支持的类型(Tuple、POJO 且字段规范)、给 Kryo 注册类、或对热点类型自定义 serializer。
5. 数据倾斜解法:把热 key 打散
数据倾斜(§4.1)是 AdTech 里最常见也最难缠的反压源,值得单独一节。核心思路一句话:既然一个 key 的数据不能拆到多个 subtask,那就想办法「制造更多的 key」,把热 key 的数据摊开。 这和缓存系列里应对热 Key 的思路一模一样——热点打散。
5.1 两阶段聚合 + 加盐打散
对可结合(associative)的聚合(count、sum、max 这类),最有效的解法是两阶段聚合(local + global),并对热 key 加盐(salting):
两阶段聚合 + 加盐:热 key 拼随机盐分散到多个 subtask 做局部预聚合,再去盐汇总。单点流量降为 1/N,global 阶段只处理已缩量的局部结果。
第一阶段(local):对 key 拼一个随机后缀(盐),keyBy(campaignId + "#" + random(0..N-1)),把原本压向一个 subtask 的热 key 数据分散到 N 个 subtask,各自做局部聚合。第二阶段(global):把局部结果去掉盐、按原始 campaignId 再 keyBy 一次,做全局合并。因为进入 global 阶段的已经是大幅缩量后的局部聚合结果(N 个局部 count 而不是几百万条原始记录),global 阶段不再倾斜。
// 阶段一:加盐 keyBy → 局部预聚合(把热 key 摊到 SALT_NUM 个 subtask)
final int SALT_NUM = 64;
DataStream<CampaignPartial> local = events
.map(e -> new SaltedEvent(
e.campaignId + "#" + ThreadLocalRandom.current().nextInt(SALT_NUM),
e.campaignId, 1L))
.keyBy(se -> se.saltedKey) // 带盐的 key → 分散到多个 subtask
.window(TumblingEventTimeWindows.of(Time.seconds(10)))
.aggregate(new CountAgg()); // 局部 count
// 阶段二:去盐 keyBy → 全局合并(只处理已缩量的局部结果,不再倾斜)
DataStream<CampaignCount> global = local
.keyBy(p -> p.campaignId) // 用原始 key 汇总
.window(TumblingEventTimeWindows.of(Time.seconds(10)))
.reduce((a, b) -> new CampaignCount(a.campaignId, a.count + b.count));
盐的粒度 SALT_NUM 要 ≥ 下游并行度,才能把热 key 真正摊到所有 subtask。进阶做法是只对识别出的热 key 加盐(普通 key 不加,省掉一次 shuffle)——用一个旁路统计出 Top-N 热 key,动态决定谁需要打散。
5.2 Flink SQL 侧的开关
如果用 Flink SQL 做聚合,不用手写加盐,直接开内置的倾斜优化:
table.exec.mini-batch.*:攒微批再聚合,减少状态访问次数;table.optimizer.agg-phase-strategy = TWO_PHASE:自动把聚合拆成 local + global 两阶段;table.exec.state.ttl配合窗口,控制状态膨胀。
对 count(distinct) 这类还有 table.optimizer.distinct-agg.split.enabled = true 专门拆分热点。
5.3 rebalance / rescale:无 key 场景的重分布
如果倾斜不在 keyBy 之后、而是上游数据本身分布不均(比如某些 Kafka 分区数据量远大于其它),或者你在做的是无状态的 map/filter,可以用重分区算子强制均衡:
rebalance():轮询(round-robin)把数据均匀打散到所有下游 subtask——最彻底,但会引入全量 shuffle(跨网络)。rescale():只在本地 TaskManager 内的下游 subtask 间轮询,避免跨网络,代价是均衡范围小一些。shuffle():随机分发。
注意:rebalance 只对无 key或不要求同 key 聚集的场景有效——它会打乱 key 的归属,不能用在需要 keyBy 语义的聚合前。
6. 外部访问优化:AsyncFunction 异步维表 join
§4.2 说过,同步查外部系统(维表 join、写库)会把 subtask 死死拖住。破解它的核心武器是 Async I/O:用 AsyncFunction 把「同步等一个 RTT」变成「同时在途多个请求」。
同步 vs 异步的差距是数量级的:假设一次维表查询 RTT 5ms,同步下单 subtask 每秒最多处理 200 条(串行等待);异步下同一时刻可以有几百个请求在途,吞吐直接拉高一到两个数量级——瓶颈从「串行等待」变成「外部系统的并发承载能力」。
// 用支持异步的客户端(如 Lettuce/异步 JDBC/异步 HTTP),发请求后立刻返回 future
class AsyncDimJoin extends RichAsyncFunction<Event, Enriched> {
private transient RedisAsyncClient client;
@Override
public void asyncInvoke(Event e, ResultFuture<Enriched> rf) {
// 不阻塞:注册回调,RTT 期间该线程继续处理下一条
client.getAsync("dim:campaign:" + e.campaignId)
.thenAccept(dim -> rf.complete(
Collections.singletonList(new Enriched(e, dim))));
}
@Override
public void timeout(Event e, ResultFuture<Enriched> rf) {
rf.complete(Collections.emptyList()); // 超时兜底,别让整条链卡死
}
}
// unorderedWait:不保证结果顺序,吞吐最高;capacity=在途请求上限
DataStream<Enriched> enriched = AsyncDataStream.unorderedWait(
events, new AsyncDimJoin(),
/*timeout*/ 200, TimeUnit.MILLISECONDS,
/*capacity*/ 500);
几个要点:
- 必须用真正异步的客户端:在
asyncInvoke里包一个同步调用 + 线程池不是真异步,只是把阻塞挪了个地方。要用 Lettuce、异步 HTTP、异步 JDBC 这类原生返回 future 的客户端。 unorderedWaitvsorderedWait:不在乎输出顺序就用 unordered,吞吐更高;需要保持 watermark 之间的顺序才用 ordered。capacity(在途上限):太小限制吞吐,太大压垮外部系统——它本身也是一个反压旋钮。- 一定要处理
timeout:否则一个卡住的请求会一直占着 capacity,最终拖垮整个算子。 - 叠一层缓存维表:维表数据读多写少,用 Caffeine 本地缓存 或 broadcast state 把热维表缓存在算子内,命中就不发外部请求——和缓存系列的 L1 思路完全一致。异步 + 缓存是维表 join 的标准组合。
这套「异步 + 攒批 + 缓存」削掉外部访问峰值的思路,和异步化削峰、限流是一脉相承的:都是不让下游被瞬时洪峰打穿,只是 Flink 把它做进了算子模型里。
7. 其它调优旋钮
倾斜和外部访问是两个大头,下面这些是配套的「细调」。
7.1 并行度
- 对齐 source 分区:Kafka source 的并行度 > 分区数时,多出来的 subtask 拿不到分区、纯空转;建议 source 并行度 = Kafka 分区数(或其约数)。
- 给瓶颈算子单独提并行度:不必全作业统一并行度。用
.setParallelism(n)只给瓶颈算子加资源,其它算子维持低并行度省资源。 - 注意 keyBy 后并行度受 key 基数限制:key 只有几个时,提再高并行度也用不满(多数 subtask 分不到 key)——这时候要回到 §5 打散。
7.2 Network buffer
反压严重、吞吐上不去时,可适度增大网络内存(taskmanager.memory.network.fraction / min / max)和每 channel buffer 数,让流水线里能缓冲更多在途数据、吸收抖动。但别盲目调大:buffer 越多,数据在缓冲里排队越久,端到端延迟越高,checkpoint barrier 也要更久才能推进(见 §7.4)——这是吞吐与延迟的经典权衡。
7.3 RocksDB 调优
大状态作业的反压经常来自 RocksDB 的磁盘 IO:
- 给足托管内存:
state.backend.rocksdb.memory.managed = true,让 Flink 统一管理 block cache + write buffer;托管内存不够,block cache 命中率就低。 - 开启 bloom filter:
state.backend.rocksdb.bloom-filter.enabled,减少无谓的磁盘查找(point lookup 不命中时快速返回)。 - 用本地 SSD:
state.backend.rocksdb.localdir指到 SSD/NVMe,别放网络盘或慢盘。 - 增量 checkpoint:
state.backend.incremental = true(RocksDB 特有),只传增量 SST,大幅降低 checkpoint 开销,间接缓解「checkpoint 拖慢作业」。
7.4 算子链、slot 与 checkpoint 的相互影响
- 算子链(operator chaining):Flink 默认把可以链在一起的算子融合进同一个 task 线程,省掉序列化和网络传输——这本身就是优化。但如果链里某个算子是重 CPU 瓶颈,可以用
disableChaining()/startNewChain()把它拆出来单独给并行度。 - Slot 共享:默认同一个 slot 里能跑一条流水线的多个算子,提升资源利用率;资源不均时可用 slot sharing group 隔离。
- 反压与 checkpoint 的死循环(重点):反压时,checkpoint barrier 会被堵在满载的缓冲区里流不动,对齐(alignment)迟迟完不成 → checkpoint 超时失败 → 作业更不健康,形成恶性循环。破解办法是 Unaligned Checkpoint(
execution.checkpointing.unaligned.enabled = true):barrier 不再等对齐,可以超越缓冲区里排队的数据直接快照(把在途 buffer 一起存进 checkpoint)。它让 checkpoint 在重反压下也能按时完成,代价是 checkpoint 体积变大。这是「反压 + 大状态」作业的救命开关,机制细节见 Exactly-Once 篇。
参考
- Apache Flink 官方文档. Monitoring Back Pressure:Web UI 反压监测、busy/backPressured/idle 指标与火焰图的一手说明。
- Apache Flink 官方文档. Network Buffer Tuning & Credit-based Flow Control:网络内存与基于信用的流控参数调优。
- Apache Flink 官方文档. Asynchronous I/O for External Data Access:
AsyncFunction/AsyncDataStream异步维表 join 的权威用法。 - Apache Flink 官方文档. Unaligned Checkpoints:反压下如何用 unaligned checkpoint 保住容错。
- Nico Kruber, Piotr Nowojski. Flink Network Stack: A Deep-Dive & Credit-based Flow Control:Flink 网络栈与 credit-based 反压机制的官方深度博文。
- Apache Flink 官方文档. State Backends & RocksDB Tuning:RocksDB 托管内存、bloom filter、增量 checkpoint 等状态后端调优。