接手过 Flink 线上任务的人,大概都经历过这样的场面:半夜手机连续震动,爬起来打开监控面板,checkpoint 已经红了大半个小时,Kafka 消费延迟像体温计一样往上窜,任务重启记录跳了七八次,另一边数据倾斜导致几个 subtask 的 CPU 飙满、另外几个却在摸鱼。这种时候最怕的就是病急乱投医,今天调个参数、明天加个并行度,表面上好像缓解了,实际上过几天又复发。
这篇 Day19-20 的内容,我把 Flink 线上最常见的四类故障——CK超时、任务重启、Kafka积压、数据倾斜——放在一起做一次系统梳理。每一类都按“现象判断 -> 根因分析 -> 排查步骤 -> 解决方案 -> 应急预案”的顺序来展开,涉及到的参数、命令和配置都会给出实际可用的版本。适合正在维护 Flink 生产任务的开发、平台运维和数仓工程师,也适合刚接触实时计算、想了解线上问题到底长什么样的新手。坦白说,这四类问题单独出现时都还好办,怕的是它们连在一起,比如倾斜引发了反压、反压又把 checkpoint 拖超时,你根本分不清谁是因、谁是果。
1. 线上问题排查的整体思路:先分清因果,再动手改参数
1.1 四类故障的内在联系
我刚开始维护 Flink 集群的时候,遇到故障第一反应就是翻 StackOverflow,然后照着帖子调参数,结果越调越乱。后来带我的老前辈说了一句让我印象非常深的话:线上故障从来不跟你讲武德,它只给你一个表象,但真正的原因往往藏在另外一层。
这四类故障的关系,我用一个实际案例来说明。有一次我们的订单实时分析任务报警 Kafka 积压严重,我第一反应是加并行度。结果加了之后积压没解决,反而把 checkpoint 搞到超时,紧接着任务开始重启。回头看日志才发现:真正的根因是订单数据里某个商家的 key 特别多,导致数据倾斜,倾斜的那个 subtask 处理不过来,反压向上游传导,读 Kafka 的速度自然就慢了下来。积压只是结果,倾斜才是诱因。而我盲目加并行度,导致 barrier 在倾斜的链路里更难对齐,checkpoint 反而超时了。
这个案例说明,排查的第一步不是修,而是判断。要搞清楚当前问题是“果”还是“因”。
| 故障现象 | 常见诱因 | 可能被触发的次生故障 |
|---|---|---|
| Checkpoint 超时 | 数据倾斜、反压、存储慢、GC 频繁 | 任务重启(超时触发 failover) |
| 任务重启 | 代码异常、OOM、心主机超时、CK 超时 | Kafka 积压(停机导致消费停滞) |
| Kafka 积压 | 反压、并行度不足、下游写库慢、消费不均 | Checkpoint 超时(积压加剧延迟) |
| 数据倾斜 | key 分布不均、join 倾斜、窗口数据集中 | 反压、CK 超时、部分节点过载 |
所以我个人排查时有一套固定的顺序:先看任务是否在稳定运行,再看背压和繁忙度,然后看 checkpoint 耗时和频次,最后分析 key 分布和记录数在各 subtask 上的偏差。这套顺序能避免“头痛医头、脚痛医脚”,也是本文全篇的一个基本框架。
1.2 排查前必须做好的三手准备
有人说,排查问题不就是打开日志看吗?这话对,也不对。等你在生产环境里面对几百万行日志时,才发现“看日志”是最后一步,前提是你提前准备好了解析日志的工具和判断基线。
第一手准备是监控指标要有历史数据。Flink 提供了丰富的 Metrics,包括numRecordsInPerSecond、checkStartDelay、alignedProcessingTime、busyTimePerSecond等,但这些指标光有当前值没用,必须有历史曲线,才能判断“异常”是从什么时候开始的、变化是突增还是渐变。没有历史曲线,你连“是不是这个版本升级引入的回归”都看不出来。
第二手准备是日志收集和关键字告警。生产上至少要做到 JobManager 和 TaskManager 日志统一收口到一套日志中心,并且对Exception、CheckpointTimeoutException、OutOfMemoryError、Connection refused这些关键字配置告警。很多任务其实提前给出了征兆,只是你没看到。
第三手准备是明确任务的血缘和拓扑图。排查的时候必须知道自己这条任务的 source 在哪、sink 在哪、中间有哪几个关键算子。我见过太多人排查了半天,才发现瓶颈根本不在 Flink 内部,而在下游 ClickHouse 写不进去了。搞清楚链路边界,有时候比搞清楚代码逻辑更重要。
2. Checkpoint 超时:最常见的“隐性杀手”
2.1 CK 超时的核心机制与判断标准
Checkpoint 超时并不是一个很复杂的机制:JobManager 每隔一定时间间隔(checkpoint.interval)发起一次 barrier 对齐,所有算子完成状态快照之后,checkpoint 才算完成。如果整个过程超过了checkpoint.timeout(默认 10 分钟),这次 checkpoint 就会被判定为失败。连续失败多次之后,根据配置的重启策略,任务会自动重启。
理解了这个机制,你就明白了 CK 超时本质上是由两件事决定的:barrier 能不能及时走完,以及状态快照能不能快速落盘。
在 Flink Web UI 里面,超时通常表现为 checkpoint 历史里大面积红色失败记录,监控面板上alignedDuration(对齐耗时)和checkpointDuration(快照总耗时)指标明显升高。我判断超时的严重程度,一般看两个阈值:单次 checkpoint 耗时超过 checkpoint interval 的 70%,或者连续 3 次失败,就需要停下来认真查。低于这个水平,偶尔抖动可以先观察。
很多同学一看到超时,第一反应就是把checkpoint.timeout调大,比如 10 分钟改 30 分钟。我跟你说,这个操作可以作为临时缓兵之计,但绝不能作为根治方案。超时的根因不解决,把超时时间调再大,也只会让失败暴露得更晚、任务恢复得更慢。
2.2 按图索骥:CK 超时的三类根因与处置
我总结了实际生产中反复出现的三类根因。
第一类是反压导致的 barrier 传播延迟。barrier 要随数据流向下游传播,如果某个算子处理速度跟不上,barrier 就会被卡在缓冲区里迟迟无法到达下游。排查方法是先看 Web UI 里每个算子是否出现背压,或者看busyTimePerSecond是否长期高于 80%。如果是,按数据倾斜和反压的解法来处理,这一篇后面会专门讲。
第二类是状态后端持久化慢。这时候 barrier 对齐没问题,但状态落盘卡住了。常见原因包括:RocksDB 所在的本地磁盘 IO 高、Checkpoint 存储目录所在的 HDFS 小文件多导致写入慢、S3 等对象存储带宽受限。排查方法是在 TaskManager 日志里看快照耗时明细,或者挨个检查存储端的 IO 和文件数量。处理办法包括给 RocksDB 换 SSD、合并 HDFS 上的小文件、把 checkpoint 存储从 HDFS 切到与计算节点同地域的 S3 等。
第三类是 GC 问题。Flink 的堆内存如果频繁 Full GC,整个 TaskManager 的线程都会暂停,barrier 自然就走不动了。我踩过最典型的一个坑是:使用 HeapStateBackend(虽然现在官方不推荐生产使用)存了大 key 的数据,结果每次快照都要序列化大量对象,频繁触发 Full GC。换到 RocksDB 之后,GC 压力明显下降。排查可以通过监控看 Full GC 的频次和耗时,也可以加-XX:+PrintGCDetails看详细日志。
这里我给出一份可直接参考的参数基线:
| 参数 | 推荐值 | 说明 |
|---|---|---|
execution.checkpointing.timeout | 10-15 分钟 | 过大只会拖延问题暴露时间 |
execution.checkpointing.interval | 1-5 分钟 | 根据数据量和恢复时间目标定 |
state.backend.rocksdb.localdir | SSD 目录 | 避免与系统盘互相干扰 |
taskmanager.memory.managed.fraction | 0.4-0.6 | RocksDB 可用内存比例 |
state.checkpoints.num-retained | 5-10 份 | 保留过多会占用存储 |
提示:调参之前先确认哪一环节慢。我要强调的是,不管哪一类根因,都要先做一次“单次 checkpoint 耗时拆解”,看时间花在 aligned 还是 snapshot 上,这能帮你少走一半弯路。
3. 任务频繁重启:从日志到恢复策略的完整链路
3.1 按重启类型快速定位
Flink 任务重启不是随机事件,它背后一定有一个触发源。我把重启分为三类,每一类的排查路径差别很大。
第一类是逻辑异常导致的定时失败。比如某条脏数据触发了 NullPointerException、反序列化失败,或者下游连接池耗尽导致 Sink 写入报错。这类重启的特点是:日志里有明确的异常堆栈,重启时间和异常出现时间高度吻合。排查方法很简单,去 JobManager 日志里搜Job has been submitted之前最近的异常。
第二类是资源不足导致的被动退出。常见的有容器内存超限被 YARN/K8s kill、堆外内存溢出、文件描述符耗尽。这类重启的特点是 TaskManager 日志里可能没有 Java 异常,但从系统层面能看到 OOM-Killer 的记录。排查时使用dmesg -T | grep -i killed查看内核日志,或者看监控里的内存水位。解决办法包括调大 TaskManager 内存、优化代码减少对象分配、排查堆外内存泄漏。
第三类是checkpoint 失败导致的主动 failover。Flink 会按照restart-strategy的配置,在 checkpoint 连续失败后自动重启任务。这类重启的特点是在日志里能看到Checkpoint timeout或者CheckpointCoordinator的相关报错,重启前 checkpoint 已经失败了好几次。
我个人的经验是,第一眼看 ExitCode,第二眼看堆栈,第三眼才查配置。不看堆栈就去调重启策略,那就是把炸弹藏起来,不是拆炸弹。
3.2 常见的重启策略配置与踩坑记录
Flink 提供的重启策略有 fixed-delay、failure-rate 和 exponential-delay 三种。线上最稳妥的配置方案是按任务的“重要程度”和“优雅恢复成本”来选。
对于核心交易类任务,我建议用 fixed-delay,间隔给长一点,比如 60 秒,给外部依赖恢复留足时间。对于下游有幂等保护的任务,可以用 failure-rate 限制单位时间内的重启次数,比如 10 分钟最多重启 3 次,超过就标记为失败,避免无限重启造成更大的数据延迟和下游压力。
这里的配置如下:
restart-strategy: failure-rate restart-strategy.failure-rate.max-failures-per-interval: 3 restart-strategy.failure-rate.failure-rate-interval: 10 min restart-strategy.failure-rate.delay: 30 s有个问题你得注意:重启策略不是“容灾”,只是兜底。如果任务每 30 秒重启一次,就算策略允许,Kafka 积压也一定会失控。所以我的底线是,同一个任务连续重启超过 5 次、且每次存活不超过 5 分钟,就停止靠自动恢复,直接走人工介入流程,把任务先停了,分析原因修完代码再恢复。
另外,从 Flink 1.15 之后,exponential-delay 策略逐步被推荐,因为它支持“启动慢时退避更长,稳定后恢复探测”的节奏,比较适合需要长时间恢复外部依赖的场景。但它的参数比 fixed-delay 多,上线前要在一个模拟环境里验证好。
4. Kafka 积压:从消费延迟反推瓶颈链路
4.1 先判断积压是“上游问题”还是“下游问题”
Kafka 积压本身不是 Flink 故障,而是 Flink 链路对外部世界呈现出来的“症状”。排查的第一步必须是先回答一个问题:是 Kafka 生产端发不出数据,还是 Flink 消费端读得慢?
判断方法很简单,看两个指标的组合。第一个是records-lag-max,表示当前积压消息量;第二个是 Flink 侧每个 subtask 的currentFetchEventTimeLag或numRecordsInPerSecond。如果积压在增长,但 Flink 读取速率已经接近分区数允许的上限,说明 Flink 侧就是瓶颈,需要往 Flink 内部查。如果 Flink 读取速率很低,但 Kafka 分区也没有异常,多半是反压导致读取被暂停,本质是下游处理不动了。
有一回我们的任务积压了 5 个小时的数据,我盯着 Flink Web UI 看了半天没看出问题,最后发现是下游 ClickHouse 在备份期间写入非常慢,把积压传给了 Kafka。所以这里必须强调一个容易被忽略的原则:积压排查的终点,不一定在 Flink 进程内,很可能是下游存储成为瓶颈。
4.2 消费链路优化和安全恢复方案
确定瓶颈方向之后,再来谈优化。如果是 Flink 内部问题导致消费慢,常见解决路径如下:
- 检查是否目标算子存在反压。在 Web UI 看到某个算子被压成红色,优先排查这里。
- 调整并行度。Kafka Source 的并行度理论上不超过分区数,如果并行度小于分区数,先增加 Source 并行度。如果并行度已经等于分区数,再加也无效,要查内部的 keyBy 重分布。
- 优化序列化和对象分配。比如使用
TimestampAssigner时避免在每条数据上创建新对象,PojoSerializer会退化等等,这些都会拖慢消费速率。 - 对 Sink 启用批量写入。例如 JDBC Sink 打开
batchSize,HBase Sink 用 BufferedMutator。
关于积压期间的恢复,我一直是这么操作的:先确认任务的 checkpoint 是否还能正常完成。如果能,优先使用savepoint 手动重启并恢复,避免从头消费 Kafka 导致下游被打垮。如果不能,就先把 Sink 断掉,让 Flink 只消费数据不写入外部,等积压水位降到可控范围再恢复写。
注意:任一情况下,都不要一上来就把并行度翻倍甚至翻三倍。积压恢复后,数据洪峰到达下游,容易把数据库击穿。恢复期间监控下游存储的写入延迟,比监控积压本身更重要。
5. 数据倾斜:最难缠的分布式顽疾
5.1 用“记录数偏差”和“繁忙度曲线”识别倾斜
数据倾斜的本质是某几个 key 的数据量远大于其他 key,导致同一算子下不同 subtask 的负载天差地别。它的排查诊断相对直接,前提是你平时习惯看两种图:每个 subtask 的numRecordsInPerSecond,以及每个 Task 的busyTimePerSecond。
判断标准比较收敛:如果一个 Task 的多个并行 subtask 里,某个(或某几个)的输入记录数是同组其他实例的 3 倍以上,而且这个趋势保持了好几个窗口,基本可以认定存在倾斜。还有一种隐蔽的表现:所有 subtask 的输入记录数看起来都一样,但其中一个 busy 时间明显偏高,这种情况多半是窗口计算里某个 key 携带的 state 特别大,处理单条记录的耗时远超平均。
我在实际项目中见过一次很有意思的倾斜:数据源的 key 是用户 ID,绝大部分用户的事件量差别不大,但有一个测试账号在刷数据,单日事件量是普通用户的几百倍。这种“尖刺型”倾斜不像“偏态型”那么好发现,必须靠分组统计 Top N key 来做。
5.2 三种经典场景的解法与代码示意
第一类是 groupBy 聚合中的倾斜。处理方案首选两阶段聚合,思路是先在 key 后面加上随机前缀打散分区,做一轮 partial agg,再去掉前缀做最终聚合。Flink SQL 中其实内置了 Local-Global 聚合,配置table.optimizer.agg-phase-strategy=TWO_PHASE就能开启,但要注意它只在窗口和状态分区命中原子的场景有效。
用 DataStream API 实现两阶段聚合,核心逻辑大致是这样:
// 一阶段:加盐聚合 keyedStream .map(record -> new Tuple2<>(record.getKey() + "#" + random.nextInt(100), record.getValue())) .keyBy(Tuple2::f0) .reduce((v1, v2) -> /* 分阶段聚合 */) // 二阶段:去盐再聚 .map(record -> new Tuple2<>(stripSalt(record.getKey()), record.getValue())) .keyBy(Tuple2::f0) .reduce((v1, v2) -> /* 最终聚合 */);第二类是 join 中的倾斜。最经典的场景是维表关联时热点 key 打到一个节点上。三种常见应对方式值得写在笔记里:把维表做成广播流,让每个 partition 各持一份;对热点 key 分流后单独和倾斜维表分片做二次 join;或者提前按热点维度对维表进行分桶预处理。线上优先级一般是广播 > 分桶 > 二次 join,因为广播最简单,代价是内存占用。
第三类是 window 中的倾斜。如果数据量集中在个别 key,比如秒杀场景里一个商品的点击流远大于其他商品,直接开窗口会导致该 key 所在 subtask 成为瓶颈。除了加盐分散之外,还可以把 window 的 trigger 和 evictor 单独定制,对热点 key 做高频 partial 输出,再外部合并。
需要特别提醒的是,加盐把数据打散之后,如果后面还有依赖原始 key 的最终结果,必须保证第二阶段能完整拿到同一个原始 key 的全量数据。我见过不少同学加盐之后第二阶段忘了去盐,最后结果的粒度要么丢掉明细、要么错误聚合,比倾斜本身还麻烦。
6. 排查工具箱与实战经验汇总
6.1 常用命令与指标速查
最后把我平时排查问题会用到的工具和命令整理成一张速查表,方便大家在现场快速对照。
| 场景 | 关键指标 / 命令 | 说明 |
|---|---|---|
| 反压 | Flink Web UI Backpressure 面板 | 便捷的实时反压视图 |
| 倾斜 | numRecordsInPerSecond分组对比 | 对比不同 subtask 的速率 |
| 频繁 GC | jstat -gcutil <pid> 1000 10 | 查看 YGC/FGC 频率 |
| 堆内/堆外内存 | jmap -heap <pid>+ 容器内存监控 | 排查 OOM 场景 |
| Kafka 积压 | Kafka Consumer Group 的records-lag | 结合 lag 变化趋势判断 |
| CK 超时 | Checkpoint 历史详细页 | 看耗时分钟段 |
| 内核 OOM | dmesg -T | grep -i killed | 定位被 K8s/YARN kill 的实例 |
顺便说一句,这套工具的组合使用非常看场景。比如你看到一个任务反压很严重,但busyTimePerSecond不高,那大概率不是计算慢,而是存在同步等待(比如等下游写库、等外部接口返回)。这时候你去看mailbox或者request的等待耗时,比盯着 CPU 更有效。
6.2 我曾经踩过的几个“逻辑陷阱”
排查线上问题多了,你会发现最坑人的不是技术难度,而是你以为自己已经找到了真相,实际上只是看到了一层表象。我梳理了三个我踩过的“逻辑陷阱”,供各位避坑。
第一个是**“加内存能解决一切”**。有一段时间我们的一个任务频繁重启,排查发现是堆外内存溢出,于是直接给 TaskManager 加了 4G 内存。结果第二天又炸了。后来仔细排查才发现,问题出在一个 ListState 不断往里 add,数据无限增长,堆内空间不足以容纳完整快照,并且序列化时间过长导致 CK 超时。内存只是帮垃圾桶装更多垃圾,真正要做的是清理垃圾——把无限增长的 ListState 改成 TTL 窗口存储或按分钟清理。
第二个是**“积压严重就加并行度”**。前面说了,积压的根因可能是下游写不动。这种情况下加并行度只会让下游更快被打爆,形成恶性循环。我后来养成的习惯是:任何优化动作之前的 10 分钟,先看下游存储的写入耗时和线程池积压情况,确认下游健康,再动 Flink 侧的并行度。
第三个是**“重启就好使”**。有一次任务的失败原因是外部 Redis 连接超时,手动重启后任务恢复了,但第二天同样的问题又出现了。原因很简单,代码里拿 Redis 写缓存时没有重试机制,超时一次就抛异常。重启只能暂时绕开当时的抖动,不能修复代码的健壮性。后来我要求所有外部依赖调用必须设置合理的超时时间和重试策略,这才真正降住了故障率。
7. 写在最后的实战体会
我在实际运维 Flink 的这几年里,最大的感受是:绝大多数线上问题不是靠“绝顶聪明”解决的,而是靠“把基础功做扎实”解决的。所谓基础的功夫,指的就是每天盯着监控面板看 10 分钟、每次上线前把 checkpoint 和重启策略的配置逐项过一遍、每次故障复盘后把排查结论沉淀到文档里。这套活儿看起来琐碎,但正是这些琐碎的东西,决定了你在凌晨三点被电话叫醒时,是胸有成竹地几步定位问题,还是手忙脚乱地到处乱翻日志。
另外还有一个小技巧值得分享:每处理完一次线上故障,我都会新建一个笔记,记录现象、初步判断、排查过程、最终根因和参数变更。几个月下来,我发现至少有 30% 的“新故障”其实是旧问题换了件马甲。把历史案例库做好,面对新问题时你可以第一时间找到相似案例做参照,比从零开始排查快得多。
这套 Flink 排查体系我已经在多个实时任务上验证过,覆盖了从几万条每秒到几十万条每秒的不同规模场景。你如果刚开始维护 Flink 任务,不必一上来就追求精通所有细节,先把我说的“先判因果、再看监控、再动参数”这个顺序刻在脑子里,配合这篇里的具体操作步骤去实践,慢慢就能形成自己的排查手感。如果后续你在实际排查中遇到了这里没覆盖到的特殊情况,我很乐意继续分享更多的排查实例。