Spark Hash Shuffle中Reducer数量计算逻辑及相关疑问
关于Spark Hash Shuffle的核心问题解答
1. "reduce侧的分区数"具体指什么?
这里的"reduce侧的分区数",就是Shuffle之后下游Stage(即处理Shuffle输出数据的Stage)的任务并行度——每个分区对应一个Reducer任务。简单说就是Shuffle完成后,数据被拆分成了多少份,每份由一个独立的Reducer去处理。
2. 它如何获知reduce侧的分区数?
Spark在构建作业DAG时,会自动解析依赖关系:当遇到groupByKey、join这类需要Shuffle的操作时,作业会被划分为Map Stage(上游)和Reduce Stage(下游)。Reduce Stage的分区数就是Hash Shuffle要使用的Reducer数量,这个信息由Spark调度器在DAG解析阶段确定,直接传递给Map Stage的任务。
3. 是否等于spark.sql.shuffle.partitions?
分两种场景:
- Spark SQL场景(用DataFrame/Dataset执行聚合、关联操作):默认等于
spark.sql.shuffle.partitions的配置值(默认200); - RDD场景:取决于你调用Shuffle操作时指定的分区数(比如
rdd.groupByKey(100)里的100),如果没指定,则使用spark.default.parallelism的配置值(默认和集群CPU核数相关)。
4. 所谓的“计算”又指什么?
这里的“计算”是Hash Shuffle的核心分配逻辑——通过哈希算法把Map端的每条数据映射到对应的Reducer分区:
- 对每条数据的Key计算哈希值;
- 用哈希值对Reducer数量(下游分区数)取模,得到该数据要发送的Reducer编号;
- 将这条数据写入对应编号的输出文件。
示例说明
场景1:RDD API(指定分区数)
// 构建一个有2个分区的RDD(对应2个Mapper任务) val ordersRDD = sc.parallelize(Seq( ("user1", 100.0), ("user2", 200.0), ("user1", 150.0), ("user3", 300.0), ("user2", 50.0), ("user3", 100.0) ), 2) // 指定下游分区数为3(即Reducer数量=3) val groupedRDD = ordersRDD.groupByKey(3)
此时Hash Shuffle的工作流程:
- 2个Mapper任务,每个都要为3个Reducer各创建1个输出文件,总文件数=2*3=6;
- 对每个用户ID计算哈希取模:比如
hash("user1")%3=0、hash("user2")%3=1、hash("user3")%3=2,对应数据会写入各自编号的Reducer文件; - 下游3个Reducer分别读取对应编号的文件,完成分组求和。
场景2:Spark SQL场景(默认配置)
import spark.implicits._ val ordersDF = Seq( ("user1", 100.0), ("user2", 200.0), ("user1", 150.0), ("user3", 300.0), ("user2", 50.0), ("user3", 100.0) ).toDF("user_id", "amount") // 默认使用spark.sql.shuffle.partitions=200作为Reducer数量 val resultDF = ordersDF.groupBy("user_id").sum("amount")
此时每个Mapper任务会生成200个输出文件,如果Mapper数量是5,总文件数就是5*200=1000。
内容的提问来源于stack exchange,提问作者figs_and_nuts
相关产品推荐
相关产品推荐

