TT Lab
开始
学习 学习路径 课程

Apache Spark — 慢作业的答案在执行计划和事件日志里

用四种策略执行同一个连接,并在执行计划中确认

在 TT Lab 中继续学习

目标

把订单和商品的连接分别用广播哈希、排序合并、shuffle 哈希三种策略运行,比较计划和结果,并通过事件日志确认 AQE 在运行中把排序合并改成了广播。还要做键重叠的表的连接让行数膨胀的实验,以及挑出另一侧没有配对的行的连接。

为什么重要

连接常常是 Spark 作业里最贵的运算,而它的成本是由策略决定的。如果一侧较小,就把较小的一侧复制给所有任务(广播),不对大表做 shuffle 就完成连接。如果两张都很大,就把两侧按键 shuffle 之后排序,边对齐边合并(排序合并)。shuffle 之后用一侧建哈希表来连接的 shuffle 哈希,跳过了排序,但要占用内存。 Spark 根据统计信息猜测大小,从而选择策略。判断为“小”的标准是 spark.sql.autoBroadcastJoinThreshold(默认 10MB)。猜错了,策略也会错——过滤之后实际上很小,却按原始大小来猜而选了排序合并;或者相反,广播了一张大表,结果驱动器内存溢出。AQE 会用 shuffle 结束之后的实际大小重新判断,修正这个问题。 连接结果的行数由键的重复情况决定。本以为某一侧的键是唯一的,结果不是,连接就会在没有任何错误的情况下把行成倍放大,后面的合计全都膨胀。

步骤

  1. 在 /root/spk/join/common.py 中放入读取四张表的函数和按类别销售额的函数,用 /root/spk/join/auto.py(应用 spk-join-auto,默认设置)把按类别销售额以带表头行的 CSV(category、revenue)写入 /root/spk/join/out/by_category。
  2. 用 /root/spk/join/smj.py(应用 spk-join-smj,spark.sql.autoBroadcastJoinThreshold=-1、spark.sql.adaptive.enabled=false)把同样的结果写入 /root/spk/join/out/by_category_smj。
  3. 在 /root/spk/join/hint.py(应用 spk-join-hint,与第 2 步相同的设置)中对商品一侧加上 broadcast 提示,写入 /root/spk/join/out/by_category_hint。
  4. 在 /root/spk/join/shash.py(应用 spk-join-shash,与第 2 步相同的设置)中对商品一侧加上 shuffle_hash 提示,写入 /root/spk/join/out/by_category_shash。
  5. 用 /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。
  6. 用 /root/spk/join/dup.py(应用 spk-join-dup),把已完成支付订单与促销表直接连接得到的行数,以及用 left_semi 连接得到的行数,以 {"naive": 정수, "semi": 정수}(占位符为整数)写入 /root/spk/join/out/dup.json。
  7. 用 /root/spk/join/anti.py(应用 spk-join-anti),把一次都没有下过订单的客户的 customer_id 以 CSV 写入 /root/spk/join/out/no_orders。
  8. 在 /root/spk/join/report.md 中以 ## 네 가지 전략(韩文,意为“四种策略”)、## AQE 의 전환(韩文,意为“AQE 的转换”)、## 키 중복(韩文,意为“键重复”)三节写成报告。第三节放入第 6 步的两个数字。

参考

较小的一侧会自动被广播

在 /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、把什么装进内存;第二节写最初的计划和最终计划有什么不同;第三节写行数膨胀了多少。