恒美微站 Logo 恒美微站
  • 首页
  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心
  • 联系我们

Flink数据倾斜从识别到化解:两阶段聚合、加盐与实战排查全复盘

  • 首页
  • 资讯中心
  • /
  • Flink数据倾斜从识别到化解:两阶段聚合、加盐与实战排查全复盘

相关资讯

MFAC无模型自适应控制:动态线性化三种方法的非线性系统应用 2026/10/8 9:26:35
数字孪生与工业设计衔接:工具链选型与落地路径解析 2026/10/8 9:26:35
旺店通WMS对接实战指南:从订单推送到库存回传的完整链路 2026/10/8 9:26:35

最新资讯

Android文件存在却打不开?Ext4底层状态诊断指南
WorkBuddy装完先补哪5个技能?find-skills、humanizer-zh等实战指南
DeepSeek Harness桌面版+Obsidian:本地知识库搭建与RAG检索实战
预算有限如何找到便宜又靠谱的AI工具渠道?低成本调用方案与避坑指南
AI编码技能框架实战:从提示词到结构化协作的完整指南
对话系统Skill命中率下降的四层归因与治理方法

今日推荐

context-mode实战指南:从全量塞入到结构化裁剪与检索增强
大模型对话上下文管理实战:三种模式与Token优化
抖音用户主页视频数据爬虫详解:点赞、收藏、分享字段抓取与 TaoToken 统一 Key 配置

本周热门

MR25H40CDF + PIC18F65K40:工业记录仪高可靠存储实战
基于STM32的数控恒压恒流电源设计:从硬件到PID调参全解析
LT9211 MIPI重定时器原理与双路扇出实战指南

本月精选

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证
2026 大模型集体涨价:用 Python 做企业 Token 成本测算与选型避坑(附配置)

Flink数据倾斜从识别到化解:两阶段聚合、加盐与实战排查全复盘

