Apache Spark — 慢作业的答案在执行计划和事件日志里
连接的配对方式能让代价相差十倍
一句话总结
Spark 可以用三种方式(广播哈希、排序合并、shuffle 哈希)完成同一个连接,选哪一种取决于其中一侧有多小,而这个判断会发生两次:计划时的估算,和运行中的实测(AQE)。
为什么要了解连接策略
连接是把两张表里键相同的行汇集到同一个地方。问题在于,这些行一开始分散在不同的分区、不同的执行器上。要配对,就必须有人移动。移动什么、移动多少,就是连接策略的全部。
假设要把 1,000 万条订单和 300 个商品连接。如果把订单按键重新划分再移动,1,000 万条要经过网络和磁盘。如果把 300 个商品给每个任务各发一份,订单原地不动,只复制 300 个。结果完全一样,移动的量却相差数万倍。实际工作中,“连接很慢”这句话有一半是这个选择选错了。
工作原理
广播哈希连接(BroadcastHashJoin)。驱动器把较小一侧的表收集起来,给每个执行器发一份。执行器用它建哈希表,大表的每个分区就原地去查这个哈希表。大表一侧没有 shuffle。所以它最快,但较小的一侧必须真的很小——因为它要整体加载到驱动器和所有执行器的内存中。
Spark 自动选择它的依据是性能调优文档中的 spark.sql.autoBroadcastJoinThreshold。默认值是 10MB,根据统计信息得出的表大小比它小就会广播。设为 -1,就会关闭自动广播。同一份文档中的 spark.sql.broadcastTimeout 是等待广播的时间,默认是 300 秒。
排序合并连接(SortMergeJoin)。把两侧都按连接键的哈希做 shuffle,让相同的键汇集到相同编号的分区。然后在每个分区里把两侧按键排序,并排扫描两行来配对。两侧都要付出 shuffle 和排序的代价,但排序在内存不足时可以溢写到磁盘,所以在任何大小下都能走到最后。这就是两张表都很大时,它成为默认选择的原因。
shuffle 哈希连接(ShuffledHashJoin)。shuffle 与排序合并相同。不同的是,在每个分区里,不做排序,而是用较小的一侧建哈希表。跳过了排序,所以可能更快,但一个分区里较小的一侧必须装得进内存。所以 Spark 只在条件合适时才选它。
from pyspark.sql import functions as F
orders.join(products, "product_id").explain() # 작으면 BroadcastHashJoin
orders.join(F.broadcast(products), "product_id") # 크기와 상관없이 브로드캐스트
orders.join(products.hint("shuffle_hash"), "product_id") # ShuffledHashJoin 요청
orders.join(products.hint("merge"), "product_id") # SortMergeJoin 요청
读计划的时候,比起连接的名字,要看它下面挂着什么。如果是排序合并,两侧的分支各会挂上一个 Exchange hashpartitioning(product_id, …) 和一个 Sort——两次 shuffle 和两次排序。如果是广播,只有较小一侧挂着 BroadcastExchange,大表一侧的分支没有 Exchange。如果是 shuffle 哈希,两侧都有 Exchange,但没有 Sort。只要熟悉这三种形态,不用去找连接的名字,也能看出移动了什么。
提示及其局限
估算有时会出错。像 CSV 这样没有统计信息的原始数据,或者经过多次过滤的结果,很难猜出大小。这时由人来告诉它策略,就是提示(hint)。提示文档列出的连接提示有四个:BROADCAST(别名 BROADCASTJOIN、MAPJOIN)、MERGE(别名 SHUFFLE_MERGE、MERGEJOIN)、SHUFFLE_HASH、SHUFFLE_REPLICATE_NL。被加上 BROADCAST 提示的一侧,无论阈值如何都会被广播。在 DataFrame API 中,broadcast() 函数做的是同样的事。
如果在两侧加上不同的提示,按 BROADCAST、MERGE、SHUFFLE_HASH、SHUFFLE_REPLICATE_NL 的顺序,靠前的获胜。而且性能调优文档明确写道:提示不是保证。因为有些策略不支持特定的连接类型。例如在左外连接中,必须保留左侧的所有行,所以不能把左侧广播出去当作哈希表。因此加了提示之后,一定要在计划里确认实际选中了什么。
运行中改变的计划——AQE
自适应查询执行(AQE)从 3.2.0 起默认开启。AQE 会在 shuffle 结束之后,根据实际产生了多少字节,重新制定剩余的计划。对连接来说,重要的是把排序合并改成广播哈希的规则。根据性能调优文档,如果根据运行时统计信息得出的某一侧小于自适应阈值(spark.sql.adaptive.autoBroadcastJoinThreshold,默认与自动阈值相同),就会改。
文档在这里附了一条诚实的说明:它不如一开始就计划成广播那样高效。因为 shuffle 已经发生了。不过它能避免两侧的排序,而且在开启本地 shuffle 读取时,可以不经网络、原地读取 shuffle 文件。在计划里,最初是 AdaptiveSparkPlan isFinalPlan=false,运行后变成 isFinalPlan=true,并且能看到其中连接的名字发生了变化。
还有一条把排序合并改成 shuffle 哈希的规则。当 shuffle 之后所有分区都小于 spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold 时就会改,但这个值的默认值是 0,所以不单独开启就不会发生。
在现场相遇的样子
第一,连接之后行数变多了。如果连接键在一侧不是唯一的,配对就会成倍增加。一个订单对应促销表里同一个商品的三行,这个订单就会变成三行。销售额合计膨胀,报告出错,却没有任何错误。连接之前数一数较小一侧键的唯一性,这个习惯就能防止这类事故。
第二,广播会把驱动器搞垮。把阈值调得很大,或者滥用提示,会让几百 MB 的表汇集到驱动器,再复制到所有执行器。驱动器内存不足或广播超时,就发生在这个时候。
第三,找“没有的”时候,不要用外连接再过滤 null,而要用反连接。连接文档把反连接定义为返回右侧没有配对的左侧行的连接。找出一次都没下过单的客户,正是这种形态,结果里只保留左侧的列。
实际工作中真正重要的事
- 连接策略是对“移动什么”的选择。复制较小的一侧,大表就不用动。
- 自动广播的标准是 10MB。在没有统计信息的原始数据上,估算可能出错。
- 提示是请求,不是命令。加了之后要在计划里确认实际名字。
- AQE 会在运行中把排序合并改成广播。不过已经付出的 shuffle 要不回来。
- 连接之前要数键的唯一性。重复的键会在没有任何错误的情况下让行数膨胀。
下一项实验要做什么
把订单和小的商品表连接,通过计划确认 Spark 自己选择了广播。关掉阈值,看它变成排序合并,再用 broadcast 和 shuffle_hash 提示亲手改变策略。在把阈值调低的情况下,从最终计划里抓住 AQE 在运行中把排序合并改成广播的那一刻;把促销表里重复的键让行数膨胀的情况,与 left_semi 连接的行数对比着数一数,再用反连接找出没有订单的客户。