ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Flink线上故障排查实战:CK超时、Kafka积压与数据倾斜的根因定位指南

Flink线上故障排查实战:CK超时、Kafka积压与数据倾斜的根因定位指南 凌晨两点被告警电话叫醒打开Flink Web UI一看Checkpoint超时飘红、作业卡在RESTARTING、Kafka LAG曲线直线拉升——这个画面估计每个维护Flink线上作业的人都不陌生。我搞Flink生产环境运维这五年发现所谓的线上故障来来回回就是几个老面孔CK超时、任务反复重启、Kafka积压、数据倾斜。而且更坑的是这四个问题经常不是独立出现的而是像多米诺骨牌一样连着倒数据倾斜导致某些subtask处理不过来继而Kafka积压积压引发反压反压让Barrier没法对齐最终Checkpoint超时任务开始重启……如果只盯着表面现象去救火很可能按下葫芦浮起瓢。这篇文章就把我这几年在线上排查这类问题的完整思路和工具链梳理一遍不写空对空的理论全部是从告警触发到根因定位再到修复验证的实操链路。内容包括Checkpoint超时到底是卡在对齐还是卡在持久化、任务重启时如何快速区分该崩溃的异常和不该抖动的系统问题、Kafka积压时怎么定位是消费能力不足还是上游压力传导、以及数据倾斜的几种典型形态和对应的加盐/预聚合处理方案。最后用一个真实复盘把四个故障的联动关系串起来讲清楚。1. CK超时先分清是Barrier对齐慢还是状态持久化慢Checkpoint超时是Flink线上最容易被误判的问题之一。很多人在Web UI上看到Checkpoint失败第一反应是状态后端是不是有问题或者HDFS是不是挂了但实际上大部分CK超时的根因根本不在持久化这一层而是在Barrier对齐这一层。1.1 从Checkpoints面板读出关键线索排查CK超时第一步是打开Web UI的Checkpoints面板不要急着看Summary先看最近几次失败的Checkpoint详细情况把几个关键指标抓出来End to End Duration整个Checkpoint从开始到完成的总耗时。如果这个值动不动就超过分钟级通常不是持久化慢而是Barrier等待时间过长。Checkpointed State Size状态大小。这个值决定了持久化阶段的理论下限如果状态几个GB那同步阶段耗时几十秒甚至几分钟都算正常。Failures列表里的具体异常信息是Checkpoint expired before completing还是Task not found还是Exception while performing checkpoint。我见过太多人一看到Checkpoint expired before completing就怀疑RocksDB写入慢实际上这个异常在绝大多数情况下是因为Barrier没能在超时时间内走完整个拓扑——说白了就是有数据反压Barrier被堵在某个算子前面排队。判断技巧把鼠标点到某个Checkpoint的详情里能看到每个Task的Checkpoint耗时里面通常会拆成Sync Duration和Async Duration两段。如果Async Duration异步持久化阶段很短而整个Checkpoint的End to End Duration很长那问题几乎可以锁定在Barrier对齐阶段。1.2 Barrier对齐慢反压是最大的元凶Flink的Checkpoint机制是依靠从Source端 inject 的Barrier沿着数据流一直往下走的。Barrier走到每个算子算子要等所有输入channel的Barrier都到了之后才开始做快照。如果某个channel前面有大量数据堵着Barrier只能排队慢慢挪这时候Checkpoint就卡住了。这就是为什么说反压是Checkpoint超时最常见的推手。数据倾斜比如某个key的数据量是其他key的几十倍、下游Sink写入慢、某个算子计算逻辑复杂这些都会造成反压反压一来Barrier就走不动CK就超时。解决思路有两个方向预算和资源充足的话首先做拓扑调优和并行度调优把反压源头掐掉。如果短期内没法改拓扑可以考虑开启Unaligned Checkpoint。Flink 1.11之后Unaligned Checkpoint已经成为生产环境比较成熟的选项它允许Barrier不需要等所有channel对齐超时N毫秒之后直接越过队列中的数据做快照。注意Unaligned Checkpoint不是银弹。开启后Barrier会携带沿途的buffer数据一起持久化导致Checkpoint的状态和文件体积变大。我在实测中发现数据量大时Checkpoint文件体积可能翻几倍对于状态后端到远端存储的网络带宽会形成额外压力。参数是这样调的在flink-conf.yaml里execution.checkpointing.unaligned: true execution.checkpointing.aligned-checkpoint-timeout: 30s含义是前30秒先尝试对齐超时后强制启用Unaligned。这种方式能兼顾对齐的精准性和超时的兜底。1.3 持久化慢RocksDB的大状态陷阱如果确认Async Duration本身就很高那才是持久化阶段的问题。线上最常见的持久化慢是RocksDB状态后端带来的。RocksDB本质是LSM树状态一大Compaction合并就会频繁发生带来两个后果一是写放大导致写入吞吐下降二是后台Compaction线程占CPU/IO抢走正常数据处理的资源。我排查过不少CK超时伴随着CPU偏高的案例最后发现是RocksDB的Compaction和主流程在抢IO资源。RocksDB场景下几个调优思路增大state.backend.rocksdb.soft-limit和hard-limit相关配置让触发Compaction的阈值更温和把state.backend.rocksdb.writebuffer.size适度调大减少写放大如果状态能精简优先精简比如把逻辑上不需要保留的字段从state里摘掉还有一个容易被忽略的细节如果Checkpoint持久化到远端用的是HDFS那HDFS的NameNode抖动、网络跨机房带宽限制都会直接体现为Async Duration飙升。这种时候先从HDFS侧确认是不是有节点在滚动重启或者网络限速。2. 任务反复重启先区分该崩溃的异常和不该抖动的系统问题Flink任务重启是另一个高频线上告警。但很多人处理时犯一个错误不管什么异常先想着加restart-strategy的容忍次数把作业搞成怎么都打不死。这个思路很危险——有些异常就该让作业停下来而不是糊弄过去。2.1 重启策略的语义你真的搞清了吗Flink的重启策略主要有三种fixed-delay、failure-rate、exponential-delay。它们的核心区别在于什么样的失败之后允许重启以及重启多少次之后放弃。fixed-delay失败后等固定时间重启重试次数上限到了就放弃failure-rate在固定时间窗口内失败超过N次才放弃exponential-delay重启间隔依次递增1s、2s、4s……适合那种外部依赖抖动比较频繁的场景很多线上问题的根源在于重启策略配置和实际业务不匹配。比如一个作业配置的是failure-rate窗口期默认是1分钟maxFailuresPerInterval是3。外部依赖一旦抖动1分钟内失败4次作业直接进入FAILED状态然后告警。这算不算配置不当算。但如果你把maxFailuresPerInterval调整到10掩盖的可能是另一个更严重的问题。我的做法是先定位失败原因再定夺重启策略顺序不要反过来。2.2 TM容器OOM堆内还是堆外要分清楚Flink任务频繁重启绕不开的一个话题就是OOM。但OOM要细分场景不同场景的修复方式完全不同第一类是Container OutOfMemoryErrorTaskManager进程被Kubernetes或Yarn直接杀掉。日志里往往是Container killed by Yarn for exceeding memory limits但你翻JVM日志却看不到明显的堆OOM。这种通常是堆外内存或RocksDB的块内存开销超过了容器限额。我遇到过一次很典型的开RocksDB状态后端state.backend.rocksdb.memory.managed: true没开结果RocksDB的block cache自己在堆外乱吃内存吃着吃着就把容器内存吃爆了。后来把managed memory配置开掉让RocksDB的内存使用纳入Flink的托管内存统一管理问题直接消失。第二类是Java heap OOM日志里能看到java.lang.OutOfMemoryError: Java heap space。这种定位相对容易看堆栈就行。多数情况下是某些算子内部缓存了太多数据比如为了做关联在内存里积了一个大Map或者开启了Window但窗口内数据量大得吓人。第三类是Direct Memory OOM日志里是Direct buffer memory。常见于用了大量Netty缓冲区或Kafka client在高吞吐下吃满了Direct MemoryTCP的send/receive buffer也在这里。调整taskmanager.memory.task.off-heap.size或JVM参数MaxDirectMemorySize可以缓解但最好还是从源头控制并发和缓冲区的使用。2.3 用户代码的半致命异常与恢复失败的连环坑还有一种很常见的重启场景任务在重启但日志里没有明确的OOM而是用户代码抛出来的异常。比如连接第三方系统的连接池被打满、调用外部API超时、处理某条脏数据时出现了运行时异常。对于这类问题我先说一个重要排查原则如果作业从Checkpoint恢复失败了日志里可能会连续报错表面看起来是状态恢复问题实际根因却在上一个异常导致状态不一致。我踩过一次坑作业A从ck恢复时一直报Failed to rollback to checkpoint我以为是状态文件损坏花了半天去排查RocksDB的SST文件最后发现是上一个Checkpoint时用户代码里有一个非幂等的数据库写入操作Barrier之后部分记录被写入了两次恢复时校验失败才连环报错。所以排查重启问题的时候不要只看最新的日志要把从第一个异常开始的时间线拉出来。我的习惯是在Yarn/K8s上把TaskManager的日志都收集到ELK或S3然后用时间戳查第一个ERROR出现前后5分钟的上下文通常根因就在那里。另外线上业务逻辑里有一个需要特别注意的点异常捕获要不要吞掉。很多开发在算子内部包了一层try-catch把异常吞掉当作脏数据跳过。这种做法的隐患是如果异常本身是因为状态或资源导致的吞掉之后作业不会重启但状态可能已经处于半损坏状态后面的计算结果就不可信了。我的建议是数据解析类的脏数据异常可以捕获并侧输出系统资源类的异常一律向外抛交给Flink的失败恢复机制处理。3. Kafka积压先分清消费能力不足还是上游把压力传下来了Kafka积压大概是Flink运维群里被问得最多的问题。每次大促、秒杀、高峰期第一波告警基本都是Kafka LAG。一看到积压就慌一慌就开始盲目扩容这是最常见的错误操作。积压这个现象本身不是问题真正的问题是消费能力追不上生产速率或者消费能力明明够但因为某些原因被卡住了。3.1 三个指标配合看定位积压到底在哪一层先说一个最容易踩的坑只看Kafka侧的LAG不看Flink侧的消费速率。你看到LAG曲线在涨以为Flink消费慢了赶紧加并行度结果加上去LAG还在涨。这时候就要回头去看Flink背压了。我建议同时开三块面板Kafka侧的Consumer Lag看总量和单分区的lag分布确认积压的规模Flink UI的BackPressure面板看哪些算子处于HIGH状态Flink UI的Task Metrics看每个Source/算子/算子的recordsConsumedRate和busyTimeMsPerSecond关键判断逻辑是这样的如果Source task的busyTime极高接近1000ms/s说明Source在拼命拉数据但下游根本吃不下这是反压——问题在下游不在上游。如果Source task不怎么忙碌但Kafka LAG还在涨那才是Source本身消费能力不足比如并行度低于分区数、单分区吞吐上限等。简单说反压指标之前的所有算子都可能被拖慢所以先找HIGH背压最靠下游的那个算子那里才是根因。3.2 下游Sink写入瓶颈十次积压八次是它的锅根据我的经验Flink消费Kafka的作业出现积压七八成情况是Sink端写入外部系统太慢。以MySQL和ClickHouse这两种最常见的Sink端举例。MySQL场景如果目标表没有走批量写入每条记录都单独INSERT在高峰期上千QPS就能把数据库连接池打穿。解决办法是把JDBC连接器换成带batch size的写入方式比如org.apache.flink.connector.jdbc里设置batch.size 100或者使用table.exec.sink.upsert-materialize结合buffer flush。我还见过一个场景目标表主键冲突频繁导致的写失败Flink侧不断重试Sink的背压直接拉高反压传回SourceKafka开始积压。ClickHouse场景最常见的坑是每个批次过小clickhouse连接器默认的batchSize如果只有1000在百万级qps下写入效率非常低。ClickHouse适合大批次高频flush而不适合一条一条插把批大小调到5万、10万配合两三秒的flush间隔写入吞吐能提升一个数量级。处理建议先看一下Sink task的numRecordsOutPerSecond和currentSendTime如果sendTime持续高位、输出速率明显低于上游就把Sink的批量参数和并行度先调起来等LAG掉了再回头优化逻辑。注意扩容Sink并行度时要同步增加外部系统的连接池上限否则并行度加了连接池不够照样堵塞。3.3 扩容之前先确认并行度和分区数的匹配关系很多人一看到Kafka积压就直接把Flink作业并行度翻倍然后发现LAG纹丝不动。这里面有一个隐藏因素Kafka Source的并行度如果大于Topic的分区数超出的并行度是空转的。Flink的FlinkKafkaConsumer每个并行子任务会负责消费一个或多个分区并行度超过分区数时多出来的task根本没有数据可分。所以正确操作是先看Kafka Topic有几个分区比如8个分区把Flink Source并行度设置为8或略小于8确认其他算子没有成为新的瓶颈如果分区数本身不够需要在Kafka侧扩容分区才能让Flink横向扩展生效这里还有个细节即便你有32个分区Source并行度是32但下游某个算子并行度只有4那反压依然会在那个算子上堵住Kafka照样积压。所以要顺着拓扑把每个算子的并行度都检查一遍找到最窄的那个瓶颈。补充一个实操经验临时处理积压时我不建议随便改作业的并行度然后从旧Checkpoint恢复。并行度变更会导致状态重分布恢复耗时比正常恢复长很多搞不好恢复期间积压更严重。我更常用的一种临时方案是在不改变并行度的情况下先把Sink的分批参数调大、把有问题的脏数据任务临时跳过用最小的变更先让LAG降下来。4. 数据倾斜从Web UI的指标矩阵一眼看出病灶数据倾斜是Flink流式计算里一个很微妙的问题——它不直接报错也不直接让作业挂掉它只是安静地让某些subtask跑得慢然后引发一连串连锁反应这个subtask反压、Checkpoint等待它对齐、整个作业背压、Kafka开始积压。所以你会看到一台机器CPU跑满其他机器闲得摸鱼整个作业的吞吐却被一台机器卡住。4.1 三个被很多人忽略的UI细节在Flink Web UI里点开一个任务的Metrics可以看到每个subtask的指标。我判断是否存在数据倾斜通常会同时看三个数值recordsReceived / recordsSent每个subtask接收和发送的记录数taskmanager的线程CPU利用率通过JMX或Grafana这个subtask的checkpoint size如果某个subtask的状态比其他subtask大很多说明它承接的数据量或key范围明显偏离均值如果发现同一个Task的各个subtask之间recordsReceived相差5倍以上或者某个subtask的CPU利用率长期是其他节点的两三倍基本可以下结论倾斜了。4.2 四种典型倾斜场景和对应解法场景一keyBy热点Key倾斜这是最常见的倾斜。比如用户日志里某个机器ID、某个店铺ID占了全量数据的40%按这个key去keyBy那40%的数据全压在一个subtask上。解法是加盐业界叫两阶段聚合第一层keyBy用(原始key 随机后缀)做key把热点key先打散到多个subtask上去做局部聚合第二层keyBy按原始key聚合合并局部结果举例用户行为按shop_id聚合统计PV某个爆款店铺的shop_id是9527。第一层keyBy时把9527加上1~10的随机后缀变成9527_1、9527_2……这样原来一个subtask的压力被分摊到10个第一层算完之后再按去掉后缀的原始key做第二次聚合得到最终PV。加盐的坑在于这个方案对增量型聚合sum/count是精准的对去重型聚合distinct不精准。因为你把一个key的数据拆到了多个临时桶不同桶里的同一个元素可能被重复统计。如果业务对精确去重有硬性要求需要把去重逻辑放在第二层聚合再做代价是第二层压力大或者用HyperLogLog之类的近似去重方案。场景二双流Join时的热点Key两张表Join明细数据量大维表数据量小但维表里某一个key比如某个爆款商品、某个大V用户被命中的次数特别多。两流Join时那个热点key导致同一个subtask的重负载。处理思路是热点拆分广播先用统计数据比如从维表里分析查出来哪些key是热点把维表分流热点key的维表数据广播到所有subtask冷key的维表数据走正常keyBy明细流也对应分流命中了热点key的数据直接在本地跟广播的维表做join其余数据走常规join最后union结果这个方案在实现上比加盐复杂一些但思路很清晰。如果你没有热点key的预先统计数据可以先从监控指标里找哪个key的流量异常大再把它配置进一个外部规则表里Flink定时加载这个规则。场景三窗口聚合倾斜比如开一个10分钟的滚动窗口本来按key分布还行但窗口在某个时间点整体数据量很大所有计算都堆到窗口触发的那一瞬间。这个其实不完全是key分布问题而是窗口触发时间点的计算密集冲突。解法是在窗口算子前做一层本地预聚合。比如先按(key 窗口起始时间)进行keyBy把数据先算一个小聚合到窗口真正触发时把预聚合结果再合并一次。实际上Flink的增量聚合函数本身就有预聚合效果——如果你用的是reduce或aggregate数据在进入窗口时就逐条计算了窗口触发时只是输出最终结果不需要全量遍历。场景四Sink端倾斜数据分布均匀但写入外部系统时某个分区或某个连接特别慢。比如写ClickHouse时目标表按天分区某个业务分区的数据量特别大写入hot partition的请求在ClickHouse侧排队。这种属于外部存储的分区热点需要在Sink侧做分摊可以按shuffleKey把数据均匀打散到多个Sink子任务然后由外部系统合并或者调整目标表的存储策略把热点分区预拆分。4.3 Rebalance到底要不要用Rebalance转发到下一个可用subtask在倾斜场景中是很多人优先想到的手段但我要提醒一句Rebalance解决的是并行度分配不均的问题解决不了key本身计算量不均衡的问题。如果计算本身是按key绑定的比如keyBy之后的聚合你前面Rebalance把数据均匀发送到了各subtask但到了该keyBy聚合的地方还是按key哈希该倾斜还是倾斜。我建议的用法是在Sink端不需要按key聚合、且外部系统写入能力有差异的场景用Rebalance把写入压力均匀分摊到每个Sink并行度上这个是有意义的。但在算子内部做聚合的场景不要指望Rebalance救命老老实实加盐或做预聚合。5. 四类故障的联动一次真实复盘的排查顺序上面分开讲了四类问题但线上真实场景里它们经常是串在一起出现的。我挑一个自己实际处理过的案例还原完整的排查链路你会发现很多事情是连锁的。5.1 时间线还原一次大促流量下的多米诺骨牌某日19:00大促流量开始上涨值班同学发现Kafka LAG告警LAG从几千涨到几十万。19:15任务开始出现Checkpoint超时。19:30任务进入反复重启状态一小时内重启了5次。19:45作业彻底FAILEDKafka LAG飙到几百万。从表面看这是一次Kafka积压→CK超时→任务重启的组合故障。但按上面说的排查方法我们先看背压打开BackPressure面板发现HIGH背压出现在Sink到ClickHouse的算子上。再点进去看Sink算子的busyTimeMsPerSecond显示990以上几乎满负载而Source端的recordsConsumedRate并不高——说明Source消费得动是下游Sink写入太慢把整条链路堵住了。而ClickHouse在这个时段因为目标表在做一个合并操作写入性能下降了30%左右。链路是这样的ClickHouse写入变慢 → Sink背压Sink背压 → 反压传导到Source之前的每个算子反压导致Kafka消费速率下降 → LAG积压反压导致Barrier对齐时间拉长 → Checkpoint持续超时Checkpoint超时失败达到restart-count上限 → 任务重启重启恢复期间不消费数据 → 积压进一步恶化表面上的CK超时其实只是下游写入性能抖动在水面下推出来的结果。5.2 复盘后的统一排查顺序经过这次复盘我梳理出了线上故障处理的优先级现在处理这类问题已经不动摇了第一步看背压找到HIGH背压最靠下游的那个算子它就是整个链路变慢的起点。接近90%的问题到这个步骤就能锁定根因方向。第二步看CK如果第一步定位出来是反压导致的CK超时那你需要做的是解决反压而不是调大checkpoint.timeout。调大超时只会推迟失败掩耳盗铃。第三步看日志如果任务已经在重启把从第一次异常开始的时间线拉出来看重点排查第一个异常后面的往往只是连锁反应。第四步看资源确认TaskManager的内存、CPU、GC是否健康。有时候慢不是因为代码慢而是机器上的资源被抢了。5.3 线上的预防性配置清单经历过被凌晨告警支配的恐惧后我给自己负责的Flink作业定了几个底线配置写在这里供参考Checkpoint状态大且延迟敏感的作业开启Unaligned Checkpointaligned-checkpoint-timeout30s普通作业保持对齐模式但把checkpoint.timeout设置为min(10min, 期望CK时长*2)Restart策略全部用failure-rate窗口1分钟最大失败3次。既能容忍外部抖动又不会让作业无限重启掩盖致命错误Kafka lag告警不只盯总量还要盯单分区最大lag防止单分区倾斜导致的局部积压关键算子指标Web UI里重点盯busyTimeMsPerSecond和mailbox throughput这两个指标能第一时间反映反压源头回到开头的那个场景。凌晨两点被叫醒看着红彤彤的告警面板现在我的第一反应不是惊慌而是按背压→CK→日志→资源的顺序去定位。这几个故障的排查链路说到底讲究的是先分清楚因果再动手处理。很多人栽跟头不是因为技术不够而是把因果关系搞反了——把Checkpoint超时当成根因去调参数把Kafka积压当成消费能力不足去扩容结果一通操作猛如虎问题反而更严重。最后再分享一个小技巧我在排查这类问题时习惯先把关键指标的截图按时间线保存下来尤其是背压面板和Metrics历史的截图。故障处理完复盘时这些时间线截图能帮你非常准确地还原哪个指标先变化、哪个指标后变化而因果关系的判断恰恰就藏在这些时间先后里。毕竟Flink线上问题排查方法论对了剩下的就是耐心和细心了。
返回列表