Skip to content
Charles Shao
Go back

Flink 深挖 · 反压定位与性能调优

views

前面几篇把 Flink 的运行时架构时间与 Watermark端到端 Exactly-OnceDataStream状态容错 一路拆到了底。这些讲的都是「正确性」——作业算得对、故障能恢复、结果不重不丢。但真正把一个流作业推上生产、扛住广告晚高峰几百万 QPS 的洪峰时,你面对的第一个敌人往往不是正确性,而是性能:作业跑着跑着吞吐塌了、端到端延迟从秒级涨到分钟级、Kafka 消费 lag 一路飙红。绝大多数这类事故,根因都指向同一个词——反压(backpressure)

本篇聚焦「容错运维」里的性能这一环:把反压是什么、怎么定位、怎么调优(尤其是数据倾斜这个 AdTech 里的头号杀手)讲透。

一句话定位:反压不是故障,而是 Flink 用「让上游慢下来」来保护自己不被压垮的自愈信号。调优的目标从来不是「消灭反压」,而是找到那个拖慢全链路的最慢算子并把它救回来——在广告实时链路里,这个最慢算子十有八九是被热 key 打爆的聚合算子。

TL;DR

Table of contents

Open Table of contents

1. 反压是什么、为什么会发生

流处理是一条算子流水线source → map → keyBy 聚合 → sink,数据像水一样自左向右流。每两个相邻算子之间都有有限大小的缓冲区(network buffer / 队列)。只要下游某个算子的处理速度跟不上上游的产出速度,它们之间的缓冲就会被慢慢填满。

缓冲满了之后会发生什么?这正是反压的核心机制:

Flink 反压逐级传导示意图。一条算子流水线从左到右依次是:Source(Kafka 消费)、Map 算子(解析/清洗)、keyBy 聚合(窗口统计)、Sink 算子(写外部库,处理变慢)。上方标注数据正常时自左向右流动。当最右侧 Sink 算子处理变慢时,下方一条红色的反压传导带自右向左逐段回传:Sink 输入 buffer 被占满,向上游发放 credit=0;聚合算子输出被堵、输入随之满,credit=0 继续向上游传递;Map 算子输入 buffer 满,credit=0 一直传到 Source;最终 Source 拿不到 credit,暂停从 Kafka 拉取,消费 lag 上涨。底部说明这就是基于信用(credit-based)的网络流控:接收端按空闲 buffer 数发放 credit,发送端只在有 credit 时才发数据,下游一慢就逐级把上游逼停。

反压逐级传导:Sink 变慢 → 输入 buffer 满 → 逐段回传 credit=0 → 最终逼停 Source。整条链路的吞吐,由那个最慢的算子决定。

关键认知:反压是一件好事。 设想一个没有反压机制的系统:上游只管猛发、下游处理不过来的数据全堆在内存里,迟早 OOM 崩溃。Flink 的反压机制让整条流水线自动降速到最慢算子的节奏——用「暂时慢下来」换「不崩溃、不丢数据」。所以你在 Web UI 看到红色的反压,不该慌着「消灭它」,而应该顺着它去找那个真正的瓶颈算子

反压真正的坏处在于它的连带效应:一旦发生,端到端延迟会显著上涨(数据在缓冲里排队)、Kafka 消费 lag 累积、Checkpoint 因为 barrier 对齐被卡住而超时(§7 会讲 unaligned checkpoint 如何破解)。这些才是需要治的,而治它们的前提是先定位瓶颈

反压能「逐级精准传导」,靠的是 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 粒度

这样一来:哪个 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,即每秒有多少毫秒处于该状态),这是定位反压的主力

三者理论上加起来约等于 1000。用它们判断瓶颈的黄金法则:

瓶颈算子 = 第一个「busy 高、但 backPressured 低」的算子。 它自己很忙、却没被下游拖着,说明「慢」就出在它身上。它上游的那些算子会是「backPressured 高」(被它拖着),下游算子则是「idle 高」(它喂不饱下游)。

举例:source(idle高) → map(backPressured高) → agg(busy高、backPressured低) → sink(idle高),一眼就能锁定 agg 是瓶颈——反压从它这里发源。

3.3 火焰图:瓶颈算子「慢在哪一行」

锁定瓶颈算子后,还要知道它慢在哪里——是卡在状态访问?序列化?还是某个正则/JSON 解析?Flink 1.13+ 内置了 Flame Graphrest.flamegraph.enabled: true 开启),对选中算子的线程做栈采样,画成火焰图。宽度越大的栈帧 = 占用 CPU 越多的调用。常见「宽帧」:

三层工具连起来用:Web UI 看哪条链红 → busy/backPressured 指标锁定起点算子 → 火焰图看它慢在哪行代码。定位准了,调优才不会瞎试。

4. 常见成因:反压从哪来

瓶颈算子的「慢」通常逃不出这四类。

4.1 数据倾斜(AdTech 头号元凶)

