news 2026/10/8 20:20:49

Flink数据倾斜实战:从定位到治理,两阶段聚合与Sink背压排查

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Flink数据倾斜实战:从定位到治理,两阶段聚合与Sink背压排查

1. 从一次任务卡死说起:数据倾斜到底是什么

先说个我自己的真实经历。有次线上跑一个实时指标计算任务,数据量一天也就几亿条,并行度开到32,结果每天到了晚高峰,整条链路就开始疯狂反压,Kafka消费Lag飙到几百万,任务重启了好几次都没用。打开Flink Web UI一看,某个Subtask的忙率直接拉到100%,而旁边的几个Subtask忙率只有10%左右,CPU根本没跑满。当时我就知道,这又是数据倾斜在作妖。

数据倾斜这个坑,做大数据的人基本都踩过。简单说就是:明明给了你32个并行度,但数据就是不均匀地挤到了某个Task上,导致那个Task成为整个链路的短板。尤其是Flink这种实时计算引擎,处理的是无界流,一旦某个并行子任务因为数据量大而处理不过来,背压就会顺着算子一路往上传递,最终把整条链路堵死。

那篇文章要解决的核心问题就两个:怎么定位数据倾斜,以及怎么治理数据倾斜。这篇文章我会从定位手段、常见场景、不同阶段的改造方案、Sink端写入保障这几个维度,把我在实战里用过、验证过、踩过坑之后沉淀下来的方法论完整梳理出来。不管你是在用Flink同步MySQL数据到ClickHouse,还是用SpringBoot整合Flink做实时任务开发,只要你的任务里存在KeyBy、窗口聚合、双流Join这些算子,这篇内容都值得完整看一遍。

先说清楚一件事:数据倾斜没有一个“万能解药”。它是分场景的,同样的倾斜现象,可能来自Key分布、Join字段选择、窗口策略,也可能是下游写入导致的背压假象。所以这篇内容不适合当字典查,更适合按顺序读一遍,建立一个完整的排查思路。

2. 三招定位数据倾斜,别再靠猜

2.1 从Web UI的忙率和积压数据量判断

很多同学一上来就喜欢看日志、翻异常,其实数据倾斜的特征非常明显,根本不用猜。打开Flink Web UI的Task Manager页面,重点看两个指标:忙率(busy%)和积压记录数(backPressure)。

正常情况下,各并行子任务的忙率应该是大致均匀的,差异通常在10%以内。如果有某个子任务的忙率长期在90%以上,而其他子任务只有20%、30%,并且该算子的输出缓冲区积压数据不断上涨,那基本可以断定这个子任务就是热点Task。

注意:忙率低不一定是没有倾斜,也可能是子任务在等待下游(比如等待Join的另一条流),这时候要结合recordIn/recordOut的速率来综合判断。最稳妥的办法是把Web UI的指标面板切到Chart模式,观察10分钟以上的趋势图,别只看瞬时值。

2.2 从日志和反压监控里寻找热点算子

如果任务已经卡死或者Lag报警了,说明问题已经很严重,这时候不要再等Web UI慢慢加载了,直接看反压监控。Flink的BackPressure监控从1.5版本之后就内置了,但它的采样周期默认是100次采样,每次间隔50ms,对于瞬时倾斜可能不够敏感。我更推荐的做法是:

  • 检查Task Manager日志里是否有频繁的“Buffer pool exhausted”记录
  • 检查Kafka Topic的消费者Lag,如果某个Partition的Lag明显高于其他Partition,说明问题已经回溯到源头
  • 用curl命令直接拉取JobManager的REST API,实时看每个Subtask的busyTimePerSecond指标

我个人的判断标准是:某个Subtask的忙率持续3分钟以上超过其他Subtask两倍以上,并且该算子的numRecordsInPerSecond远高于均值,那基本就可以确定为数据倾斜,不需要再做其他复杂的分析了。

2.3 解析上游数据分布,确认根因

定位到热点算子之后,还需要搞清楚“倾倒是怎么形成的”。这一步很多人会跳过,但恰恰是关键所在。数据倾斜的根因说白了只有三类:

  • Key分布不均匀:比如按用户ID分桶,1%的大用户贡献了90%的数据量
  • 数据膨胀:Join时一对多关联,某条大Key对应的关联数据特别多
  • 计算本身存在集中性:比如窗口结束时的集中触发,或者全局聚合点

