什么是数据倾斜,会带来哪些问题?如何解决?
一则或许对你有用的小广告
欢迎加入小哈的星球,你将获得:专属的实战项目(4个项目都能学) / 1v1 提问 / 简历修改 / Java 学习路线 / 社群讨论 / 学习打卡 / 每月赠书
《Spring AI 项目实战(问答机器人、RAG 智能客服、联网搜索)》已完结,基于
Spring AI + Spring Boot 3.x + JDK 21...,查看介绍《从零手撸:仿小红书(微服务架构)》 已完结,基于
Spring Cloud Alibaba + Spring Boot 3.x + JDK 17...,查看介绍;演示链接:http://116.62.199.48:7070/《从零手撸:前后端分离博客项目(全栈开发)》 2 期已完结,演示链接:http://116.62.199.48/
新开坑项目:《从零手撸:秒杀系统高并发优化实战》 正在更新中...,查看介绍
截止目前,星球内专栏累计输出 150w+ 字,讲解图 5110+ 张,还在持续爆肝中.. 后续还会上新更多项目,已有 4700+ 小伙伴加入学习,欢迎点击围观
面试考察点
- 基础概念掌握度:面试官想确认你是否真的理解分区 / 分组操作在分布式计算中的本质——shuffle 之后数据如何重新分布,而不只是会背一句“数据分布不均”。
- 问题诊断能力:能否从症状(任务卡在 99%、某个 task OOM、整体执行时间被拖长)反推出是数据倾斜,而不是网络、GC 或代码 bug。说白了,看你有没有线上排障经验。
- 解决方案的深度与广度:是否能根据不同场景(
group by、join、count(distinct))给出对应的解法,而不是只会一句“加分区数”。这个是拉开差距的关键。
核心答案
数据倾斜:在分布式并行计算中,数据按某个 key 分区或分组后,绝大部分数据集中到了少数几个 task / 节点上,导致负载严重不均衡的现象。
核心危害一句话:木桶效应——整个任务的执行时间被最慢的那个 task 拖死,其他节点干等着。
解决思路也一句话:让分布不均的 key 重新打散,或者绕开 shuffle。
| 场景 | 典型症状 | 首选解法 |
|---|---|---|
group by 倾斜 |
某 reduce task 数据量爆炸 | 两阶段聚合(加盐 + 局部聚合 + 全局聚合) |
| 大表 join 小表 | join 阶段卡死 | Map Join(广播小表) |
| 大表 join 大表 | reduce 端 join 拖死 | 扩容 + 随机前缀打散 |
| 空值 / 异常值集中 | 某 key 是 null 占比 80% | 提前过滤或单独处理 |
| 框架级兜底 | 不想改代码 | Spark AQE / Hive 倾斜开关 |
深度解析
一、为什么会发生数据倾斜
分布式计算里只要触发 shuffle(group by、join、distinct、orderBy 等都会触发),数据就要按 key 重新分发到下游。分发规则一般是 hash(key) % N。
问题就出在 key 的分布上:如果某个 key 的数据量远超其他 key(比如一个电商平台里某头部主播的订单数据占全部订单的 30%),那么这个 key 所在的分区的数据量就是其他分区的几十倍甚至上百倍。
下面这张图对比了正常分布和倾斜分布:
上图对比了两种数据分布情况:
- 正常分布:100G 数据大致均匀打散到 4 个分区,每个分区处理 25G,所有 task 几乎同时完成,资源利用率高
- 倾斜分布:某个热点 key 占了 95G 数据,几乎全部涌进同一个分区,这一个 task 处理 95G,其他三个 task 早早结束干等着
关键认知:分布式计算的总耗时 ≈ 最慢 task 的耗时,不是平均耗时。所以倾斜的代价是线性放大的。
二、数据倾斜带来的具体问题
这块很多人只会说“慢”,其实问题有 4 个层次,面试时分得越清楚越显深度:
- 执行时间被拖长:最直观的问题,整个 Job 卡在最后 1% 长时间不动,进度条上表现为 “卡 99%”。我之前排查过一次,99% 卡了 40 分钟,最后发现就是某个用户 ID 的数据占了 60%。
- 内存溢出(OOM):倾斜 task 处理的数据量远超预期,
Reduce端维护的 HashMap / 排序缓冲区被撑爆,直接抛 OOM。 - 资源利用率极低:99 个 task 早早完成、CPU 闲置,1 个 task 满载运行。你申请了 100 个核心的资源,实际只在用 1 个——纯烧钱。
- 任务最终失败:拖到超时阈值,或者 OOM 直接挂掉,整张链路失败重跑,陷入死循环。
三、解决方案:分场景拆解
数据倾斜没有“万能药”,必须对症下药。下面按场景拆开讲。
场景 1:group by / 聚合类倾斜 → 两阶段聚合
这是最常用、面试也最常考的解法。核心思路是把一次聚合拆成两次:
- 第一阶段:给每个 key 加上随机前缀(比如
0~9的随机数),把热点 key 打散到多个分区做局部聚合 - 第二阶段:去掉前缀,恢复原始 key,做全局聚合
上图展示了两阶段聚合的核心流程,可以分三步理解:
- 加盐打散:给原本集中的 key 拼接一个
0~N的随机前缀,让 1 亿条hotkey 被切分成 10 份,每份 1 千万,分散到不同分区 - 局部聚合:每个分区在本地做一次
group by,把同一前缀的数据先聚合一次,输出0_hot: xxx、1_hot: xxx…… - 全局聚合:去掉前缀后,数据量已经从 1 亿骤减到 10 条(每个前缀一条),再做一次
group by就能拿到最终结果
说白了就是 MapReduce Combiner 思想的扩展。
伪代码示意:
// 第一阶段:加盐 + 局部聚合
// 原始数据 (userId, 1),userId 为热点 key
JavaPairRDD<String, Long> stage1 = data
.mapToPair(t -> {
// 加随机前缀 0~9
int prefix = new Random().nextInt(10);
return new Tuple2<>(prefix + "_" + t._1, t._2);
})
.reduceByKey(Long::sum); // 局部聚合
// 第二阶段:去前缀 + 全局聚合
JavaPairRDD<String, Long> stage2 = stage1
.mapToPair(t -> {
// 去掉前缀,还原原始 key
String originalKey = t._1.substring(t._1.indexOf("_") + 1);
return new Tuple2<>(originalKey, t._2);
})
.reduceByKey(Long::sum); // 全局聚合
坑点提醒:
count(distinct)这种操作不能直接套用两阶段聚合,因为 distinct 的语义不能拆分。一般的做法是先 group by 去重,再 count,或者在数据量极大时用 近似计数(HyperLogLog) 牺牲一点精度换性能。
场景 2:大表 join 小表 → Map Join(广播 Join)
如果一边是大表、一边是小表(比如维表、字典表),根本不应该走 reduce join。直接把小表广播到每个 Map 节点,在 Map 端完成 join,彻底绕开 shuffle,倾斜也就不存在了。
// Spark 广播 Join 示例
Broadcast<Table> smallTableBC = sc.broadcast(smallTable.collectAsMap());
JavaRDD<ResultRow> result = largeTable
.map(row -> {
Table small = smallTableBC.value();
// 在 Map 端直接匹配,无 shuffle
return joinWith(row, small.get(row.getKey()));
});
或者用框架的自动广播:在 Spark SQL 里,只要小表小于 spark.sql.autoBroadcastJoinThreshold(默认 10MB),就会自动转 Map Join。
场景 3:大表 join 大表 → 扩容 + 随机前缀
两边都很大、没法广播,怎么办?思路是人为把小一点的那张表“扩容”,大表加随机前缀,让 join 能打散:
- 给大表的 key 加
0~N的随机前缀 - 把小表的每条数据复制 N 份,每份加上对应前缀
0~N-1 - 这样大表的任意一个 key 都能在小表里找到 N 个匹配项
代价是小表数据量扩大 N 倍,所以 N 一般取 5~10,别太离谱。
场景 4:空值 / 异常值集中 → 预处理
这个最容易被忽略。我见过一个 case:日志数据里 80% 是 userId=null(爬虫、未登录用户),结果按 userId 分组时全部挤到一个分区。
解法:
- 能过滤就过滤:
where user_id is not null - 不能过滤就打散:给 null 值分配一个随机 key,比如
concat('null_', rand()),让它们均匀分散
场景 5:框架级兜底 → 开启自适应优化
如果不想改代码(线上老系统动不得),可以靠框架的内置倾斜处理:
- Spark 3.0+ AQE(自适应查询执行):开启
spark.sql.adaptive.enabled=true和spark.sql.adaptive.skewJoin.enabled=true,运行时自动检测倾斜并拆分大分区,这是目前生产环境最省心的方案。顺带说一句,从 Spark 3.2 开始 AQE 已经默认开启了,所以你只要用相对新的版本,啥都不用配就能享受这个福利。 - Hive:
set hive.groupby.skewindata=true;(自动两阶段聚合)、set hive.optimize.skewjoin=true;(倾斜 join 优化),再加一个hive.skewjoin.key(默认 100000 行)控制“超过多少行才算倾斜键”。 - MapReduce:合理设置
Combiner,减少 shuffle 数据量。
面试高频追问
- 追问一:如何快速定位是不是数据倾斜?
- 看 任务进度——卡在 99% 长时间不动是典型症状;看 任务日志——有没有某个 task 处理的数据量是其他 task 的几十倍;看 Spark UI / ResourceManager——是不是只有一两个 Executor 满载,其他全闲着。
- 追问二:两阶段聚合为什么有效?
- 因为它把单次聚合变成 “先分散聚合、再汇总聚合”。第一次 shuffle 时热点 key 已经被打散,没有倾斜;第二次 shuffle 时数据量已经被局部聚合压缩到很小,即便集中也不会再倾斜。
- 追问三:Map Join 和 Reduce Join 的本质区别是什么?
- Map Join 在 Map 端完成,不需要 shuffle,小表广播;Reduce Join 必须经过 shuffle,按 key 重分发后在 Reduce 端匹配。前者适合大表 join 小表,后者是大表 join 大表的无奈之选。
- 追问四:Spark 3.0 的 AQE 是怎么自动处理倾斜的?
- 运行时统计每个分区大小,发现某个分区远大于中位数时,自动把这个分区按 splits 大小拆分成多个子分区,由不同的 task 并行处理,最后合并结果。本质上是把倾斜处理从“代码层”下沉到了“执行计划层”。
常见面试变体
- “Spark 任务卡在 99% 不动,可能是什么原因?”(经典切入式提问,数据倾斜是头号嫌疑)
- “
group by和join的倾斜分别怎么处理?” - “为什么
count(distinct)容易倾斜?有什么替代方案?” - “分库分表场景下会不会有数据倾斜?怎么解决?”(热点分片问题,比如订单表按
userId分片,头部主播数据全集中到一个库)
记忆口诀
“先诊断、再过滤、能广播就广播、不能广播就打散、最后靠框架兜底”
按场景选解法:
- 聚合倾斜 → 两阶段(加盐 + 局部聚合 + 全局聚合)
- 小表 join → Map Join(广播)
- 大表 join → 扩容 + 随机前缀
- 空值集中 → 预处理
- 不想改代码 → AQE / Hive 倾斜开关
总结
数据倾斜的本质是 shuffle 后 key 分布不均,带来的危害一句话——木桶效应,最慢的 task 拖死整个 Job。解决思路就两条主线:打散热点 key,或者干脆绕开 shuffle。把每个场景的解法对应上,再加上 Spark 3.0 AQE 这个现代框架的兜底方案,这题基本能答得有深度。
