1、 join 算法
join算法是 一般是reduce阶段执行时选择的join算法
1.1、BroadcastHashJoin
最快的join算法,只需map端匹配,无需shuffle过程
- 小表完整广播到所有executor,在excutor中,大表执行遍历,无需shuffle;
- 小表阈值有参数设置;spark.sql.autoBroadcastJoinThreshold 默认 10mb
1.2、ShuffleHashJoin (SHJ)
ShuffleHashJoin 需要把整张表按照shuffle 后分区均分后数据全部加载进内存建哈希表,另一张表去hash表中匹配
因此需要 表的大小数据量大小< spark.sql.autoBroadcastJoinThreshold * spark.sql.shuffle.partitions
默认值: < 10MB * 200
- 如果遇到数据倾斜hash表装不下,会溢写到磁盘,io极限导致溢写慢且超出executor的内存会oom
1.3、SortMergeJoin(SMJ)
spark3.2.1 默认算法,稳定
流式逐行读取、仅存两个指针行,不需要缓存整块数据,内存占用恒定极低,不会爆内存
分区排序 + 归并匹配,两段都不会持有大量数据
分区排序阶段
- 先读取一批数据放入内存缓冲区排序;
- 缓冲区满了就溢写到磁盘,分成多个有序小文件;
- 最后多路归并磁盘上的有序文件,输出整体有序分区的一个文件
归并匹配阶段
两张有序数据集,只用双指针流式读取
- 只在内存保留 A 表当前一行、B 表当前一行;
- 比较 join key,小的那条直接丢弃,读取下一行;
- 相同 key 就拼接输出,不需要缓存任何历史数据
2、shuffle算法
shuffle主要解决数据如何从上游task风阀到下游task的算法,常用算法有HashShuffle,SortShuffle,
spark.shuffle.manager=sort; 值有hash / sort / tungsten-sort
shuffle主要是按照一定的规则进行分发数据,map端一般不计算数据,但是也有参数可以在map进行局部计算
读取 spark.shuffle.manager
├─ hash → HashShuffleManager(淘汰)
├─ tungsten-sort
│ └─ 满足条件 → UnsafeShuffleWriter
│ └─ 不满足 → 普通 SortShuffleWriter
└─ sort(默认 SortShuffleManager)
├─ 判断:下游分区数 ≤ bypassMergeThreshold && 无 Map 端聚合(mapSideCombine == false)
│ └→ BypassSortShuffleWriter
└─ 其他情况
└→ 标准 SortShuffleWriter
2.1、HashShuffleManager算法(已淘汰)
每个 Map Task,根据key进行hash分区,分区后为每一个下游 Reduce Task单独生成一个文件。假设有 M 个 MapTask,R 个 ReduceTask,输出总文件数=M*R 个,容易导致服务器层面文件句柄不够用问题
# 比如一个有20map, 20reduce
# 生成map端shuffle 生成文件有
map1 执行生成的文件part-0, part-1 … part-19 20个
map2 执行生成的文件part-0, part-1 … part-19 20个
.
.
.
map19 执行生成的文件part-0, part-1 … part-19 20个
总共400个
优化参数:spark.shuffle.consolidateFiles=true,但是也只是优化了map阶段,并没有优化reduce端,如果有200个reduce,仍然有很大概率出现文件句柄问题
2.2、SortShuffle(兜底的shuffle算法)
内存充足
- map读取数据,根据关键key的hash值进行分区,存放到PartitionedAppendOnlyMap中
- 内存充足,缓存大小默认5mb,写满申请内存,申请成功,一直往对象中写,完成后写入落磁盘生成shuffle_xxx.data, shuffle_xx.index两个文件
内存不足时
- map读取数据,根据关键key的hash值进行分区,存放到PartitionedAppendOnlyMap中,放在内存在
- 写满5mb,申请内存,申请不到开启溢写磁盘,溢写文件大小是缓存中的大小;
- 溢写时先排序,便于合并时保障最后再合并文件shuffle_xxx.data是有序的
- 分批写入,溢写的临时文件都是有序的
- 溢写完成后,清空内存数据结构,继续接收后续数据
- 合并生成shuffle_xxx.data,shuffle_xxx.index文件
Map 端接收数据
↓
写入内存结构 (AppendOnlyMap 或 PairBuffer)
↓ (内存满了)
[溢写循环] 按分区(和Key)排序 -> 写入磁盘临时文件 (spill-1, spill-2...)
↓ (Map 任务结束)
内存剩余数据 + 所有 spill 文件 进行多路归并 (Merge)
↓ (有聚合则在此处进行最终 Combine)
按分区 ID 顺序写入 最终 data 文件
↓
生成配套的 index 索引文件
↓
Map 端 Shuffle 完成,等待 Reduce 端拉取
2.3、SortShuffle的BypassSortShuffleWriter优化
走BypassSortShuffleWriter 优化的条件
无 Map 端聚合(mapSideCombine == false)
即 ShuffleDependency 中没有指定聚合器(Aggregator)。这意味着像 reduceByKey、aggregateByKey 这类带有 Map 端预聚合的操作不会走 Bypass 机制;而 groupByKey、join 等操作则可能触发。Reduce 端分区数小于等于阈值
dep.partitioner.numPartitions <= spark.shuffle.sort.bypassMergeThreshold
该阈值默认值为 200。你可以通过 spark.sql.shuffle.partitions 或 spark.default.parallelism 调整分区数,或通过 spark.shuffle.sort.bypassMergeThreshold 参数调整这个阈值。
优化:不需要经过派去,数据写入是,直接根据hash 算出分区id,写入不同独立临时文件中,最后将这些文件合并成一个文件,并生产索引
缺点:map任务需为每个分区维护一个独立文件写入器和缓冲区,如果分区过大,会打开大量文件句柄和内存缓冲区,因此默认限制200
2.4、Tungsten-Sort
在它诞生之前,标准的 SortShuffleWriter 有两个致命的痛点:大内存时GC(垃圾回收)停顿
与标准的SortShuffle 逻辑一样,只是采用堆外内存,不存在gc问题
内存不足溢写,充足则不溢写;最终生成有序的shuffle_xx.data + shuffle_xx.index两个文件