确认根因的方法也不复杂:在热点算子上游加一个临时Sink,把数据按Key做一次count统计,10分钟就能看到Top N的Key分布情况。如果Top 10 Key贡献的数据量超过总量的50%,那这就是典型的Key倾斜。

3. 定位之后怎么办:分场景的改造方案

3.1 KeyBy重分区:核心思想是打散和二次聚合

定位到Key倾斜,最常见的解决手段就是给Key“加盐”。我没法给你一个“使用什么加盐策略最好”这种一刀切的答案,因为加盐方案跟业务语义是强绑定的,关键在于你的下游聚合能不能接受“先局部聚合、再全量聚合”的结果。

我的常规做法是分两步走:

第一步,局部聚合。给Key加一个随机后缀,比如userId变成userId + "_" + random.nextInt(10),把大Key的数据打散到多个子任务上做第一轮聚合。这一步能有效降低单点的处理压力。

第二步,全量聚合。去掉随机后缀,只保留原始Key,开一个窗口或者直接做KeyBy聚合,把局部聚合的结果再做一次合并。

这里要说清楚一个关键问题:加盐能解决的是“大Key导致单点压力”的场景,但有一个前提——业务上必须接受一定程度的延迟。如果你加了10个盐,第一轮聚合意味着有10个并行子任务在同时处理同一个大Key的数据,这会引入额外的合并开销和网络传输,如果任务本身的时延要求非常苛刻,那么加盐方案就不适用。

3.2 Join场景倾斜:处理思路完全不同

很多人在做双流Join时遇到倾斜,第一反应也是给Key加盐。但我必须说一句:在Join场景里加盐,往往会带来严重的正确性问题。因为Join的语义要求同一Key的数据必须同时到达同一个子任务才算精准匹配,加盐之后直接相当于把左流和右流拆到了两个不同的子任务上,结果就是大量关联不上的数据被丢弃。

这种情况下,更实用的方案是:

  • 维表场景:优先用广播代替普通Join,把维表做成Async I/O的异步查询,或者用广播状态加载到每个子任务本地,避免数据Shuffle
  • 大表和小表Join:如果关联的数据量本身差异巨大,可以考虑把大表按业务维度裁剪再用广播方式处理
  • 数据膨胀场景:分析关联字段,如果一条A表记录对应了B表的几十万条记录,这时候无论怎么并行都没用,必须从源头上压缩数据量,比如提前在内层做过滤group,或在扩容前先做pre-aggregation

要说清楚的是:不是所有Join场景都能靠Flink层面解决。有些关联膨胀是业务模型本身的问题,你就得回到上游SQL、回到数据同步链路去改。这种问题属于“即使并行度拉到512也没救”的情况,等到Web UI已经卡得没法看再排查,等于是在给死亡任务做临终关怀。

3.3 大Key场景的特殊处理:拆分算子和局部缓存

还有一种倾斜不太容易被定位到,因为它不是发生在KeyBy阶段,而是发生在状态读写阶段——大Key对应的状态太大,导致RocksDB读写都卡在同一个Task上。这个场景在Flink里极其隐蔽,因为Web UI上看忙率分布是均匀的,但整体吞吐就是上不去。

我实战中用过的有效方案是:把大Key对应的数据源单独拆出来,用一条独立的Flink任务链路去处理,和普通Key的任务链路分开跑。比如一个大商户贡献了90%的交易量,那就给这个商户单独映射一个分桶,把它路由到独立的子任务上,配合单独的并行度资源配置。

另一个思路是用旁路缓存。把大Key的内容存到Redis里,窗口计算时不走状态读写,直接从Redis里读取历史数据做合并计算。这个方案适用于场景本身允许数据有秒级延迟的情况,好处是绕开了RocksDB状态读写的单点瓶颈。

4. 直接从源头优化:两阶段聚合改造实战

4.1 第一阶段:窗口前先做一次分区聚合

上一节讲的都是“倾斜已经发生之后怎么治”,但更高级的做法是提前做好架构设计,在一开始就让倾斜发生的概率降到最低。

