精选·Spark

Spark Join 有哪几种实现?生产上怎么选型?

91学AI·2026/7/13·7 阅读

考察点

join 是 Spark 里最贵的操作之一,这题考察的是你对执行引擎的理解而不是 API 记忆。面试官想听:三种 join 各自的执行流程和代价模型(谁要 shuffle、谁要排序、内存放不放得下);为什么 Sort Merge Join 是默认兜底;小表广播的阈值和实际收益;以及遇到 join 慢的时候怎么从 explain 判断当前走了哪种策略、能不能干预。追问常往 broadcast 的坑、AQE 动态切换 join、大表 join 大表的倾斜处理方向走。

参考答案

三种实现的总览

Spark SQL 的 join 物理策略主要三种,按代价从低到高排:

策略是否 shuffle是否排序前提条件
Broadcast Hash Join (BHJ)不 shuffle 大表小表能放进单 executor 内存
Shuffled Hash Join (SHJ)双方 shufflebuild 侧每个分区的 hash 表放得下
Sort Merge Join (SMJ)双方 shuffle无,兜底策略

下面逐个说执行流程。

Broadcast Hash Join

小表整体收集到 driver,再广播到每个 executor,在内存里建成 hash 表(build side);大表不 shuffle、不排序,每个 task 流式读自己分区的数据(probe side),逐条去 hash 表里探测匹配。

代价:广播小表的 IO(但广播用的是 BitTorrent 式的 TorrentBroadcast,比点对点快),大表一次顺序扫描。省掉了大表的 shuffle 和排序,这是它快的原因。适用场景就是典型的大小表 join:维度表(几万到几百万行)join 事实表(几十亿行)。

阈值由 spark.sql.autoBroadcastJoinThreshold 控制,默认 10MB——统计的是表的物理大小(Spark SQL 依据元数据或逻辑计划估算),不是行数。生产上调到 100MB~500MB 很常见,但要想清楚:广播变量会进每个 executor 的内存,大广播 + 高并发 task 可能挤占执行内存。

Shuffled Hash Join

两张表都按 join key 做 shuffle(保证同 key 落在同一分区),然后每个分区内把 build side 的数据建 hash 表,probe side 流式探测。

它比 SMJ 省掉了排序,但要求 build side 单个分区的数据能放进内存建 hash 表——注意是按分区说的,不是全表。Spark 选择 SHJ 的条件比较保守(要求平均分区大小小于 autoBroadcastJoinThreshold 的一定倍数,且优先选小的一侧做 build),所以它出场频率不高,但在「两表都大于广播阈值、但分区后单分区 build side 可控」时比 SMJ 划算。

Sort Merge Join

默认兜底策略,三段式:

  1. shuffle 阶段:两表按 join key 分区重分布,同 key 到同一分区;
  2. sort 阶段:每个分区内两侧数据各自按 key 排序(外部排序,可 spill);
  3. merge 阶段:双指针归并,像归并排序的 merge 步,O(n) 扫完。

代价是全量 shuffle + 全量排序,但内存安全:排序可 spill、merge 是流式的,两张任意大的表都能 join 完,这是它当兜底的原因。数据倾斜时它的弱点暴露——热点 key 全压在一个 task 上,排序和归并都卡在那一个分区。

Catalyst 怎么选,以及怎么干预

Catalyst 的选择逻辑大致是:先看 join hint(broadcast / merge / shuffle_hash),没 hint 时检查一侧是否 ≤ autoBroadcastJoinThreshold 且有统计信息支持,能广播就广播;否则看能不能 SHJ;都不行就 SMJ。

干预手段:

-- 强制广播(Spark 3.x hint 写法)
SELECT /*+ BROADCAST(dim) */ * FROM fact f JOIN dim ON f.id = dim.id

-- 强制 SMJ(大表 join 大表时防止误广播)
SELECT /*+ MERGE(f, d) */ ...

RDD 层面没有自动策略,自己掌握:小表 join 大表用 sc.broadcast() 手动广播 + map 端探测;两个大表只能 cogroup/join 走 shuffle。

选型口诀与排查方法

  • 小表(< 广播阈值)join 任意大表:BHJ,收益最大。
  • 中表 join 中表,单分区可控:SHJ 可能比 SMJ 快 20-30%(省排序)。
  • 大表 join 大表、或有倾斜风险:SMJ 兜底,倾斜时配合倾斜处理(加盐、AQE skew join 优化)。
  • 不确定走了哪种:df.explain(true) 或 Spark UI 的 SQL tab 看物理计划,找 BroadcastHashJoin / SortMergeJoin / ShuffledHashJoin 字样,BroadcastExchange 节点确认广播发生。

一个常见误判:表的真实大小超过阈值(统计信息过期或未 ANALYZE TABLE),该广播的没广播,走了 SMJ 慢十倍。解法:ANALYZE TABLE 刷新统计,或用 hint 强制。

可能的追问

  • Broadcast 的实现细节,会不会压垮 driver? collect 到 driver 时确实占 driver 内存,超过 spark.driver.maxResultSize 会失败;广播分发走 TorrentBroadcast,把数据切块后 executor 之间互相拉取,driver 只发种子,压力可控。
  • AQE 对 join 有什么帮助? AQE 拿到 shuffle 后的真实统计,能把原计划中的 SMJ 动态降级成 BHJ(发现一侧实际很小),还能自动拆分倾斜的 join 分区(OptimizeSkewedJoin),3.x 之后效果明显。
  • 大表 join 大表且一侧有热点 key 怎么办? 经典方案是热点 key 加盐打散:热点侧 key 加随机前缀(0~N),另一侧对应膨胀 N 份,join 完再去盐。或者 Spark 3.x 直接开 spark.sql.adaptive.skewJoin.enabled 让 AQE 处理。

评论 (0)

暂无评论,快来抢沙发吧!

91学AI

© 2026 91学AI · 按岗位学 AI 与大数据. All rights reserved.