什么是数据倾斜,会带来哪些问题?如何解决?


一则或许对你有用的小广告

欢迎加入小哈的星球,你将获得:专属的实战项目(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+ 小伙伴加入学习,欢迎点击围观

面试考察点

  1. 基础概念掌握度:面试官想确认你是否真的理解分区 / 分组操作在分布式计算中的本质——shuffle 之后数据如何重新分布,而不只是会背一句“数据分布不均”。
  2. 问题诊断能力:能否从症状(任务卡在 99%、某个 task OOM、整体执行时间被拖长)反推出是数据倾斜,而不是网络、GC 或代码 bug。说白了,看你有没有线上排障经验。
  3. 解决方案的深度与广度:是否能根据不同场景group byjoincount(distinct))给出对应的解法,而不是只会一句“加分区数”。这个是拉开差距的关键。

核心答案

数据倾斜:在分布式并行计算中,数据按某个 key 分区或分组后,绝大部分数据集中到了少数几个 task / 节点上,导致负载严重不均衡的现象。

核心危害一句话:木桶效应——整个任务的执行时间被最慢的那个 task 拖死,其他节点干等着。

解决思路也一句话:让分布不均的 key 重新打散,或者绕开 shuffle

场景 典型症状 首选解法
group by 倾斜 某 reduce task 数据量爆炸 两阶段聚合(加盐 + 局部聚合 + 全局聚合)
大表 join 小表 join 阶段卡死 Map Join(广播小表)
大表 join 大表 reduce 端 join 拖死 扩容 + 随机前缀打散
空值 / 异常值集中 某 key 是 null 占比 80% 提前过滤或单独处理
框架级兜底 不想改代码 Spark AQE / Hive 倾斜开关

深度解析

一、为什么会发生数据倾斜

分布式计算里只要触发 shufflegroup byjoindistinctorderBy 等都会触发),数据就要按 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 亿条 hot key 被切分成 10 份,每份 1 千万,分散到不同分区
  • 局部聚合:每个分区在本地做一次 group by,把同一前缀的数据先聚合一次,输出 0_hot: xxx1_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。

Map Join 广播小表
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=truespark.sql.adaptive.skewJoin.enabled=true,运行时自动检测倾斜并拆分大分区,这是目前生产环境最省心的方案。顺带说一句,从 Spark 3.2 开始 AQE 已经默认开启了,所以你只要用相对新的版本,啥都不用配就能享受这个福利。
  • Hiveset hive.groupby.skewindata=true;(自动两阶段聚合)、set hive.optimize.skewjoin=true;(倾斜 join 优化),再加一个 hive.skewjoin.key(默认 100000 行)控制“超过多少行才算倾斜键”。
  • MapReduce:合理设置 Combiner,减少 shuffle 数据量。

Spark AQE 倾斜优化
Spark AQE 倾斜优化

面试高频追问

  1. 追问一:如何快速定位是不是数据倾斜?
    • 任务进度——卡在 99% 长时间不动是典型症状;看 任务日志——有没有某个 task 处理的数据量是其他 task 的几十倍;看 Spark UI / ResourceManager——是不是只有一两个 Executor 满载,其他全闲着。
  2. 追问二:两阶段聚合为什么有效?
    • 因为它把单次聚合变成 “先分散聚合、再汇总聚合”。第一次 shuffle 时热点 key 已经被打散,没有倾斜;第二次 shuffle 时数据量已经被局部聚合压缩到很小,即便集中也不会再倾斜。
  3. 追问三:Map Join 和 Reduce Join 的本质区别是什么?
    • Map Join 在 Map 端完成,不需要 shuffle,小表广播;Reduce Join 必须经过 shuffle,按 key 重分发后在 Reduce 端匹配。前者适合大表 join 小表,后者是大表 join 大表的无奈之选。
  4. 追问四:Spark 3.0 的 AQE 是怎么自动处理倾斜的?
    • 运行时统计每个分区大小,发现某个分区远大于中位数时,自动把这个分区按 splits 大小拆分成多个子分区,由不同的 task 并行处理,最后合并结果。本质上是把倾斜处理从“代码层”下沉到了“执行计划层”。

常见面试变体

  • “Spark 任务卡在 99% 不动,可能是什么原因?”(经典切入式提问,数据倾斜是头号嫌疑)
  • group byjoin 的倾斜分别怎么处理?”
  • “为什么 count(distinct) 容易倾斜?有什么替代方案?”
  • “分库分表场景下会不会有数据倾斜?怎么解决?”(热点分片问题,比如订单表按 userId 分片,头部主播数据全集中到一个库)

记忆口诀

“先诊断、再过滤、能广播就广播、不能广播就打散、最后靠框架兜底”

按场景选解法:

  • 聚合倾斜 → 两阶段(加盐 + 局部聚合 + 全局聚合)
  • 小表 join → Map Join(广播)
  • 大表 join → 扩容 + 随机前缀
  • 空值集中 → 预处理
  • 不想改代码 → AQE / Hive 倾斜开关

总结

数据倾斜的本质是 shuffle 后 key 分布不均,带来的危害一句话——木桶效应,最慢的 task 拖死整个 Job。解决思路就两条主线:打散热点 key,或者干脆绕开 shuffle。把每个场景的解法对应上,再加上 Spark 3.0 AQE 这个现代框架的兜底方案,这题基本能答得有深度。