精选·Spark

reduceByKey 和 groupByKey 有什么区别?为什么生产上基本不用 groupByKey?

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

考察点

经典的性能对比题,但很多人只会背「reduceByKey 有 combine」。面试官想听的是:combine 到底在什么位置发生、削减的是哪段的开销;什么场景下两者差别会大到 OOM 和毫秒级的差距;以及能不能举一反三说出 combineByKey、aggregateByKey、foldByKey 的家族关系。追问常往「groupByKey 有没有正当用途」「为什么 map 端 combine 不保证执行」「数据倾斜时 reduceByKey 还会 OOM 吗」这些方向走。

参考答案

一句话版本

reduceByKey 在 shuffle 之前、map 端先做本地聚合(combine),传给下游的是聚合后的结果;groupByKey 不做任何预聚合,把同 key 的所有原始记录原样 shuffle 到 reduce 端再分组。两者结果可以等价,但 shuffle 的数据量完全不是一个量级。

用数字感受差距

假设一个日志聚合场景:1 亿条 (word, 1) 记录,去重后只有 100 万个不同 word,分布在 1000 个 map 分区上。

  • groupByKey:1 亿条记录全部序列化、写 shuffle 文件、跨网络传输、reduce 端再反序列化。网络搬运的是 1 亿条原始记录。
  • reduceByKey:每个 map 分区先在内存里用 hash 表做本地累加,一个分区里同一个 word 最终只剩一条 (word, count)。最坏情况(每个词在每个分区都出现)map 端输出 1000 × 100 万 = 10 亿条?不对——每个分区输出的是该分区内去重后的词,上限是「分区数 × 词数」,但实际日志里单个分区的词量远小于全局词表,比如单分区 20 万个词,那 shuffle 总量就是 1000 × 20 万 = 2 亿条?仍看着多,但对比 1 亿条全量原始记录(每条还短),实际生产里 combine 削减率轻松到 10 倍以上——一个分区 10 万条日志可能只聚出几千个词。

削减率取决于单个 map 分区内的 key 重复度:日志计数、PV/UV 统计这类重复度极高的场景,combine 能砍掉 90% 以上的 shuffle 量;如果 key 几乎不重复(比如按订单 ID 聚合,一单一 key),combine 没收益,但也无损失。

内存与稳定性的差别更致命

数据量差异只是其一,内存行为才是关键:

  • reduceByKey 的 reduce 端拿到的是已经 combine 过的数据,聚合可以用外部排序 + spill 完成,内存可控。
  • groupByKey 的 reduce 端要把同 key 的所有 value 聚到一起形成 Iterator。如果实现需要把这个组物化(Spark 某些路径会这样),而某个 key 的数据量特别大(热点 key),单个 executor 的内存根本装不下,直接 OOM。这就是为什么「生产上禁用 groupByKey」不是教条,是血的教训——它把一个可以流式处理的聚合退化成全量物化。

算子家族关系

这四个算子底层都是 combineByKey 的不同封装:

// reduceByKey:初值就是第一条记录本身
reduceByKey(f) = combineByKey(identity, f, f)

// aggregateByKey:自带零值,map 端 combine 和 reduce 端 merge 可以不同
aggregateByKey(zero)(seqOp, combOp) = combineByKey(seqOp(zero, _), seqOp, combOp)

// foldByKey:aggregateByKey 的特例,seqOp == combOp
// groupByKey:combineByKey 但 combine 不做归约,只收集

aggregateByKey 是最实用的:比如求平均值需要 (sum, count) 二元组,零值 (0, 0),seqOp 累加,combOp 合并。foldByKey 适合有幺元的操作(求和、求积)。

groupByKey 的正当用途

不是永远不能用。两类场景它是对的:

  1. 确实需要全组数据:比如按用户分组后要取该用户行为序列做 session 分析、或要保留全部明细做窗口计算,这时本来就要全部数据,combine 无从谈起——但注意这时更常见的是用 DataFrame 的 window/groupBy + collect_list,有 Catalyst 优化。
  2. combine 语义不存在:聚合逻辑不满足结合律/交换律(比如拼接保序),预聚合无意义。

另外 Spark SQL 的 groupBy 走的是另一套实现(HashAggregateExec → SortAggregateExec),有 code-gen 和溢出保护,和 RDD 的 groupByKey 不是一回事,别混着说。

一个反直觉的点:combine 不保证执行

map 端 combine 是尽力而为的——spark 在 shuffle write 时如果内存缓冲装不下会先 spill,combine 只在缓冲内发生。所以业务逻辑绝不能依赖「combine 一定执行过一次」这种假设,combine 函数必须满足结合律和交换律,且结果和只做一次 reduce 端聚合完全一致。这也是 combineByKey 要求 mergeValue 和 mergeCombiners 语义一致的原因。

可能的追问

  • 热点 key 下 reduceByKey 还会 OOM 吗? 会。combine 削减的是记录条数,但单个超大 key 在 reduce 端聚合时如果中间结果本身巨大(比如 collect 类语义),仍可能撑爆。这时要配合数据倾斜处理手段:加盐打散、两阶段聚合。
  • aggregateByKey 和 reduceByKey 怎么选? 聚合结果的类型和输入 value 类型不同、或需要显式零值时用 aggregateByKey;同类型简单归约用 reduceByKey,语义更直观。
  • DataFrame 里对应的优化是什么? Catalyst 会把 groupBy().agg() 规划成 HashAggregate,支持 partial aggregate(相当于 combine)和 spill;Spark 3.x 还有 AQE 的倾斜 join/聚合优化。写 SQL 时这套是自动的,这是 DataFrame API 比 RDD API 省心的典型例子。

评论 (0)

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

91学AI

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