我做过一个典型的实时大屏任务,场景是统计各省份每分钟的订单金额。最初的实现是直接按省份ID做KeyBy,结果某个电商大省在晚间高峰期的数据量能占到全平台40%以上,那个省份对应的子任务每个月都要宕几次。

改造方案就是两阶段聚合:第一步先做分区内预聚合,即按“省份 + 随机后缀”的方式打散,窗口内部先做一轮计算,把结果量从千万级压缩到百级别;第二步再按真实省份ID做汇总。

代码实现并不算复杂,核心逻辑大致是这样:

// 第一阶段:加盐打散,局部聚合 DataStream<OrderRecord> keyedStream = source .keyBy(order -> order.getProvinceId() + "_" + ThreadLocalRandom.current().nextInt(10)) .window(TumblingProcessingTimeWindows.of(Time.seconds(10))) .aggregate(new ProvinceAggFunction()); // 第二阶段:按真实省份聚合 DataStream<ProvinceMetric> resultStream = keyedStream .keyBy(metric -> metric.getProvinceId()) .window(TumblingProcessingTimeWindows.of(Time.seconds(10))) .aggregate(new MergeProvinceAggFunction());

两阶段聚合的效果立竿见影:同一条链路,改造前热点Task的忙率是95%,改造后Top Task忙率降到50%左右,整条链路的吞吐量直接翻了一倍。不过要提醒的是,两个窗口会带来额外的延迟,如果业务要求的延迟是秒级以内的,这个方案就要再斟酌了。

4.2 添加盐的粒度怎么定:参考QPS和状态规模

加盐粒度的选择,是最容易被人忽略的细节。盐加少了,热点Task可能还是吃不下;盐加多了,第二阶段合并的开销又上来了。我个人的经验公式是:盐的个数 = 热点Key的数据量 / 单个子任务可承载的QPS。

假设某个大省高峰期的数据量是每秒12万条,单个并行子任务稳定处理大概是每秒2万条,那么盐的个数取6到8个比较合适。注意这里是“并行度 × 盐个数”对应的整体处理能力要能覆盖峰值,并且最好留30%的余量,因为Flink的窗口触发、序列化开销都会挤占一部分处理能力。

还有一个容易被忽略的点:加盐切分的逻辑最好放在最上游的Source端。如果把加盐逻辑放在KeyBy之前,等于先做了一次全量Shuffle,再打散,这样虽然热点子任务被分担了,但Shuffle本身可能成为新的瓶颈。

4.3 从源头pre-aggregation,彻底绕开倾斜算子

最近Flink社区里很火的一个话题,是把倾斜治理的思路往前移到“源端”而不是“算子端”。比如你用Flink做MySQL同步到ClickHouse,其实就可以利用Flink CDC里的shuffle优化策略,直接在Source阶段根据表结构提前做一次数据合并,然后再往下游分发。

SpringBoot整合Flink做实时任务时,我也踩过一个类似的坑:把Kafka的数据直接接入Flink,不加任何预处理,结果一个按订单ID分桶的KeyBy算子每天固定倾斜。后来我在接入层加了一个轻量的内存聚合结构,把1秒内的相同订单ID数据先合并成一条再进入主链路,倾斜问题不治而愈。

这个思路其实值得我们反思:很多倾斜不是Key分布造成的,而是上游数据本身存在冗余。如果每条数据都足够“瘦”,就算某个Key的数据量大,处理压力也会小很多。所以在设计数据入链路的阶段就做好压缩和预聚合,比在链路内部想各种办法去打散更高效。

5. Sink端写入的“背压假象”:JDBC连接器异常深度排查

5.1 Sink背压和数据倾斜的区别:先看队列堆积方向

在大部分Flink任务里,数据倾斜的锅都让上游的KeyBy算子背了,但实际上有相当一部分情况,真正的瓶颈发生在Sink端。尤其是做MySQL同步到ClickHouse、用JDBC连接器做批量写入的任务,Sink的写入能力跟不上,就会在上游算子形成反压,表现看起来跟数据倾斜一模一样。