keyBy 会按 key 的哈希把数据分发到下游各个 subtask。理想情况下 key 分布均匀、每个 subtask 分到差不多的量。但现实里 key 的分布经常是长尾的——少数 key 占了绝大多数流量。

Flink 数据倾斜示意图。左侧一个 keyBy(campaignId) 算子按 key 哈希把数据分发到右侧四个聚合 subtask。四个 subtask 的负载严重不均:subtask 0、subtask 1、subtask 3 都很空闲(绿色,各约 18 万条/秒、busy 约 20%),而 subtask 2 是热 key 所在的 subtask(红色,约 760 万条/秒、busy 约 98%),负载条几乎被打满,旁边标注「热 key 打满单 subtask」。底部说明:按 campaignId keyBy 时,某个爆量大广告主的 campaign 会形成超级热 key,数据全被哈希到同一个 subtask,该 subtask 过载并向上游传导反压,而其余 subtask 大量空闲;整个作业的吞吐由这一个最慢的 subtask 决定,加机器、提并行度都救不了单点。

数据倾斜:一个超级热 key 把单个 subtask 打满、其余 subtask 空转。作业吞吐 = 最慢单点吞吐,横向扩容对单 key 无效。

在广告场景里这几乎是必然:按 campaignId / advertiserId 做实时统计时,某个爆量大广告主(比如大促期间猛砸预算的头部客户)的 campaign 会瞬间成为超级热 key,它的曝光/点击事件全被哈希到同一个 subtask。结果:那个 subtask busy=1000、疯狂反压,其余 subtask idle 很高、闲得发慌。这时候加并行度、加机器都没用——因为一个 key 的数据不会被拆到多个 subtask,瓶颈永远是那一个点。解法见 §5。

4.2 外部系统慢:sink 与同步维表 join

第二类元凶是算子里同步调用了外部系统,被外部的 RTT 死死拖住:

同步外部调用的火焰图特征很明显:大量时间卡在 socketRead / jdbc 栈上。解法是异步 IO(§6)。

4.3 状态访问慢与 GC

有状态算子(窗口、聚合、去重)每处理一条都要读写状态。当状态后端是 RocksDB(大状态场景的默认选择,见 状态与 Checkpoint 篇)时:

此外,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)

Flink 数据倾斜解法对比图,分左右两个面板。左面板「① 直接 keyBy 聚合:热 key 压垮单点」(红色主题):热 key campaign=42(该广告主爆量、占比 70%+)的全部记录,向下汇聚到单个聚合 subtask 2,该 subtask 过载,整个作业吞吐等于这个最慢单点的吞吐,其余 subtask 空闲。右面板「② 两阶段聚合 + 加盐打散」(绿色主题):先给热 key 拼上随机后缀 salt(0..3),用 keyBy(campaign + salt) 打散,向下分发到多个 local 预聚合 subtask 并行做局部聚合(局部 count/sum),再去盐、用 keyBy(campaign) 把各局部结果汇聚到一个 global 聚合算子做二次合并。底部说明:热 key 被摊到 N 个 subtask,单点压力降为约 1/N,global 阶段只处理已被大幅缩量的局部聚合结果,因此不再倾斜;代价是引入一次额外 shuffle 与两段聚合逻辑,且只适用于可结合的聚合。

两阶段聚合 + 加盐:热 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,动态决定谁需要打散。

如果用 Flink SQL 做聚合,不用手写加盐,直接开内置的倾斜优化:

对 count(distinct) 这类还有 table.optimizer.distinct-agg.split.enabled = true 专门拆分热点。

5.3 rebalance / rescale:无 key 场景的重分布

如果倾斜不在 keyBy 之后、而是上游数据本身分布不均(比如某些 Kafka 分区数据量远大于其它),或者你在做的是无状态的 map/filter,可以用重分区算子强制均衡:

注意: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);

几个要点:

这套「异步 + 攒批 + 缓存」削掉外部访问峰值的思路,和异步化削峰限流是一脉相承的:都是不让下游被瞬时洪峰打穿,只是 Flink 把它做进了算子模型里。

7. 其它调优旋钮

倾斜和外部访问是两个大头,下面这些是配套的「细调」。

7.1 并行度

7.2 Network buffer

反压严重、吞吐上不去时,可适度增大网络内存(taskmanager.memory.network.fraction / min / max)和每 channel buffer 数,让流水线里能缓冲更多在途数据、吸收抖动。但别盲目调大:buffer 越多,数据在缓冲里排队越久,端到端延迟越高,checkpoint barrier 也要更久才能推进(见 §7.4)——这是吞吐与延迟的经典权衡。

7.3 RocksDB 调优

大状态作业的反压经常来自 RocksDB 的磁盘 IO:

7.4 算子链、slot 与 checkpoint 的相互影响

参考


views
Share this post on:

Previous Post
Flink 深挖 · 部署模式与 Flink on K8s
Next Post
Flink 深挖 · 状态管理与 Checkpoint 容错