You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark AQE中SMJ转BHJ时Local Shuffle Reader相关疑问

Spark AQE中SMJ转BHJ后Local Shuffle Reader的疑问与理解

我的理解

假设将两个文件读入DataFrame,执行过滤后开启AQE进行Join操作。基于自动广播阈值和ShuffleMapStage生成的shuffle文件大小,Spark将SortMergeJoin(SMJ)转换为BroadcastHashJoin(BHJ),流程如下:

  1. AQE提交初始执行计划→DAGScheduler提交2个ShuffleMapStage→两个executor执行shuffle任务并将shuffle文件写入各自磁盘→状态发送至mapOutputTracker;
  2. AQE重新规划为BHJ的物理执行计划并重新提交→DAGScheduler提交缺失的stage;
  3. 由于ShuffleMapStage已完成且shuffle文件可用,跳过这些stage→读取shuffle文件以获取待广播的Join关联数据,构建RDD并收集到Driver后广播至持有另一数据集分区的executor;
  4. 在广播数据上构建哈希表,执行经典哈希Join。

附Spark UI SQL标签页执行计划图:
SQL执行计划图1
SQL执行计划图2

我的疑问

  1. 上述理解是否正确?
  2. 我看到左右关联均使用了CustomShuffleReader (local),原因是什么?CustomShuffleReader (local)是否意味着会从ShuffleMapStage执行所在executor的本地存储读取数据?
  3. 是否是这样:若禁用spark.sql.adaptive.localShuffleReader.enabled配置(对应配置截图:配置项截图),待广播数据的shuffle文件会通过sortShuffleManager常规的shuffleClient从reducer端executor读取;若启用该配置,则会使用CustomShuffleReader从shuffle文件写入节点读取,而非通过reducer ID从reducer端读取?(对应执行计划截图:执行计划截图)

解答

1. 整体理解的正确性

你的流程描述基本正确,但有一处关键细节需要修正:待广播的数据不需要先收集到Driver再广播。AQE将SMJ转为BHJ时,广播的是shuffle后的小数据集,Spark会直接在executor端读取本地shuffle文件并构建哈希表,或者将小数据集的所有shuffle文件数据拉取到当前executor后完成广播逻辑,全程无需经过Driver中转,这样能避免Driver成为性能瓶颈。

其余步骤的逻辑完全准确:AQE先提交初始SMJ计划并执行ShuffleMapStage,待获取shuffle文件的实际大小后触发计划重写,转为BHJ,跳过已完成的ShuffleMapStage,直接利用已生成的shuffle文件执行哈希Join。

2. CustomShuffleReader (local)的作用与含义

  • 含义:CustomShuffleReader (local)确实代表从ShuffleMapStage执行所在的executor本地存储读取shuffle文件,而非通过常规shuffle拉取方式从远程节点读取。
  • 出现原因:初始SMJ计划中,两个数据集都按Join键做了哈希分区;转为BHJ后,这种分区逻辑不再需要,此时需要读取小数据集的全部shuffle文件。Local Shuffle Reader会让executor优先读取本地磁盘上的shuffle文件,本地没有的部分才从远程节点拉取,能大幅减少跨节点数据传输的开销,提升读取效率。

左右关联都出现这个Reader,是因为两个数据集在初始计划中都执行了shuffle;转为BHJ后,不管是作为广播端(读取全部数据用于广播)还是被匹配端(读取自身分区数据做哈希匹配),都需要读取各自的shuffle文件,因此都会用Local Shuffle Reader来优化读取性能。

3. 配置开关对shuffle读取方式的影响

你的理解完全正确:

  • 禁用spark.sql.adaptive.localShuffleReader.enabled时,Spark会采用常规shuffle读取逻辑:按reducer ID拉取对应分区的shuffle数据,数据会从map端节点传输到reducer端节点,存在跨节点网络开销。
  • 启用该配置时,CustomShuffleReader会直接从shuffle文件的写入节点(即执行ShuffleMapTask的executor)读取数据,优先读本地文件,远程文件直接从原节点拉取,跳过了按reducer ID分区的逻辑,避免了不必要的数据传输,在转为BHJ的场景下能最大化利用本地存储,提升整体效率。

内容的提问来源于stack exchange,提问作者akhil pathirippilly

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 16:55:21