精选·Flink

Flink 数据倾斜的表现与处理方案

91学AI·2026/7/13·7 阅读

考察点

数据倾斜是分布式计算的老大难,这道题考的是「识别 + 对症下药」的完整方法论。面试官想听到:怎么确认是倾斜、不同场景(聚合/Join)的解法差异、加盐打散的具体做法和代价。追问会往两阶段聚合的适用前提、Join 倾斜的特殊处理走。

参考答案

怎么确认是倾斜

Web UI 上看同一个算子的各个 subtask:Records Received / Sent 差异悬殊(有的几千万,有的几十万), busyTime 也两极分化——忙的 subtask 100% 满载,闲的在睡觉。作业整体吞吐被最慢的那个 subtask 拖死(分布式任务的木桶效应),同时伴随反压和 Checkpoint 对齐时间长。倾斜的本质是 key 分布不均:keyBy 之后某个热点 key(大促时的头部商家、明星用户的 ID、空值或默认值 key)的数据全落到一个 subtask 上。

聚合场景的解法

轻度倾斜:两阶段聚合(local-global)。先在本地做预聚合再 shuffle,类似 MapReduce 的 combiner。DataStream 里没有自动 combiner,要手动实现:map 里用短窗口小状态攒批,或者干脆拆成「按 (key, 桶号) 先聚一次、再按 key 聚第二次」。Flink SQL 可以开 table.optimizer.agg-phase-strategy=TWO_PHASE 加 localAgg,对 SUM/COUNT 这类可拆分聚合很有效。

重度倾斜(单个超级热点 key):加盐打散。给 key 拼上随机后缀 key + "_" + random(0, N),热点 key 被拆成 N 份散到 N 个 subtask 各自聚合,得到部分结果后再去盐按原 key 二次聚合汇总。N 取多少看热点倍数,热点占 30% 流量时 N=10 通常足够。代价是多一轮 shuffle 和聚合,代码复杂度也上去了,只用于确认的热点。

Join 场景的解法

Join 倾斜更麻烦,因为状态也要按 key 存。几个方案按优先级:

  • 空值/脏 key 过滤。很多倾斜是空 key、默认值 key("-"、"unknown")造成的,上游过滤或单独处理,一行代码解决一半问题。
  • 热点 key 单独处理。把 Top N 热点 key 名单(从离线统计或配置中心拿)拆出来单独走一条链路——比如热点侧 broadcast,非热点正常 keyBy,最后 union。常见在头部 KOL、旗舰店铺场景。
  • 小表 broadcast。一侧是小维表就别 keyBy 了,广播到所有实例本地 Join,倾斜问题直接消失。
  • 加盐 Join。两流 Join 都倾斜时,热点侧加盐打散,另一侧数据复制 N 份对应打散后的 key,能匹配上但数据量放大 N 倍,要算好成本,属于没办法时的办法。

治本与治标

最后说个视角:加盐、两阶段都是治标,治本往往在数据源头——能不能在上游就按更细的粒度预聚合(埋点侧先按分钟聚合再上报)、能不能把大 key 业务上拆分(大商家按子店铺拆)。面试里把「先过滤脏 key → 再考虑两阶段/加盐 → 极端情况热点单独处理」的决策顺序讲清楚,比罗列方案更有说服力。

可能的追问

  • 加盐后结果怎么保证正确?——要求聚合函数满足结合律(SUM/COUNT/MAX 都行),部分结果二次聚合数学上等价;UV 类(HyperLogLog 或 distinct count)要用支持合并的近似结构,不能直接 sum。
  • 倾斜只在高峰期出现怎么办?——流量洪峰放大热点效应,确认热点 key 是否随时间变化,用动态名单方案;同时保证资源能扛峰值时段最热 subtask 的量。
  • Window All(并行度 1)是不是必然倾斜?——是单点不是倾斜,所有数据过一个实例,只适合小流量;大流量的全流聚合应该先 keyBy 打散再二次汇总。

评论 (0)

暂无评论,快来抢沙发吧!

91学AI

© 2026 91学AI · 按岗位学 AI 与大数据. All rights reserved.