前面几篇把 Flink 的运行时架构、时间与 Watermark、端到端 Exactly-Once、DataStream API 与状态容错 一路拆到了底,讲的其实都是同一件事——正确性:作业算得对、故障能恢复、结果不重不丢。可真把一个流作业推上生产、让它去扛广告晚高峰那几百万 QPS 的曝光洪峰时,先找上门的往往不是正确性,而是性能:吞吐跑着跑着就塌了,端到端延迟从秒级涨到分钟级,Kafka 消费 lag 一路飘红,实时大盘卡在十几分钟前不再更新。绝大多数这类事故,根因都指向同一个词——反压(backpressure)。
本篇聚焦「容错运维」里的性能这一环,沿着「现象 → 机制 → 定位 → 动作」走一遍:反压这个现象到底是什么,credit-based 流控在底下怎么运转,怎么用 Web UI、指标和火焰图把瓶颈算子揪出来,以及针对每类成因该下什么手——尤其是数据倾斜这个 AdTech 里的头号杀手。
一句话定位:反压不是故障,而是 Flink 用「让上游慢下来」来保护自己不被压垮的自愈信号。调优的目标从来不是「消灭反压」,而是顺着它找到那个拖慢全链路的最慢算子并把它救回来——在广告实时链路里,这个最慢算子十有八九是被热 key 打爆的聚合算子。
TL;DR
- 反压 = 下游处理不过来、逐级向上游传导的「减速信号」:某个算子消费速度跟不上,它的输入缓冲被占满,会反向逼停上游,一级级传到 source,最终让 source 放慢甚至暂停从 Kafka 拉取。
- 反压本身是好事:它是 Flink 的天然限流机制,避免数据无限堆在内存里把 TaskManager 撑爆(无反压的系统只会直接 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()重分布——本质就是Redis 热 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 判断下游还收不收得下,反压传导更快、更精准,随 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 是瓶颈,你会看到它上游的 source、A、B 一串全红,而 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(backPressured 高) → map(backPressured 高) → agg(busy 高、backPressured 低) → sink(idle 高),一眼就能锁定 agg 是瓶颈——反压正是从它这里发源、再逐级往上游传导的。
3.3 火焰图:瓶颈算子「慢在哪一行」
锁定瓶颈算子只是完成了一半,还得知道它究竟慢在哪里:是卡在状态访问,是序列化拖后腿,还是某段正则 / JSON 解析吃光了 CPU?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 掉到零,上游瞬间被反压。这和本地缓存的 GC 预算里「缓存配太大反而触发 Full GC」是同一类问题——内存即预算,托管内存、堆内缓存、状态三者是在抢同一块蛋糕。
4.4 序列化开销
Flink 在算子之间传递数据、读写状态时都要做序列化。如果你的数据类型没能被 Flink 的 TypeInformation 高效识别(比如泛型擦除后认不出的类型,或字段不规范的 POJO),就会回退到 Kryo 这个通用但慢得多的序列化器,火焰图里出现大片 Kryo 栈就是明确信号。解法也直接:改用 Flink 原生支持的类型(Tuple、字段规范的 POJO)、给 Kryo 注册类,或者对热点类型手搓一个自定义 serializer。
5. 数据倾斜解法:把热 key 打散
四类成因里,数据倾斜(§4.1)既最常见也最难缠,值得单独拿一节来治。回到那句反复出现的判断:反压只是信号,真正要救的是最慢算子;而在倾斜场景里,这个最慢算子就是被热 key 压住的那个 subtask。核心思路一句话:既然一个 key 的数据不能拆到多个 subtask,那就想办法「制造更多的 key」,把热 key 的数据摊开。 这和Redis 应对热 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,而不是几百万条原始记录。
// 阶段一:加盐 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」换成「同时在途多个请求」。
这中间的差距是数量级的:仍按一次维表查询 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 调优
呼应 §4.3,大状态作业的反压经常直接来自 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 等状态后端调优。