我排查过的案例里,最典型的是这样的情况:Web UI显示的瓶颈在Map算子,忙率很高,但Map算子的逻辑只是做一行格式转换,根本不可能成为瓶颈。后来一查Sink端,发现是JDBC连接器批量写入超时疯狂报错,每一条失败的数据都会重试,重试期间阻塞了后面的数据,反压一路传到了上游。

判断究竟是倾斜还是Sink背压,一个很简单的方法:看队列堆积方向。倾斜导致的堆积是“某些Subtask输入堆积”,而Sink背压导致的堆积是“所有上游Subtask输出堆积”,后者的分布是均匀的。如果你看到全体Subtask都在堆积,且瓶颈算子的聚合指标没有明显的“高个子”,那一定是下游反压,别再花时间怀疑倾斜了。

5.2 JDBC连接器写不动的三个隐藏原因

说到JDBC连接器,很多做了很久的工程师都会栽在这里。我梳理一下三个最隐蔽的原因:

第一个是批量大小配置不当。Flink JDBC Sink默认的批量写入大小是batchSize=1000、batchIntervalMs=1000,如果你的每次写入涉及的数据比较大,比如一条记录就有几KB的JSON字段,那一次写1000条就会导致网络包过大、数据库锁竞争加剧,频繁超时。这种场景调小batchSize反而更高效,我一般会先压到100试试效果,再逐步往上调。

第二个是事务隔离级别冲突。Flink的JDBC Sink在开启setAutoCommit(false)之后,如果目标ClickHouse或MySQL本身的隔离级别和连接池配置不匹配,会导致写入失败。这类报错有时候是间歇性的,网络热词里的“flink的jdbc连接器异常”很大程度上就是这类问题。排查手段是打印连接获取的Detail日志,看哪个环节耗时最长。

第三个是并发写同一个表时的行锁竞争。如果你的Sink并行度和目标表的写入热点冲突(比如多个Subtask同时写同一个范围的主键),WAL日志会疯狂等待,体现出来就是背压。解决思路是把写入层的并行度适当调低,或者落库之前先做一次按目标表主键的排序。

5.3 解决Sink反压的落地配置:从超时到批次

下面给出一套我调优过多次的JDBC Sink参数模板,可直接参考:

参数项推荐值说明
sink.buffer-flush.max-rows100-500按单条体积灵活调整,大对象字段取小值
sink.buffer-flush.interval1s-2s放宽一点可以帮助批量提交合并
sink.max-retries5重试次数太多反而加剧背压
jdbc.fetch-size500避免单次拉取过大内存溢出
connector.write.flush.interval1s适用ClickHouse连接器
sink.parallelism实测后调低目标是让每个分片单次写入的数据量稳定

重要提示:调低Sink并行度听起来违背直觉,但很多时候确实有效。因为目标数据库能承受的总事务/秒是有限的,并行度太高会导致大量请求同时打过去,数据库反而疲于处理冲突和锁等待。一个比较稳的做法是先测目标库当前能承受的写入TPS,再反推Flink Sink应该开多少并行度。

5.4 MySQL同步到ClickHouse场景的写入倾斜

把MySQL的数据实时同步到ClickHouse,也是Flink JDBC连接器使用最热的场景之一。我做过一个电商订单同步任务,源库的单表数据按订单时间增长,如果直接开KeyBy订单ID再写入ClickHouse的Distributed表,置换逻辑会经常出问题。

这个场景的倾斜往往不在Key算子,而在ClickHouse的MergeTree表引擎层。ClickHouse在建表时如果分了3天、7天的分区,数据写入会发生写倾斜:最新分区的写入量很大,而历史分区几乎不写入。这时候Flink Sink再开高并行度,只是想当然的配置,实际反而会加剧写入冲突。

实操建议是:ClickHouse表的分区键尽量选择成本更低的维度,比如日期或者哈希字段。如果实在要按业务主键做分区,那么Sink端要开启“固定字段路由”,尽量让同一分区的数据都落到同一个或少数几个子任务上连续写入,避免频繁切换分区导致的锁等待和文件碎片。

6. 我被问最多的问题,列个排查速查表

6.1 为什么明明加了盐,倾斜还是没缓解?

先说结论:大概率是盐加在了不该加的位置。很多人的实现是:

