Flink 调优的核心,不是把参数调到最大,而是找到最慢的一环并恢复平衡。

前几个月,我们团队一个核心实时作业在凌晨挂了。不是宕机,是比宕机更烦人的那种——它还在跑,但 Kafka Lag 以每分钟几百万条的速度往上蹿。值班同学上来就三板斧:加并行度、扩内存、调大 Checkpoint timeout。折腾了快一小时,Lag 纹丝不动。

最后发现,Sink 端有个用户 ID 占了全量流量的 40%,一个 SubTask 被打满,整条链路跟着憋死。改了下 Key 的分区策略,十分钟恢复。那天晚上我深刻体会到一件事:Flink 调优大多数时候不是资源不够,而是某个环节卡住了,你没找到而已。

这篇不讲原理,前面几篇已经写过了。这里只聊线上踩过的坑,以及我们后来总结的一套排查习惯。

一、先别急着改参数,看看监控

很多工程师(包括我自己)遇到性能问题,第一反应是打开 flink-conf.yaml。这其实是最后一步。我们现在排查前,会先看这几项数据,形成基线:

  • Kafka Lag / Source TPS:如果 Lag 在涨,说明消费跟不上。是 Source 慢,还是下游反压导致的?
  • BackPressure 指标:Flink 1.19 之后 Web UI 的 BackPressure Tab 已经很好用了。哪个算子的比值高,就从哪里开始查。
  • Busy / Idle Ratio:如果大部分 SubTask 的 Busy 不到 30%、Idle 很高,说明并行度可能给高了,数据被切太碎,每个 SubTask 分不到多少活。反过来,如果 Busy 接近 100% 还伴随着 BackPressure,才说明并行度可能不够。
  • Checkpoint Duration:如果一次 Checkpoint 的时间已经接近甚至超过 Interval,那它实际上在持续堆积,早晚要崩。
  • GC 情况:我们用 JMX Exporter 指标上报来观察 GC 时间。Young GC 超过 100ms 或者出现 Full GC,就要警惕了。

没有这些数据,调优就是盲打。我们吃过这个亏——有一次改了半天参数,最后发现是下游 MySQL 整点在跑批处理,跟 Flink 没关系。

二、并行度:不是越大越好,也不是够用就行

一个隐蔽的坑:分区数不匹配

我们曾经遇到一个作业,Kafka Topic 48 个分区,Source 并行度设了 20。看起来没问题,跑起来也正常,但偶尔会有延迟抖动。后来仔细看才发现,20 个 SubTask 分 48 个分区,意味着有些 SubTask 消费 3 个分区,有些只消费 2 个——天然的数据倾斜,跟业务逻辑完全无关。

后来我们把 Source 并行度改成 24(48 的约数),抖动消失了。

我的习惯是:Kafka Source 的并行度最好跟分区数成整数倍关系,1:1 或者 1:2 都可以。如果一定要不成倍数,至少保证不要出现"某些 SubTask 多消费一个分区"的情况。

并行度过高也有代价

另一个极端是并行度设得太大。我们有个作业,工程师为了"保险",把并行度从 24 提到了 96。结果:

  • Shuffle 阶段的网络连接数从 576 涨到 9000+,TM 之间的 Netty 连接打满了
  • RocksDB 的 SST 文件数量暴涨,Compaction 开始跟不上
  • Checkpoint 的元数据变大了,每次 Checkpoint 多花了 30 秒

最后老老实实降回 48,反而跑得更好。

所以我现在调并行度的习惯是:按 1.5x ~ 2x 的步长调,每次只改一处,观察十分钟。别贪心。

Chaining 也要看

Flink 默认会把能 chain 的算子串在一起,这通常是好事,减少了序列化和网络传输。但如果 chain 里面某个算子成了瓶颈,整个 chain 都会慢下来。

有个简单的方法:在 Web UI 的 Job Graph 里,看看每个 Vertex 包了几个算子。如果发现一个 Vertex 里面塞了五六个算子,而它的 BackPressure 很高,可以考虑用 disableChaining() 拆开,定位到底是哪个算子在拖后腿。

三、背压:找到第一个卡住的算子

背压本身不是问题,它是 Flink 的自我保护。你要找的是第一个出现背压的算子——从 Sink 往 Source 方向回溯,那么出现背压的后一个算子就是瓶颈。

怎么快速定位

