SortMergeJoin与ShuffleHashJoin内部机制及Shuffle阶段疑问
测试场景:从HDFS读取两个JSON文件生成各2个分区的DataFrame,执行内连接并指定使用SMJ,默认Shuffle分区数为200,任务分为4个阶段,不同阶段运行6个CodeGen生成的Java函数。
问题1
我理解Stage 2和Stage 3从数据源读取数据,通过CodeGen代码执行流水线转换,为每个输入分区生成并写入1个Shuffle数据文件(共2个),每个文件包含对应200个Shuffle分区的块,块按Shuffle分区ID排序,同时为每个数据文件写入1个记录块位置的索引文件。该理解是否正确?若正确,每个块内的数据是否在此阶段已按连接键排序?因我仅在Stage 4中看到带排序逻辑的CodeGen代码。
回答
你的核心理解是正确的:
- Stage 2、3作为Shuffle的Map阶段,每个输入分区对应生成1个Shuffle数据文件(两个输入分区对应2个数据文件),每个文件内包含对应200个Shuffle分区的数据块,块按Shuffle分区ID顺序排列,同时每个数据文件会生成1个索引文件记录各块的偏移量和长度。
- 每个块内的数据不会在此阶段按连接键排序。SMJ的排序逻辑确实是在Stage 4的Reduce阶段完成的:Map阶段仅负责根据连接键计算Shuffle分区ID,将数据写入对应分区的块中,不会做连接键排序;到了Stage 4,Reducer读取对应分区的数据后,才会执行连接键排序,之后再做合并连接。
问题2
Stage 4中有200个任务执行Shuffle读取,为Stage 2和Stage 3的每个mapPartitionRDD分别创建shuffledRowRDD。Stage 2和Stage 3的每个mapPartitionRDD各有200个待读取的块,那么Stage 4中ID为0的单个Reducer任务是否会并行读取Stage 2和Stage 3的Shuffle文件中对应分区0的块?若如此,它如何通过Shuffle ID(唯一标识Stage)区分来自不同Stage的块?
回答
- Stage 4中ID为0的Reducer任务会并行读取Stage 2和Stage 3对应分区0的块,Spark的Shuffle读取器支持同时从多个Shuffle源拉取数据。
- 区分不同Stage的块依靠以下两点:
- Shuffle ID:每个Shuffle阶段有唯一的Shuffle ID,Stage 2和Stage 3对应不同的Shuffle ID,Reducer任务在拉取数据时会指定对应的Shuffle ID,以此区分来自不同Stage的Shuffle数据。
- ShuffleDependency:每个ShuffledRowRDD关联对应的ShuffleDependency,其中包含Shuffle ID、分区器等信息,Spark通过这个依赖关系精准定位到对应Stage的Shuffle输出。
问题3
当使用ShuffleHashJoin(SHJ)替代SMJ时,执行阶段类似,但Stage 4执行哈希连接而非SortMergeJoin。这是否意味着两者的Shuffle阶段完全相同?我了解到ShuffleHashJoin会为每个输入分区的每个Shuffle分区生成1个Shuffle文件,那么在此测试场景下,Stage 2和Stage 3的每个mapOutputRDD是否会生成2*200=400个Shuffle文件?SHJ是否也会生成索引文件和数据文件?
回答
- 两者的Shuffle阶段不完全相同:
- SMJ要求两边数据都按连接键分区(且后续在Reduce阶段排序),所以Shuffle分区器是基于连接键的哈希分区器。
- SHJ通常会选择将小表广播或按连接键哈希分区,大表按相同连接键哈希分区;如果是两边都做Shuffle的场景,分区器也是基于连接键的哈希分区,但Shuffle阶段的输出逻辑略有不同——SMJ的Map阶段仅分区不排序,而SHJ的Map阶段也不会排序,但后续Reduce阶段是直接构建哈希表而非排序合并。
- 关于Shuffle文件数量:你的理解有误,SHJ的Map阶段同样是每个输入分区生成1个Shuffle数据文件,而非每个Shuffle分区生成1个。在测试场景下,Stage 2和Stage 3的每个Map任务(对应1个输入分区)生成1个Shuffle数据文件,总共2个输入分区对应2个数据文件,每个文件包含200个Shuffle分区的块,和SMJ的文件数量一致。
- SHJ同样会生成索引文件和数据文件:只要是基于Spark统一ShuffleManager的Shuffle输出,都会生成数据文件和对应的索引文件,用于Reducer定位到对应分区的块位置。
内容的提问来源于stack exchange,提问作者akhil pathirippilly