发布时间:2026/10/8 9:26:35
Flink数据倾斜从识别到化解:两阶段聚合、加盐与实战排查全复盘 数据倾斜这个话题做实时的同行应该都不陌生。症状特别典型明明并行度开到了16跑到Web UI上一看只有一个Subtask负载接近100%数据量占了全流程的大头其他节点闲得跟没事人一样整个作业的吞吐被拖到脚踝延迟直线上升。我最早遇到这个问题是在一个用户行为实时统计的场景里某个头部渠道的用户量占了全量数据的七成KeyBy之后一个节点要处理千万级数据其他节点只有几十万。那时候我也没太多经验看到反压飙红第一反应就是加并行度结果加了也没用——热点数据就那么多加并行度只是复制热点并没有分散掉。后来才慢慢摸清楚Flink数据倾斜的本质是分区策略和数据分布不匹配不是资源不够也不是并行度不够是数据和分区之间的映射关系出了问题。这篇博文把我踩过的坑和验证过的方法整理一遍从怎么识别倾斜、核心化解方案到MySQL同步到ClickHouse的实战排查再加上SpringBoot整合Flink时容易踩的雷做一次系统的复盘。不论你是刚接触Flink的新手还是已经在大数据处理里折腾了一段时间的开发者这篇文章应该都能提供一些能直接上手的东西。1. 数据倾斜到底长什么样现象识别与根因分析1.1 判断作业是不是倾斜先看这四个信号很多同学一上来就怀疑数据倾斜但其实未必。我自己判断是否倾斜基本只看四个信号综合起来才能下结论。第一个信号是Web UI上Subtask的数据量差距悬殊。打开Flink的JobManager页面看每个Subtask的Records Received或者Bytes Received如果某个Subtask的数据量是其他Subtask的5倍以上基本可以确定这个算子的输入数据Key分布有问题。注意这里说的是“输入”数据量因为倾斜通常发生在KeyBy之后的下游算子不是KeyBy本身。第二个信号是反压在整个算子链上传递但上游吞吐并没有增长。反压说明算子处理不过来了但如果倾斜发生反压会集中在某个Subtask对应的那条链路上其他链路很通畅。这时候去点开深入的BackPressure监控热点TaskManager的某个线程占用率会明显高于同进程的其他线程。第三个信号是Checkpoint时间越来越长甚至超时失败。倾斜节点承担了大量数据状态也会比其他节点大做快照的时候序列化耗时显著增加。如果你发现CK的Duration从几秒钟飙到几十秒而且总是固定挂在某一个节点上那基本就是倾斜了。第四个信号是延迟监控曲线出现周期性尖峰或者说持续攀升不回落的趋势。数据倾斜导致某个分区处理不过来数据在队列里排队端到端延迟自然会增大。这个信号最容易被业务方感知也是最影响SLA的。这四个信号都出现了才基本坐实倾斜。只有反压没有数据不均可能是Sink慢只有CK超时没有数据不均可能是状态太大。定位问题一定要用数据说话不能凭感觉。1.2 倾斜是怎么产生的分区策略和Key分布的错配要理解倾斜的根因得先搞清楚Flink的分区机制。以最常用的KeyBy为例Flink会对Key做哈希然后对并行度取模计算出这条数据应该进入哪个分区。源码里用的是murmurHash对这个哈希值再做一个正态化处理尽量让数据均匀分布到各个分区。这个机制对均匀分布的Key是没问题的。比如一个订单流订单ID天然唯一哈希后分布就很均匀。但真实业务里的Key哪有那么均匀。用户行为按用户ID分组头部用户贡献大量行为订单统计按商户ID分组头部商户的订单量可能占四成点击流按商品ID分组爆款商品的点击量碾压长尾商品。这些Key的分布基本都符合Zipf分布少数Key占据大多数数据。哈希之后这些热点Key还是会落到同一个分区于是那个Subtask就变成了整个作业的瓶颈。这就是问题的本质分区算法是固定的、均匀的但业务数据是不均匀的、长尾的两者天然错配。你可以把这种情况类比成一条高速公路上90%的车都涌向同一个出口而出口只有一个收费通道后面自然排起了长队。换个收费站的名称或者多加几个车道都解决不了问题得想办法让这些车分散到不同的出口出去。除了KeyBy聚合倾斜还会出现在双流Join的场景。Join的关联键如果是用户ID或者商户ID一旦某个关联键的数据量异常巨大对应的那个Join任务就会过载。还有窗口计算比如事件时间窗口如果某个窗口内的数据量特别大而其他窗口很稀疏也会出现类似倾斜的现象。维表关联也一样热门维度被集中请求缓存频繁失效导致同一条链路的压力飙升。2. 四两拨千斤Flink数据倾斜的核心化解方案2.1 聚合场景首选两阶段聚合先说适用范围最广的方案——两阶段聚合国内技术圈喜欢叫它“局部聚合全局聚合”英文里叫LocalAggregate GlobalAggregate。它的核心思路非常朴素既然热点Key都挤在一个分区里那我先给Key加个随机前缀把一个大Key拆成N个小Key让它们分散到不同分区去做第一轮聚合第一轮聚合的结果再去掉前缀按原始Key做第二轮聚合。数学上很简单sum(sum(x)) sum(x)。每一条数据先按自己的子Key做部分求和再把这些部分和汇总就是全量结果。这个方案实现起来也不复杂用Flink DataStream API写的话大致是这样的DataStreamTuple2String, Long source ...; // 第一阶段给Key加随机前缀分散热点 DataStreamTuple2String, Long localAgg source .map(new RichMapFunctionTuple2String, Long, Tuple2String, Long() { Override public Tuple2String, Long map(Tuple2String, Long value) { int salt new Random().nextInt(10); return Tuple2.of(salt _ value.f0, value.f1); } }) .keyBy(t - t.f0) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .reduce((v1, v2) - Tuple2.of(v1.f0, v1.f1 v2.f1)); // 第二阶段去掉前缀按原始Key再做一次汇总 DataStreamTuple2String, Long result localAgg .map(t - Tuple2.of(t.f0.split(_)[1], t.f1)) .keyBy(t - t.f0) .window(TumblingEventTimeWindows.of(Time.minutes(1))) .reduce((v1, v2) - Tuple2.of(v1.f0, v1.f1 v2.f1));有几个细节值得注意。随机盐的范围一般取并行度的2到10倍。取太小了不行10倍的并行度20个盐值热点Key也才拆成20份效果有限取太大了也不行第一阶段产生的中间数据膨胀Shuffle开销反而变大。我平时习惯先取并行度的3倍左右跑一遍看效果再微调。两阶段聚合会带来额外的延迟因为第一阶段的局部聚合窗口结束后还要进入第二阶段的窗口再计算一次整体结果的产出时间会比原来多一个窗口周期。实时性要求特别高的场景要权衡一下这个代价。还有个好消息是如果你用的是Flink SQL两阶段聚合并不需要自己写代码Flink SQL的优化器内置了LocalGlobal机制。开启MiniBatch之后本地聚合会自动执行相当于框架帮你把两阶段聚合做掉了。SET table.exec.mini-batch.enabled true; SET table.exec.mini-batch.allow-latency 5 s; SET table.exec.mini-batch.size 5000; SET table.optimizer.agg-phase-strategy TWO_PHASE;MiniBatch的核心是攒一批数据再处理减少访问状态的次数配合TWO_PHASE的聚合策略倾斜缓解效果是比较明显的。但注意MiniBatch拉长了数据可见的延迟所以需要对实时性要求没那么苛刻的场景才适合。2.2 Join场景盐化Salting与广播流两阶段聚合解决的是聚合场景但数据倾斜绝不只在聚合里出现双流Join同样会倾斜而且解决思路不太一样。双流Join倾斜的套路是“Salting”中文常叫加盐。看名字就知道思路和两阶段聚合的第一步有点像把热点Key进行拆分给关联键加一个随机后缀让原本集中在一个Key上的数据散到多个Key上去。但Join有个特殊性加盐必须对两条流同步进行并且盐的规则要一致。不然你这边把A流的Key加了盐B流的Key没加或者规则不同两边就关联不上了。举个例子订单流和支付流按订单号关联订单号分布均匀通常不会倾斜。但如果是订单流和商户维度流按商户ID关联头部商户的订单量巨大这就麻烦了。处理方式是把热点商户拆成N个虚拟商户比如原本商户ID为10001拆成10001_0、10001_1、...、10001_9。订单流打散的时候对热点商户随机分配一个后缀商户维表里则把这一个商户复制成10份每份对应一个后缀。这样两边就能正常关联同时压力被分散到10个分区。// 订单流热点商户随机加盐 if (hotMerchantIds.contains(merchantId)) { int salt random.nextInt(10); emit(merchantId _ salt, order); } else { emit(String.valueOf(merchantId), order); } // 商户维表热点商户复制10份 for (int salt 0; salt 10; salt) { emit(merchantId _ salt, merchant); }另一类常见场景是大表Join小表比如实时流关联一个全量维度表。这时候最省事的方案就是广播把维度表通过BroadcastStream发送到每个TaskManager内存里Join直接在本地完成压根不经过Shuffle也就不存在分区倾斜的问题。Flink的Broadcast机制做这个非常成熟唯一的要求是小表的数据量得控制在GB级别以内太大了内存扛不住。之前在我维护的一个任务里就把一个几百MB的商户维表改成广播整体吞吐上了一个台阶延迟反而降了。但广播流也有它的软肋那就是维度表更新不及时。如果用了RocksDB或外部存储作为维表热点的维度Key频繁命中缓存导致单节点压力大此时更推荐用异步I/O加缓存的方式而不是广播。2.3 兜底手段自定义Partitioner与侧输出流前面说的两阶段聚合和加盐核心都是把热点Key打散。但还有一种场景你没法提前打散因为热点Key是动态出现的今天这个商户爆单明天那个商户爆单写成死代码肯定不行。这时候有两种兜底方案。第一种是自定义Partitioner实现Partitioner接口在分区阶段就干预数据的分发逻辑。比如维护一个动态更新的热点Key集合对于热点Key计算一个额外的哈希偏移让它们落到不同的分区。这个方案对KeyBy之前的算子有效但对KeyBy之后的窗口聚合来说下游还是按照Key路由所以更适用于那些没有分组语义、只是想均匀分发数据的场景。第二种是侧输出流Side Output。这是我觉得最灵活的思路在进入KeyBy之前先用ProcessFunction对数据做识别判断是不是热点Key。如果是热点数据用侧输出流单独分离出来正常的非热点数据则继续走主流程。分离出来的热点流可以单独设置并行度、单独开窗口、单独调优甚至落到独立的Sink。这样热点数据就不会拖累主流程整个作业的稳定性会好很多。final OutputTagEvent hotTag new OutputTagEvent(hot) {}; DataStreamEvent mainStream source.process(new ProcessFunctionEvent, Event() { Override public void processElement(Event value, Context ctx, CollectorEvent out) { if (hotKeyDetector.isHot(value.getUserId())) { ctx.output(hotTag, value); } else { out.collect(value); } } }); DataStreamEvent hotStream mainStream.getSideOutput(hotTag);侧输出流的坏处是两套链路要各自维护计算口径容易漂移。比如热点流和主流的统计逻辑各写一套如果统计口径有变动两边都要改容易漏。所以这个方案适合当临时救火手段或者配合配置中心做动态开关不适合作为长期架构。3. 实战案例MySQL同步到ClickHouse的倾斜排查与修复3.1 一个真实的同步链路与第一现场前阵子帮一个客户排查过MySQL同步ClickHouse的任务链路大概是这样MySQL的Binlog通过CDC组件采集进KafkaFlink消费Kafka后做实时清洗和聚合最后通过JDBC连接器写入ClickHouse。任务上线第一个星期跑得挺稳后来业务量上来头部商户的订单量暴涨问题就集中爆发了。症状也典型ClickHouse的Sink偶尔报错报的是Connection is not available有时候是Communications link failure任务过几分钟又自动恢复但每隔一段时间就重复一次。一开始我们以为是ClickHouse连接池配置有问题后来把Web UI打开一看问题根源找到了——Flink作业里有一个按商户ID做统计的窗口算子上某个Subtask的输入数据量是其他节点的7倍多反压一直顶到Kafka Source。下游的JDBC Sink批量写入因为上游数据涌得太多连接被拖死才报出一堆连接异常。所以很多时候JDBC连接器异常只是一个表象真正的病根还在上游的数据倾斜。遇到类似错误不要急着调连接参数先把整条链路的数据分布看一遍。3.2 修复过程两阶段聚合 Sink参数调优确认倾斜之后修复思路很清晰。因为业务统计就是按商户ID分组之后做SUM和COUNT属于典型的聚合倾斜直接上两阶段聚合。我们在这里没有改代码而是把作业从DataStream改成Flink SQL用MiniBatch LocalGlobal参数解决。CREATE TABLE order_stat_sink ( merchant_id BIGINT, order_count BIGINT, total_amount DECIMAL(20, 2), window_start TIMESTAMP(3), window_end TIMESTAMP(3) ) WITH ( connector jdbc, url jdbc:clickhouse://clickhouse-host:8123/default, table-name order_stat, username default, password ***, sink.buffer-flush.max-rows 2000, sink.buffer-flush.interval 3s, sink.max-retries 3 );这里多说一句Sink参数。JDBC连接器默认的flush触发条件是缓冲行数达到sink.buffer-flush.max-rows或者间隔时间到sink.buffer-flush.interval两个条件先到先触发。之前客户默认配置是5000行和1秒。问题是ClickHouse单次插入的行数太大容易触发服务端的part合并瓶颈反而拖慢写入还可能把连接干超时。我们调到2000行和3秒之后单批数据量更可控写入抖动明显减轻。当然Flink SQL的LocalGlobal对SUM和COUNT是有效的但对那种必须先精确去重再计数的场景效果会打折扣因为两阶段聚合的中间状态没法做去重。这种情况就得回到DataStream里手写加盐逻辑或者用HyperLogLog之类的近似算法来顶。3.3 JDBC连接器异常的排查实录趁着这个机会把常见的JDBC连接器异常一起复盘一下都是这个链路里真实遇到过的。第一种是Connection is not available。Flink的JDBC连接器内部用的是连接池连接被MySQL或ClickHouse服务端关闭之后如果连接池没有及时剔除失效连接再借出来用就会报这个错。解决思路不是修改连接串加autoReconnect那么简单而是要控制连接的闲置时间和重试逻辑。在Flink的JDBC connector里没有专门的keepalive参数我的经验是尽量缩短单批写入间隔让连接不容易闲置同时把sink.max-retries设置成3让偶发的连接失败能自动重试。第二种是BatchUpdateException或者SQLException这类通常是数据本身的问题。MySQL里某个字段是JSON类型但ClickHouse表里对应的是String写入时会因为类型不匹配失败MySQL的decimal(20,4)映射到ClickHouse的Decimal(20,4)没问题但如果是decimal(20,4)映射到Float64精度就会丢。我的建议是同步链路的表结构尽量一一对应字段类型提前核对好不要指望连接器做魔法转换。第三种是ClickHouse写入幂等性的问题。Flink任务重启或者回放Kafka数据的时候同一个统计结果可能会写两次。ClickHouse本身不支持UPDATE和DELETE重复插入就会产生重复数据。解决这个问题通常要建ReplacingMergeTree或者CollapsingMergeTree引擎表利用版本号或者正负号来做去重。这块虽然和JDBC连接器异常关系不大但实际上很多“数据变多”的异常都是从这里冒出来的值得同步留意。4. 上下游联调SpringBoot整合Flink的几个坑4.1 先想清楚你要的是哪种整合模式这几年微服务很普及很多团队希望把Flink任务纳入已有的SpringBoot体系里统一管理。但这里有个大坑Flink是跑在独立集群上的分布式计算框架SpringBoot是跑在应用服务器里的Web服务两者根本不是一回事。不少人第一次整合的时候直接在SpringBoot的进程里new一个StreamExecutionEnvironment然后execute()以为这样就是把Flink集成了。这个做法在本地调试时没问题但到了生产环境就会出现类加载冲突、资源争抢、任务重启管理混乱等一系列问题。真正的整合模式应该选这两种。第一种是SpringBoot作为任务管理平台通过Flink REST API向独立集群提交Jar包任务运行在Flink集群上SpringBoot只管把参数传过去、监控运行状态。第二种是SpringBoot服务作为数据生产者通过Kafka或HTTP把数据交给FlinkFlink任务消费数据SpringBoot不直接碰Flink的运行时。4.2 依赖冲突最经典的翻车现场我见过太多人倒在这一步。SpringBoot自带的依赖管理加上Flink客户端依赖两者会对齐到完全不同的版本。比如Flink 1.17用的Guava版本和SpringBoot内嵌的Guava版本不一致启动时直接NoSuchMethodErrorJackson也是重灾区Flink对Jackson有自己的序列化需求SpringBoot的Jackson版本可能把Flink的序列化器打崩。解决办法也很成熟要么在SpringBoot里把Flink相关依赖的scope全部设为provided把依赖交给Flink集群要么用Maven Shade插件把Flink依赖relocate到自定义的包路径下避免和SpringBoot的类冲突。从我自己的经验来看如果SpringBoot只是提交和管理任务最好的做法是根本不要引Flink客户端依赖直接用HTTP调用REST API问题从根上就没了。4.3 一个简单的Flink任务管理工具类写一个轻量的SpringBoot整合方式其实并不复杂。用RestTemplate或者OkHttp调用Flink的REST API就能完成上传、提交、查询、取消的闭环。核心就几个接口。// 1. 上传Jar包 POST /api/v1/jars/upload // 2. 提交任务返回jobId POST /api/v1/jars/{jarId}/run { programArgsList: [--job, order-stat-sync] } // 3. 查询任务状态 GET /api/v1/jobs/{jobId} // 4. 取消任务 PATCH /api/v1/jobs/{jobId}?modecancelSpringBoot里封装一下就是Component public class FlinkJobClient { private final RestTemplate restTemplate; public FlinkJobClient(RestTemplate restTemplate) { this.restTemplate restTemplate; } public String submit(String managerUrl, String jarId, ListString args) { String url managerUrl /api/v1/jars/ jarId /run; MapString, Object body new HashMap(); body.put(programArgsList, args); ResponseEntityString response restTemplate.postForEntity( url, body, String.class); return response.getBody(); } }这个方案里有个容易被忽略的点Flink的REST API默认不开启鉴权如果你的Flink集群暴露在非受信网络里一定要在Flink配置里开启安全认证或者在SpringBoot和集群之间加一层网关做白名单。还有REST API的超时时间要设置合理上传Jar包时大文件可能需要几十秒默认的HTTP连接超时太短会直接失败。5. 避坑速查表与经验心得5.1 数据倾斜排查速查表把常见现象、可能原因、定位方法和解决手段整理成一个表平时排查的时候照着看一眼比从头分析快得多。现象可能原因定位方法解决手段单个Subtask数据量远高于其他Key分布不均热点Key集中Web UI查看各Subtask的Records Received两阶段聚合、Salting、侧输出流分流整条链反压但CPU利用率参差不齐分区哈希与Key分布错配Metrics按算子查看BackPressure自定义Partitioner、调整并行度并观察效果Checkpoint反复超时倾斜节点状态过大查看Checkpoint详情的State Size开启RocksDB增量快照、加盐打散热点Sink写报Connection异常上游倾斜导致Sink瞬时压力过大查看Sink端error log与上游数据分布先治上游倾斜再调buffer-flush与重试MySQL同步ClickHouse数据重复缺少幂等机制对比主键数据量换ReplacingMergeTree/CollapsingMergeTree引擎5.2 几条实在的经验加并行度解决不了倾斜这是我反复见到有人踩的坑。有人看到Web UI里一个节点数据量大直接把这个算子的并行度从4调到16结果热点Key还是落在同一个分区里只是那个分区的任务被拆成了更多slot数据还是在同一台机器上。并行度只能让无状态且数据均匀的负载受益对热点Key的倾斜无能为力。定位数据倾斜一定要用数据说话。打开Web UI看每个Subtask的numRecordsInPerSecond按时间维度对比一下就能清楚看到热点节点和其他节点的流量差距。不要上来就猜是GC问题、网络问题还是代码问题先确认是不是倾斜再对症下药。方案选型要按场景来。聚合类倾斜首选两阶段聚合或者Flink SQL的LocalGlobal。Join类倾斜首选Salting加盐或者广播流。维表关联热点优先考虑缓存加异步I/O。动态热点无法预判用侧输出流做兜底。没有一套方案能通吃所有场景组合起来用才是正路。热点Key要做成动态配置不要写死在代码里。我在前面的实战案例中把热点商户集合放到配置中心检测到新的热点商户后配置下发任务不重启也能生效。线上跑了一段时间稳定性比原来好很多。最后补一句收尾的话。我在实际排查中踩过几次坑之后最大的体会是数据倾斜不是一次性修复的问题而是一个需要持续监控的动态问题。业务的Key分布随着时间不断变化今天的热点明天可能就凉了今天不起眼的Key明天可能就爆了。把倾斜识别做成监控项把缓解方案做成标配能力整个Flink作业的稳定性才能真正上一个台阶。希望这篇博文的经验对你有用有问题评论区一起聊。

关于恒美微站

恒美微站专注于为个体商户、工作室提供极简自助建站服务,让每个人都能轻松拥有专业网站。

快速链接

  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心

服务项目

  • 可视化建站
  • 拖拽编辑
  • 主题定制
  • SEO 优化
  • 网站托管

联系方式

  • 📍 地址:北京市朝阳区建国路 88 号
  • 📞 电话:400-888-8888
  • ✉️ 邮箱:info@hmyw.cn
  • 🕐 时间:周一至周日 9:00-18:00

© 2024 恒美微站 hmyw.cn 版权所有 | 京 ICP 备 12345678 号