Flink 1.20中 Web UI 的 BackPressure Tab 里会直接显示每个 SubTask 的 BackPressured Ratio。我的经验是:

  1. 打开 Job → BackPressure
  2. 从 Sink 开始往上看,找到第一个 BackPressured Ratio 明显偏高的算子
    1. 点进去看它的 Busy/Idle:如果 Busy 接近 100%,说明算子本身在吃满;如果 Busy 很低但背压高,说明它在等外部资源(比如数据库、接口),这两种情况的处理方式完全不同。

Sink 慢怎么办

这是最常见的背压原因。我们的处理优先级一般是:

  1. 提 Sink 并行度——先看看是不是并行度跟上游严重不匹配。我们有个作业上游 48 并发,Sink 只有 2 个,单条写入 MySQL 的 RT 大概 5ms,理论峰值也就 400 TPS,撑死也上不去。
  2. 改批量写入——JdbcSinkbatch.size 从 1 改成 500,吞吐能翻几十倍。
  3. 异步化——如果下游是接口调用,用 AsyncFunction 替代同步请求。
  4. 降事务粒度——如果下游支持,把大事务拆小,减少锁持有时间。

网络 Buffer 不足

如果中间算子出现背压,但 Busy 并不高,有可能是跨 TM 的网络 Buffer 不够。特别是 keyByrebalance 这种会触发 Shuffle 的操作。

Flink 1.20 的统一内存模型下,可以调 taskmanager.memory.network.fraction,默认 0.1,Shuffle 大的作业可以提到 0.2~0.25。不过 Network Memory 有上限 taskmanager.memory.network.max(默认 1GB),调 fraction 不会无限涨。另外,Network Memory 和 Managed Memory(给 RocksDB 用的)是两个独立区域,互相不直接挤占,但它们都是从总内存里切出来的,总内存固定的情况下,一方占多了,留给 Task Heap 的空间就少了。

外部系统抖动

有时候背压跟 Flink 完全没关系。我们遇到过下游 MySQL 整点跑批处理,写入 RT 从 5ms 飙到 200ms,Flink Sink 跟着反压。这种时候调 Flink 参数纯属浪费时间,得去协调下游的负载窗口。

四、Checkpoint 超时:别只会调 timeout

Checkpoint 持续超时,很多工程师的做法是把 execution.checkpointing.timeout 从 10min 改成 30min。这其实是掩耳盗铃——它只是让失败来得更晚一些。

我们的排查顺序

第一步:先看有没有背压。

Barrier 是跟着数据流走的。如果链路被背压堵死了,Barrier 传不到 Sink,Checkpoint 的 Alignment 阶段就会一直卡着。Web UI 里看 Checkpoint 的 Alignment Duration,如果持续增大,基本就是背压的锅。

第二步:看状态大小。

在 Checkpointing 页面里,重点看这几个数:

  • Checkpoint Size:总状态多大?如果超过了 Managed Memory 的 50%,就要警惕了。
  • Checkpoint Duration:一次 Checkpoint 花了多久?如果已经占到 Interval 的 80% 以上,说明它在临界运转,迟早要崩。
  • 单个 SubTask 的状态:如果某个 SubTask 的状态是其他的 5 倍以上,说明有倾斜。

状态膨胀的常见原因:

  • 窗口设得太长。我们见过一个作业用 7 天的 TumblingWindow,状态涨到 TB 级,Checkpoint 根本跑不完。
  • Key 设计有问题。比如用 user_id + 毫秒级时间戳 当 Key,Key 空间无限膨胀。
  • 没配 State TTL,过期数据一直赖着不走。

第三步:看 RocksDB 是否撑不住了。

如果状态大,RocksDB 的 Compaction 很容易成为瓶颈。可以在 TaskManager 日志里搜这几个关键词:

Compaction error
Write stall
Too many L0 files

出现的话,说明 RocksDB 在挣扎。可以试着调大 taskmanager.memory.managed.fraction,给 RocksDB 更多内存做缓存和 Compaction。

Unaligned Checkpoint 慎用

Flink 1.11 之后支持了 Unaligned Checkpoint,允许 Barrier 绕过阻塞的数据 Buffer。听起来很美好,但它会把 Buffer 里的数据也快照进去,Checkpoint 大小会暴涨。我们试过在一个状态 80GB 的作业上开这个,Checkpoint 直接涨到 120GB,反而更慢。

我的建议是:状态小于 10GB、背压严重的作业可以考虑开;状态大的,老老实实解决背压问题。

几套常用的 Checkpoint 配置

# 普通作业(状态 < 1GB,背压不严重)
execution.checkpointing.interval: 3min
execution.checkpointing.timeout: 10min
execution.checkpointing.min-pause-between-checkpoints: 1min
state.backend.incremental: true