keyBy(order -> order.getUserId() + new Random().nextInt(10))

把加盐逻辑写在了KeyBy算子内部,这样带来的效果是:同一个用户的同一条数据,每次进来随机盐都不同,下游的窗口聚合就会把数据打散到10个临时窗口里。表面看倾斜没了,但第二阶段合并时,由于每个临时窗口里都各自维护了一份状态,本身就把计算量放大了好几倍。

正确的做法是给每个窗口内的Key固定盐,而不是每条记录随机盐。典型写法是带时间或窗口ID的确定性加盐:

keyBy(order -> order.getUserId() + "_" + order.getWindowStartTime() % 10)

这样才能做到“同一批窗口数据固定在一个子任务上做局部聚合”,而不同窗口之间打散到不同子任务。

6.2 窗口聚合倾斜和普通KeyBy倾斜有什么不同?

窗口聚合倾斜,指的是在窗口计算结束时,所有属于同一个窗口的数据瞬间集中到同一个算子做计算。这种倾斜跟Key分布关系不大,而跟窗口的设计关系大。常见于事件时间窗口结束后有大量的Trigger和定时器在同一个瞬间触发。

排查此类问题,不要盯着Key看,要先观察是否有明显的时间周期规律,比如每整点、每整分钟固定卡一下。如果确实存在这种规律,可以考虑把窗口切割成更小的粒度(从1小时切成15分钟),或者把每个窗口的任务Schedule错开,避免同一时间点所有窗口同时触发大规模计算。

6.3 状态很大的Key倾斜,如何判断是RocksDB问题?

RocksDB的状态后端本身是没有并行度概念的,它归属于某个具体的Keyed State。所以当一个Key长期占用巨大的状态量(比如几GB甚至几十GB),RocksDB的找表、Compaction都会卡在这个Key上。

判断方法很简单:检查Task Manager指标中的rocksdb.cur-size-all-mem-tables。如果一个Task的值远超其他Task,且数字持续升高不下降,那基本就是大Key状态膨胀导致的单点问题。这个问题靠加盐是救不了命的,方法是把状态清理逻辑加好,或者调整状态TTL,把不常用的历史状态尽早淘汰。

7. 聊几点工程落地的经验心得

7.1 并行度不是越多越好:先定Sink容量再往上推

很多团队在做Flink性能调优时,第一个动作就是把并行度往大了调,仿佛并行度是万能的。但我的经验是:并行度是结果,不是原因。你要先根据最下游的承载能力(比如目标ClickHouse、MySQL、Kafka分片数),推断出Sink的合理并行度,再逆向推导上游算子的并行度。如果下游根本接不住,上游开再多的并行度也只是把压力堆积在反压链路上,没有实际意义。

举个很简单的例子:Kafka某个Topic一共12个分区,那Source端的并行度开到12就够了。你开24个并行度,多出来的12个分区根本拿不到数据,反而多了一堆空闲的Task。同理,下游数据库承受能力是5000 TPS,Flink全链路能处理5万 TPS,那瓶颈就是下游,你花再多精力调上游也没用。先把短板找到,才能谈优化。

7.2 用好Source端的慢启动消费策略

Flink在对接Kafka做实时聚合时,有可能会遇到一种奇怪的现象:任务启动的前10分钟运行正常,之后突然积压Lag。这种情况往往是Kafka的分区数据不均衡,某个分区的历史数据量特别大,而Flink的消费逻辑是全速消费的,就导致消费慢的分区成了瓶颈。

解决思路是开启Source端的有限制消费,比如设置startup-mode配合一个启动期的限速策略,让任务在启动阶段先以较低速度消费,等State恢复稳定后再逐渐放开。这个槽位配置看起来不起眼,但在SpringBoot整合Flink做在线任务的场景里,真的能避免很多无谓的“任务重启-重复消费-再倾斜”的恶性循环。

7.3 混沌测试:给任务注入故障看看它的抗性

最后讲一个很多人忽略的环节:上线前的“故障注入测试”。我在处理过很多次线上倾斜问题之后,形成的一个固定习惯是:每次新任务上线前,都会特意调一个上游Topic的某几个分区的数据量,人为制造出5倍以上的流量差异,然后观察任务的表现。

