spark shuffle

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两个文件

©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容