# 大状态作业(> 10GB)
execution.checkpointing.interval: 5min
execution.checkpointing.timeout: 15min
state.backend.incremental: true
taskmanager.memory.managed.fraction: 0.5

# 背压严重但状态小
execution.checkpointing.interval: 5min
execution.checkpointing.unaligned.enabled: true
execution.checkpointing.max-aligned-checkpoint-size: 1mb

五、CPU 很低但任务慢:作业可能在"装死"

CPU 利用率低不等于系统空闲。恰恰相反,它往往意味着作业在等待——等 IO、等网络、等下游确认。几种常见情况:

现象 大概率原因 怎么处理
CPU < 30%,TPS 也低 并行度给少了 提并行度
CPU < 30%,但有背压 外部调用慢 / IO 等待 异步化、批量、优化下游
CPU 周期性跌到零 GC 停顿 看 GC 日志,调内存或 GC 策略
个别 SubTask CPU 高,其他很低 数据倾斜 打散 Key

GC 是个隐形杀手

Flink 默认用 G1GC,大部分场景够用了。但大状态作业容易踩坑——RocksDB 的 Block Cache 和 Write Buffer 占的是堆外内存,但 JNI 调用和元数据在堆内。如果堆给小了,元数据区会频繁触发 Full GC。

我们一个作业原来堆给 4GB,频繁 Full GC,每次停顿 2~3 秒,直接导致背压。调到 8GB 之后,Full GC 消失了。

另外,有些算子会产生大量短生命周期对象(比如复杂 JSON 解析),Young GC 会特别频繁。这种情况可以考虑开启对象重用:

env.getConfig().enableObjectReuse();

但要注意:开启后如果下游算子改了对象内容,上游数据会被污染。只在纯转换(Map/Filter 无修改逻辑)的场景用。

六、数据倾斜:最隐蔽的杀手

一个热 Key 能让整个作业的吞吐掉 80%。而且它不报错,只是慢,很难发现。

怎么发现

最直接的方法:在数据倾斜页查看累积情况,执行图中的 Data Skew 看实时情况

低版本也可以在 Web UI 里看每个 SubTask 的 Records Received。如果最大值和最小值差 5 倍以上,大概率有倾斜。

另一个方法是在 Kafka 侧看分区 Lag,如果某些分区的 Lag 明显大于其他分区,说明 Source 端就已经有倾斜了。我们还干过一件事:在 KeyedProcessFunction 里按 Key 上报计数到 Prometheus,直接在 Grafana 里看 Key 分布。虽然有点重,但对排查疑难杂症很有用。

怎么治

两阶段聚合是最常用的方法,适用于 Sum/Count 这类可交换结合的聚合:

// 第一阶段:加盐打散
DataStream<Metric> partial = stream
    .map(new RichMapFunction<Metric, Metric>() {
        private Random random;
        @Override
        public void open(Configuration parameters) {
            random = new Random();
        }
        @Override
        public Metric map(Metric value) {
            value.setSaltKey(value.getUserId() + "_" + random.nextInt(10));
            return value;
        }
    })
    .keyBy(Metric::getSaltKey)
    .window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
    .aggregate(new PartialSumAggregate());

// 第二阶段:去掉盐值,全局聚合
DataStream<Metric> global = partial
    .map(value -> {
        value.setUserId(value.getSaltKey().split("_")[0]);
        return value;
    })
    .keyBy(Metric::getUserId)
    .window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
    .aggregate(new FinalSumAggregate());

如果倾斜只集中在少数几个 Key 上,也可以把热 Key 单独分流处理:

Set<String> hotKeys = loadHotKeysFromRedis();

DataStream<Metric> hotStream = stream
    .filter(x -> hotKeys.contains(x.getUserId()))
    .keyBy(x -> x.getUserId() + "_" + ThreadLocalRandom.current().nextInt(5))
    .process(new HotKeyHandler());

DataStream<Metric> normalStream = stream
    .filter(x -> !hotKeys.contains(x.getUserId()))
    .keyBy(Metric::getUserId)
    .process(new NormalHandler());

最根本的方案当然是重新设计分区键。比如把单字段 user_id 改成 user_id + device_id 的组合键,或者按日期做二级分区。但这需要动业务逻辑,不是总能落地。

七、几个容易忽视的点

序列化

Flink 默认回退到 Kryo 序列化器效率不高。我们在一个 TPS 十万级的作业里,把 Kryo 换成 Avro,CPU 直接降了 20% 左右。序列化在网络 Shuffle 和 Checkpoint 里的开销,很容易被低估。

网络