如果任务在人为注入的倾斜下还能保持20%以内的忙率波动,说明架构的拆解是健康的。如果忙率直接飙到90%,那就说明当前的设计还扛不住极端情况。这个做法的价值在于,把排查从“事后救火”变成了“事前预防”,成本比线上加急处理低太多了。

8. 最后说几句实在话

做了这么多年的实时计算,处理过数不清的倾斜问题,我的一个体会是:数据倾斜是统计规律,不是Bug。只要有分布式、有数据分布,就一定会有倾斜。与其花大量时间找一个“彻底消灭倾斜”的银弹,不如建立一套完整的定位手段和治理预案,把每次倾斜变成一次可复用的经验。

这个内容如果你完整读下来了,大概率对下面几件事有了一个系统的认知:怎么从Web UI和监控指标中精准定位热点算子,什么样的倾斜场景适合加盐、什么样的场景不能加盐,Sink端反压和上游倾斜怎么区分,以及最关键的——在设计任务之初就通过两阶段聚合、分桶策略和容量规划来规避倾斜。我希望你把这些方法拿回自己的业务里去试,试过之后你可能会发现,那些让你熬夜排查的“疑难杂症”,其实在架构层面很早就写下了答案。

最后再分享一个小技巧:每次碰到倾斜问题,先不要急着改代码,问自己三个问题——这个倾斜是不是真的发生在计算层?上游能不能做预聚合?下游Sink能不能接住全量计算的结果?把这三个问题想清楚,大部分倾斜问题其实已经解决了一半。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/8 20:20:31

TiDB国产化升级实践:从分布式架构到行业落地的选型指南

作为一个长期在数据库选型和架构改造一线折腾的人&#xff0c;最近圈子里讨论度最高的话题&#xff0c;除了国产化替代&#xff0c;就是分布式数据库到底怎么选。恰好下周要去长沙参加3月14日的TiDB社群“湘聚”活动&#xff0c;主题聚焦零售、医疗、金融、交通、智能制造这些重…

作者头像 李华
网站建设 2026/10/8 20:20:30

MFAC无模型自适应控制仿真全解析:伪偏导数估计与CFDL/PFDL/MIMO实践

最近整理了一套很实用的仿真资料&#xff0c;主题正好是“六个MFAC无模型自适应控制仿真伪偏导数估计动态线性CFDLPFDLMIMO”&#xff0c;里面除了程序&#xff0c;还配了一部分参考资料。我陆陆续续用这套东西给不同项目做数据驱动控制验证&#xff0c;踩了不少坑&#xff0c;…

作者头像 李华
网站建设 2026/10/8 20:19:33

IB Specification 2.1 实战:从报文头到QP状态机的RDMA排障指南

简介&#xff1a;IB Specification 2.1 是 IBTA 发布的 InfiniBand 架构官方规格书&#xff0c;对应 Volume 1 通用规范&#xff0c;面向 RDMA 网络研发工程师、数据中心架构师及 HPC 技术人员&#xff0c;帮助理解高速互连标准的设计与演进。这份 PDF&#xff08;共 1 个文件&…

作者头像 李华
网站建设 2026/10/8 20:18:15

仿百度网盘JavaWeb小型云盘系统:从部署到实现核心功能

简介&#xff1a;一套基于Java Web实现的轻量级云盘系统&#xff0c;面向正在学习Java后端与Web开发的初学者、以及需要快速搭建在线存储演示项目的开发者。项目模仿百度网盘的核心交互&#xff0c;涵盖文件上传、下载、分享、删除、重命名等常用操作&#xff0c;并包含用户认证…

作者头像 李华
网站建设 2026/10/8 20:16:05

MQTT工业物联网实战:从Broker搭建到设备接入与云平台对接

做工业项目的朋友应该都有同感&#xff1a;现场设备一旦要上云&#xff0c;通信协议是第一道绕不过去的坎。这几年我经手的项目里&#xff0c;MQTT几乎是出现频率最高的一个词——从485仪表、PLC采集&#xff0c;到组态软件&#xff0c;再到TLink这类物联网云平台&#xff0c;中…

作者头像 李华