Apache Spark — 慢作业的答案在执行计划和事件日志里
用四种策略执行同一个连接,并在执行计划中确认
目标
把订单和商品的连接分别用广播哈希、排序合并、shuffle 哈希三种策略运行,比较计划和结果,并通过事件日志确认 AQE 在运行中把排序合并改成了广播。还要做键重叠的表的连接让行数膨胀的实验,以及挑出另一侧没有配对的行的连接。
为什么重要
连接常常是 Spark 作业里最贵的运算,而它的成本是由策略决定的。如果一侧较小,就把较小的一侧复制给所有任务(广播),不对大表做 shuffle 就完成连接。如果两张都很大,就把两侧按键 shuffle 之后排序,边对齐边合并(排序合并)。shuffle 之后用一侧建哈希表来连接的 shuffle 哈希,跳过了排序,但要占用内存。
Spark 根据统计信息猜测大小,从而选择策略。判断为“小”的标准是 spark.sql.autoBroadcastJoinThreshold(默认 10MB)。猜错了,策略也会错——过滤之后实际上很小,却按原始大小来猜而选了排序合并;或者相反,广播了一张大表,结果驱动器内存溢出。AQE 会用 shuffle 结束之后的实际大小重新判断,修正这个问题。
连接结果的行数由键的重复情况决定。本以为某一侧的键是唯一的,结果不是,连接就会在没有任何错误的情况下把行成倍放大,后面的合计全都膨胀。
步骤
- 在 /root/spk/join/common.py 中放入读取四张表的函数和按类别销售额的函数,用 /root/spk/join/auto.py(应用
spk-join-auto,默认设置)把按类别销售额以带表头行的 CSV(category、revenue)写入 /root/spk/join/out/by_category。 - 用 /root/spk/join/smj.py(应用
spk-join-smj,spark.sql.autoBroadcastJoinThreshold=-1、spark.sql.adaptive.enabled=false)把同样的结果写入 /root/spk/join/out/by_category_smj。 - 在 /root/spk/join/hint.py(应用
spk-join-hint,与第 2 步相同的设置)中对商品一侧加上broadcast提示,写入 /root/spk/join/out/by_category_hint。 - 在 /root/spk/join/shash.py(应用
spk-join-shash,与第 2 步相同的设置)中对商品一侧加上shuffle_hash提示,写入 /root/spk/join/out/by_category_shash。 - 用 /root/spk/join/aqe.py(应用
spk-join-aqe,spark.sql.autoBroadcastJoinThreshold=100k,AQE 开启),把已完成支付的订单与tier == 'vip'的客户连接,把按城市统计的订单数以 CSV(city、orders)写入 /root/spk/join/out/vip_by_city。 - 用 /root/spk/join/dup.py(应用
spk-join-dup),把已完成支付订单与促销表直接连接得到的行数,以及用left_semi连接得到的行数,以{"naive": 정수, "semi": 정수}(占位符为整数)写入 /root/spk/join/out/dup.json。 - 用 /root/spk/join/anti.py(应用
spk-join-anti),把一次都没有下过订单的客户的customer_id以 CSV 写入 /root/spk/join/out/no_orders。 - 在 /root/spk/join/report.md 中以
## 네 가지 전략(韩文,意为“四种策略”)、## AQE 의 전환(韩文,意为“AQE 的转换”)、## 키 중복(韩文,意为“键重复”)三节写成报告。第三节放入第 6 步的两个数字。
参考
- 原始数据:订单
/data/shop/orders.csv,商品/data/shop/products.csv(400 行),客户/data/shop/customers.csv(2 万行),促销/data/shop/promos.csv(product_id, promo_code——有几个商品有两个代码)。 - 策略通过计划中的算子名来确认:
BroadcastHashJoin、SortMergeJoin、ShuffledHashJoin、BroadcastHashJoin … LeftSemi、… LeftAnti。AQE 运行之后的计划在== Final Plan ==一节中。 - 提示用
F.broadcast(df),或df.hint("broadcast")、df.hint("shuffle_hash")、df.hint("merge")来加。如果是 SQL,就是/*+ BROADCAST(p) */。 - 常见错误:在第 2–4 步中开着 AQE,让计划在运行中被改变;把像促销表这样键重叠的表,当成维度表。
- 官方文档:Join Strategy Hints · Hints · JOIN · Converting sort-merge join to broadcast join
较小的一侧会自动被广播
在 /root/spk/join/common.py 中放入按 schema 读取四张表(订单、商品、客户、促销)的函数,以及按类别销售额(已完成支付订单 × 商品,category、revenue=qty×price 之和)的函数,再以应用名 spk-join-auto(设置保持默认值)创建 /root/spk/join/auto.py,把结果以带表头行的 CSV 写入 /root/spk/join/out/by_category。
商品表只有几 KB,远小于阈值(10MB)。Spark 会把它复制给所有任务,订单一侧不做 shuffle。评分器会检查事件日志的计划里是否有 BroadcastHashJoin,以及按类别销售额是否与原始数据一致。
关掉阈值,就是排序合并
以应用名 spk-join-smj、设置 spark.sql.autoBroadcastJoinThreshold=-1、spark.sql.adaptive.enabled=false 创建 /root/spk/join/smj.py,把与第 1 步相同的结果写入 /root/spk/join/out/by_category_smj。
阈值 -1 的意思是“不做自动广播”。现在订单和商品两侧都按 product_id 做 shuffle 并排序,然后合并。结果必须与第 1 步一行都不差。之所以关掉 AQE,是为了不让策略在运行中发生改变。
用提示强制广播
以应用名 spk-join-hint、与第 2 步相同的设置(阈值 -1,关闭 AQE)创建 /root/spk/join/hint.py,但要把商品一侧用 F.broadcast(p) 包起来,写入 /root/spk/join/out/by_category_hint。
提示优先于统计信息。即使把阈值关掉,只要有提示,也会广播。反过来说,随手加在大表上的广播提示,会原样占用驱动器和所有执行器的内存。
shuffle 哈希连接——用跳过排序来换取
以应用名 spk-join-shash、与第 2 步相同的设置创建 /root/spk/join/shash.py,但要对商品一侧加上 p.hint("shuffle_hash"),写入 /root/spk/join/out/by_category_shash。
shuffle 哈希连接把两侧按键做 shuffle 之后,用较小一侧的分区建哈希表,再让大表一侧流过去。没有排序,所以可能更快,但哈希表必须装得进内存。确认一下计划里的 Sort 算子消失了。
AQE 在运行中改变策略
以应用名 spk-join-aqe、设置 spark.sql.autoBroadcastJoinThreshold=100k(保持 AQE 开启)创建 /root/spk/join/aqe.py,把已完成支付订单与 tier == 'vip' 的客户按 customer_id 连接,把按城市统计的订单数以带表头行的 CSV(city、orders)写入 /root/spk/join/out/vip_by_city。
客户表文件大于 100KB,所以最初的计划是排序合并。但按 vip 过滤之后的实际大小只有几十 KB。AQE 会在 shuffle 的 Map 阶段结束之后,根据这个大小把剩余的计划改成广播。评分器会比较同一次运行的最初计划和最终计划。
键重叠,连接就会让行数膨胀
以应用名 spk-join-dup 创建 /root/spk/join/dup.py,把已完成支付订单与促销表(/data/shop/promos.csv)按 product_id 直接连接得到的行数,以及用 left_semi 连接得到的行数,以 {"naive": 정수, "semi": 정수}(占位符为整数)写入 /root/spk/join/out/dup.json。
促销表里有代码有两个的商品。这样的商品的订单,直接连接就会变成两行。left_semi 只看“是否有配对”,所以不会增加左侧的行。想一想,按是否有促销来拆分销售额时,该用哪一种。
挑出没有配对的一侧——left_anti
以应用名 spk-join-anti 创建 /root/spk/join/anti.py,把一次都没有下过订单(与状态无关)的客户的 customer_id,以带表头行的 CSV 写入 /root/spk/join/out/no_orders。
用 not in 子查询,或者 left join 之后过滤 null,也能做到,但 left_anti 的含义一目了然,而且在键里混有 null 时也不会让人糊涂。看看计划里连接类型是否显示为 LeftAnti。
留下哪种策略在什么时候合适
在 /root/spk/join/report.md 中写出 ## 네 가지 전략(韩文,意为“四种策略”)、## AQE 의 전환(韩文,意为“AQE 的转换”)、## 키 중복(韩文,意为“键重复”)三节。第一节放入第 1–4 步中看到的连接算子名,第三节放入第 6 步的两个数字。
第一节写每种策略对什么做 shuffle、把什么装进内存;第二节写最初的计划和最终计划有什么不同;第三节写行数膨胀了多少。