跨 TM 的数据传输走 Netty,默认配置在小集群够用了。但如果集群规模大、Shuffle 多,可以调一下 taskmanager.memory.network.fraction。不过别调太高,会抢 RocksDB 的内存。

监控埋点

最后说一句,好的调优离不开好的监控。我们团队现在每个 Flink 作业都会埋这几类 Metric:

  • 业务指标:TPS、Latency、Lag
  • 资源指标:CPU、内存、GC
  • Flink 原生指标:BackPressure、Busy/Idle、Checkpoint Duration/Size
  • 自定义指标:Key 分布、外部调用 RT

没有数据,再好的调优经验也是瞎子摸象。

八、复盘两个真实的坑

案例一:Kafka 10 万 TPS,落库只有 2 万

我们的实时推荐系统,从 Kafka 读用户行为事件,清洗后写 MySQL。

下午两点业务放量,Kafka TPS 从 3 万飙到 10 万。Flink Source 端消费没问题,但 Sink 端开始背压,Kafka Lag 以每分钟 200 万条的速度累积。

值班同学第一反应是把作业并行度从 24 提到 48。没用。又加了两台 TaskManager。还是没用。

四十分钟后,有人打开 Web UI 看了一眼 Sink 算子——并行度是 2。上游 48 个并发往 2 个 Sink 里灌,每条写入 MySQL 的 RT 大概 5ms,理论峰值就 400 TPS,怎么可能撑得住 10 万?

后来把 Sink 并行度提到 16(跟 MySQL 表的分区数对齐),开启批量写入 batch.size=500,连接池从 10 调到 50。十五分钟后 TPS 恢复到 9 万+,Lag 开始收敛。

教训:加机器之前,先确认瓶颈到底在哪里。并行度不是只调作业整体的,Sink 的并行度也要看。

案例二:Checkpoint 持续失败,作业没法升级

实时风控系统,用 Flink CEP 做规则匹配,RocksDB 状态大概 80GB。

早上九点触发升级,改了个 CEP 规则,作业重启。然后 Checkpoint 就开始出问题——Duration 越来越长,第一次 10 分钟 Timeout,调大到 30 分钟还是 Timeout。

排查下来,RocksDB 状态在持续增长,已经飙到 80GB 以上。原因是 CEP 模式里用了 24 小时的超时窗口,大量过期匹配没清理,状态无限膨胀。

临时方案:开启增量 Checkpoint,Interval 从 3 分钟调到 5 分钟,先把 Checkpoint 稳住。

长期方案:给 CEP 模式加了 .within(Time.hours(6)) 限制,同时配了 State TTL,过期数据自动清理。清理完历史状态重新部署,Checkpoint Duration 回到 4 分钟左右。

Pattern<AlertEvent, ?> pattern = Pattern
    .<AlertEvent>begin("start")
    .where(evt -> evt.getType().equals("SUSPICIOUS"))
    .next("middle")
    .where(evt -> evt.getAmount() > 10000)
    .within(Time.hours(6));  // 原来没限制,默认一直等

StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Time.hours(12))
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
    .cleanupIncrementally(10, true)
    .build();

教训:Checkpoint 超时,先别调 timeout。看 Barrier 能不能对齐、看状态是不是在膨胀、看 RocksDB 是不是在挣扎。调 timeout 是最后的选择。

九、最后说几句

Flink 调优没有什么万能公式。每个作业的瓶颈可能都不一样——有的卡在 Sink,有的卡在状态,有的卡在一条 SQL 语句上。

我现在的习惯是:改参数之前,先在 Web UI 里待三分钟,把 BackPressure、Busy/Idle、Checkpoint 这几个Tab 看一遍。 大部分时候,瓶颈在哪里已经很明显了。

还有就是,每次只改一个参数,观察十分钟。别同时调并行度、调内存、改 Checkpoint 配置,那样就算好了,你也不知道是哪个改动起的作用。

最后,把之前说的整理成一个 checklist,调优前过一遍:

  • [ ] 确认了当前 TPS、Lag、BackPressure 基线
  • [ ] 定位了第一个出现背压的算子
  • [ ] 检查了它的 Busy/Idle 和 State Size
  • [ ] 确认没有数据倾斜
  • [ ] 检查了 Checkpoint Duration 和 Interval 的比例
  • [ ] 确认了 GC 时间正常
  • [ ] 验证了并行度跟 Kafka 分区数匹配
  • [ ] Sink 并行度跟下游吞吐能力匹配
  • [ ] 每次只改一个参数

 

Logo

汇聚全球AI编程工具,助力开发者即刻编程。