Spark 有哪些join方式
Spark Join 可以从两个维度理解:
- SQL 语义类型:Inner、Left、Right、Full、Semi、Anti、Cross Join
- 物理执行策略:Broadcast Hash Join、Sort Merge Join、Shuffled Hash Join 等
实际性能调优时,重点看第二类。
一、按 SQL 语义分类
1. Inner Join:内连接
只返回两边 Join Key 能匹配的数据。
df1.join(df2, "user_id", "inner")
SELECT *
FROM orders o
JOIN users u
ON o.user_id = u.user_id;
2. Left Outer Join:左外连接
保留左表全部记录;右表无匹配时填 NULL。
df1.join(df2, "user_id", "left")
SELECT *
FROM orders o
LEFT JOIN users u
ON o.user_id = u.user_id;
常用于:
订单表 LEFT JOIN 用户维表
3. Right Outer Join:右外连接
保留右表全部记录;左表无匹配时填 NULL。
df1.join(df2, "user_id", "right")
实践中较少单独使用,通常可通过交换左右表的位置改写为 Left Join。
4. Full Outer Join:全外连接
保留左右两边的全部数据,匹配不到的一侧补 NULL。
df1.join(df2, "user_id", "full")
SELECT *
FROM a
FULL OUTER JOIN b
ON a.id = b.id;
常用于:
- 两个系统的数据对账;
- 查找只存在于某一侧的数据;
- CDC 数据核对。
5. Left Semi Join:左半连接
只返回左表中“在右表存在匹配”的记录,不返回右表字段。
SELECT o.*
FROM orders o
LEFT SEMI JOIN users u
ON o.user_id = u.user_id;
可以理解为:
SELECT *
FROM orders
WHERE user_id IN (SELECT user_id FROM users);
适合做“存在性过滤”。
6. Left Anti Join:左反连接
只返回左表中“在右表找不到匹配”的记录。
SELECT o.*
FROM orders o
LEFT ANTI JOIN paid_orders p
ON o.order_id = p.order_id;
等价语义接近:
SELECT *
FROM orders o
WHERE NOT EXISTS (
SELECT 1
FROM paid_orders p
WHERE o.order_id = p.order_id
);
典型场景:
- 找没有支付的订单;
- 找维表中不存在的事实数据;
- 找增量数据中未处理过的记录;
- 数据对账、脏数据排查。
7. Cross Join:笛卡尔积
左表每一行与右表每一行组合。
df1.crossJoin(df2)
SELECT *
FROM a
CROSS JOIN b;
如果:
A 表 1 万行
B 表 10 万行
结果可能达到:
1 万 × 10 万 = 10 亿行
所以一般应避免。只有在右表非常小、确实需要全组合时才使用。
有些 Spark 配置中默认禁止隐式笛卡尔积,需要显式开启或写 CROSS JOIN:
spark.sql.crossJoin.enabled
二、按 Spark 物理执行策略分类
Spark 会根据数据大小、Join 条件、统计信息、配置及 AQE,选择实际的 Join 算法。
可以通过:
df.explain("formatted")
或 Spark UI 的 SQL 页面查看最终使用的 Join。
1. Broadcast Hash Join(BHJ,广播哈希 Join)
将较小表广播到每个 Executor,在每个 Executor 本地构建 Hash 表,再扫描大表进行匹配。
小表 -> 广播到所有 Executor
大表 -> 不 Shuffle 或尽量少 Shuffle
示意:
小表
↓ 广播
Executor 1: 大表分区 1 + 小表
Executor 2: 大表分区 2 + 小表
Executor 3: 大表分区 3 + 小表
示例
from pyspark.sql.functions import broadcast
result = orders.join(
broadcast(dim_user),
"user_id",
"left"
)
适用场景
大事实表 JOIN 小维表
例如:
订单明细表:数百 GB
省份维表:几 MB
商品分类维表:几十 MB
优点
- 通常不需要对大表 Shuffle;
- 常常是性能最好的 Join 策略;
- 减少网络传输和磁盘 I/O。
风险
小表必须能被广播到 Executor 内存中。表太大时可能造成:
- Executor OOM;
- 广播超时;
- GC 压力;
- 任务不稳定。
自动广播阈值常由以下配置控制:
spark.sql.autoBroadcastJoinThreshold
常见默认值是 10MB,但不同 Spark 版本和平台可能不同。
spark.conf.set(
"spark.sql.autoBroadcastJoinThreshold",
50 * 1024 * 1024
)
不要只看文件大小,还应考虑反序列化后对象大小和 Executor 可用内存。
2. Sort Merge Join(SMJ,排序合并 Join)
Spark 中大表等值 Join 最常见的策略之一。
执行流程:
左表按 Join Key Shuffle + Sort
右表按 Join Key Shuffle + Sort
↓
按排序后的 Key 做 Merge Join
示意:
left table -> Exchange -> Sort
right table -> Exchange -> Sort
↓
SortMergeJoin
示例
result = fact_orders.join(fact_payments, "order_id")
适用场景
- 两张大表 Join;
- 等值 Join;
- 两边都不适合广播;
- Join Key 可排序。
优点
- 能处理大规模数据;
- 内存压力通常比 Hash Join 更容易控制;
- 可进行外部排序并支持 Spill;
- 是大表 Join 的稳定方案。
缺点
- 通常两侧都有 Shuffle;
- 有排序开销;
- 网络和磁盘 I/O 较大;
- 数据倾斜会非常明显。
Spark 常见配置:
spark.sql.join.preferSortMergeJoin
通常默认偏向 Sort Merge Join。
3. Shuffled Hash Join(SHJ,Shuffle Hash Join)
左右表先按 Join Key Shuffle,使相同 Key 到同一个分区;每个分区内选择较小的一侧构建 Hash 表,另一侧进行探测。
左表 -> Shuffle
右表 -> Shuffle
↓
每个分区内:小分区建 Hash 表,大分区查找
适用场景
- 等值 Join;
- 两边都无法广播;
- 每个分区内的小表可以放进内存;
- 不希望承担完整排序的开销;
- 一边数据量明显小于另一边。
优点
- 相比 SMJ,可能省去排序;
- 特定场景下性能较好。
缺点
- 对内存要求更高;
- 如果分区数据过大或存在倾斜,容易 OOM;
- 大规模稳定生产场景中往往不如 SMJ 常见。
是否选用通常由 Spark 优化器、统计信息和 AQE 决定。
4. Broadcast Nested Loop Join(BNLJ,广播嵌套循环 Join)
当 Join 条件不是等值条件时,往往无法使用 Hash Join 或 Sort Merge Join。
如果其中一张表很小,Spark 可以广播小表,然后对大表的每条记录遍历小表记录,判断 Join 条件。
例如范围 Join:
SELECT *
FROM orders o
JOIN discount_rules d
ON o.amount BETWEEN d.min_amount AND d.max_amount;
特点
大表每条记录 × 广播小表的每条记录
若小表有 N 行、大表有 M 行,复杂度接近:
O(M × N)
适用场景
- 非等值 Join;
- 小表极小;
- 范围匹配;
- 不等于、大小比较、复杂条件关联。
注意
即使小表能广播,BNLJ 也可能非常慢。它不是普通 Broadcast Hash Join。
5. Cartesian Product(笛卡尔积执行)
没有有效 Join 条件时,Spark 可能执行笛卡尔积。
SELECT *
FROM a
CROSS JOIN b;
物理计划可能出现:
CartesianProduct
其复杂度接近:
O(M × N)
除非一边数据极小,否则应避免。
三、Bucket Join(分桶 Join)
Bucket Join 更准确说是一种减少或避免 Shuffle 的数据组织优化,不是独立的通用 Join 算法。
如果两张表:
- 都按相同 Join Key 分桶;
- 桶数兼容;
- 使用相同或兼容的分桶规则;
- Spark 可以识别并利用桶信息;
那么 Join 时可能避免部分甚至全部 Shuffle。
示例建表:
CREATE TABLE orders_bucketed (
order_id STRING,
user_id STRING,
amount DECIMAL(10, 2)
)
USING parquet
CLUSTERED BY (user_id) INTO 128 BUCKETS;
另一张表也按 user_id 进行兼容分桶:
CREATE TABLE users_bucketed (
user_id STRING,
user_name STRING
)
USING parquet
CLUSTERED BY (user_id) INTO 128 BUCKETS;
注意
现实中 Bucket Join 的使用效果依赖:
- Spark 版本;
- Hive Metastore / Catalog;
- 文件格式;
- 表是否经过正确写入;
- 是否启用了相关优化;
- 表是否被破坏分桶布局。
因此不能假设“做了分桶就一定没有 Shuffle”,要以 EXPLAIN 结果为准。
四、AQE 会动态调整 Join 策略
Spark 3 中的 AQE(Adaptive Query Execution,自适应查询执行)可以根据运行时统计信息动态改写 Join 策略。
例如原计划是:
Sort Merge Join
运行时发现一侧 Shuffle 后只有几 MB,AQE 可能改为:
Broadcast Hash Join
常见配置:
spark.conf.set("spark.sql.adaptive.enabled", "true")
AQE 还可能进行:
- 动态合并 Shuffle 分区;
- 自动处理部分倾斜 Join;
- 将 Sort Merge Join 转换为 Broadcast Hash Join;
- 将 Sort Merge Join 转换为 Shuffled Hash Join。
因此最终应在 Spark UI 的实际执行计划中确认 Join 类型,而不是只看初始 explain()。
五、如何选 Join 方式?
| 场景 | 优先策略 |
|---|---|
| 大表 Join 小维表,小表可放进每个 Executor 内存 | Broadcast Hash Join |
| 两张大表按等值 Key 关联 | Sort Merge Join |
| 等值 Join,单分区内一侧明显较小且能放入内存 | Shuffled Hash Join |
| 非等值 / 范围 Join,且一侧很小 | Broadcast Nested Loop Join |
| 无关联条件、必须全量组合 | Cartesian / Cross Join,谨慎 |
| 两张表按相同 Key 做了兼容分桶 | Bucket Join 优化,查看是否减少 Shuffle |
六、如何看 Spark 实际使用了哪种 Join?
result.explain("formatted")
重点关注物理计划中的关键词:
BroadcastHashJoin
SortMergeJoin
ShuffledHashJoin
BroadcastNestedLoopJoin
CartesianProduct
如果看到:
Exchange hashpartitioning(...)
表示此处存在 Shuffle。
例如:
+- SortMergeJoin [user_id], [user_id], Inner
:- Sort [user_id ASC]
: +- Exchange hashpartitioning(user_id, 200)
+- Sort [user_id ASC]
+- Exchange hashpartitioning(user_id, 200)
这就是典型的 Sort Merge Join:两边按 user_id Shuffle,再排序后合并。
总结
Spark Join 的 SQL 语义包括:
Inner Join
Left / Right / Full Outer Join
Left Semi Join
Left Anti Join
Cross Join
Spark 常见物理 Join 策略包括:
1. Broadcast Hash Join(小表广播,优先考虑)
2. Sort Merge Join(两张大表等值 Join 的常见方案)
3. Shuffled Hash Join(分区内小表建哈希)
4. Broadcast Nested Loop Join(非等值 Join + 小表广播)
5. Cartesian Product(笛卡尔积,应谨慎)
性能调优的核心通常是:
尽可能使用 Broadcast Hash Join;无法广播的大表等值 Join 通常使用 Sort Merge Join;同时重点处理 Shuffle、数据倾斜和不合理分区数。
更多推荐
所有评论(0)