一条 join 落在集群上可能走四种完全不同的算法代价差出一个量级甚至更多。选哪一种的不是你是 Spark 自己依据是它手里那份统计信息。统计信息大概率是估的估歪了计划不会报错只会换一条更贵的路走下去一路走到笛卡尔积也没人提醒你。下面的参数与默认值全按 Apache Spark 4.x 的口径。一. 谁在替你选什么时候定下来物理规划阶段有一个专门挑 join 算法的策略类4.x 里是SparkStrategies.scala的JoinSelection。它拿到两棵逻辑计划树和自带的统计信息产出一个物理算子定完就进执行阶段。判据读的是plan.stats.sizeInBytes那张表在这一刻被认定的字节数。它可能是昨天 ANALYZE 出来的可能是按文件大小猜的可能是过滤条件估出来的一个比例唯独不可能是刚称过的真重量。没有 hint 的时候源码注释把顺序写死了五条规则一条条往下试。有一侧小到能广播且 join 类型允许广播这一侧走 BroadcastHashJoin有一侧小到能在本地建哈希表且比另一侧小得多同时spark.sql.join.preferSortMergeJoin是 false走 ShuffledHashJoinjoin key 的类型能排序走 SortMergeJoinjoin 类型是 inner 系走 CartesianProduct前面全不走走 BroadcastNestedLoopJoin注释原话是这一步可能 OOM但没有别的选择计划里的节点名和源码算子类一一对应BroadcastHashJoin 是BroadcastHashJoinExecSortMergeJoin 是SortMergeJoinExecShuffledHashJoin 是ShuffledHashJoinExecBroadcastNestedLoopJoin 是BroadcastNestedLoopJoinExecCartesianProduct 是CartesianProductExec。没有等值 key 的 join 走另一条阶梯先看一侧能不能广播能就 BroadcastNestedLoopJoininner 系再考虑 CartesianProduct最后仍回到 BroadcastNestedLoopJoin 广播较小那一侧。两条阶梯都收在这张兜底网里。二. BroadcastHashJoin赌的是那一侧真的小自动触发靠spark.sql.autoBroadcastJoinThreshold4.x 文档给的默认值是 10MB判据是sizeInBytes不超过它设成-1就把广播这条路整个封掉。省掉的东西很诱人小表发到每个 executor 内存里大表边读边在 map 端配对没有 shuffle也没有排序和跨节点搬运。成本分三段都能在源码里指着看。driver 先把整张小表收上来走executeCollectIterator这一步受spark.driver.maxResultSize约束4.x 默认 1g超了作业直接失败。收上来要建成哈希表或数组指标是time to build建完分发出去的指标是time to broadcast等它的上限是spark.sql.broadcastTimeout默认 300 秒。到了 executor 侧这张表每个 executor 存一份是乘法关系的内存占用。4.1.0 起多了一道硬闸spark.sql.maxBroadcastTableSize默认 8589934592 字节即 8GiB广播数据超过它就抛错终止不给你慢慢 OOM 的机会。行数上也有一道写死在BroadcastExchangeExec走哈希表那条路的按BytesToBytesMap容量上限折算其余情况是 512,000,000 行。哪一侧能当广播侧不是想当然的。joins.scala里canBuildBroadcastLeft只放行 inner 系和 right outercanBuildBroadcastRight放行 inner 系、left outer、left single、left semi、left anti 和 existence。也就是说 left join 的左表不能广播左表要保行广播补不上缺失配对想广播只能改 SQL。三. SortMergeJoin 是兜底不是最优两边的 key 各自排序再像拉链一样归并配对这是 sort-merge join。它是大表 join 最常见的落点原因不在它快在于它几乎总是可用RowOrdering.isOrderable过线就轮到它。代价是两块两侧各排一次序大表排序是磁盘和外溢的常客两侧又都要按 key 做 hash 重分区分区数由spark.sql.shuffle.partitions决定默认 200。4.x 里有一件事和很多老笔记不一致spark.sql.join.preferSortMergeJoin的默认值是 true。这个键在源码里标了 internal不在公开配置表的列面上默认 true 意味着即便 ShuffledHashJoin 的条件全满足规则二也被堵死sort-merge 仍是默认落点。3.x 时代它默认 false从那时候留下来的调优清单会误导人。sort-merge 最典型的翻车姿势是倾斜。hash 分区按 key 哈希投递同一个 key 的所有行必须落到同一个分区一个热点 key 就能把那个分区撑成别人的一百倍。AQE 判倾斜是两道线spark.sql.adaptive.skewJoin.skewedPartitionFactor默认 5.0 配上spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes默认 256MB两道同时过才动手拆。一道都没过的热点分区sort-merge 原样交给一个 task199 个分区先跑完剩它一个啃撑爆的那份stage 的墙钟被它定住。认它有三个位置。Stages 页最慢 stage 的 Summary Metrics 里task 耗时的 max 和 median 差出一个量级是信号再对一眼这些慢 task 的 Shuffle Read Size 是不是同时也是最大的那几个。4.x 的 stage 详情页没有替你标好倾斜的徽标只能自己读数字。AQE 真动手拆过的分区露面在 SQL 页执行计划里节点是AQEShuffleRead描述里写着 skewed。四. ShuffledHashJoin三条线同时过才轮到它它和 sort-merge 的区别只有一个不排序。两侧照样 shuffle 到同一批分区然后在分区内直接建哈希表配对省掉排序代价是每个 task 要在自己内存里立一张哈希表。joins.scala里这条路的判据是三个条件叠起来的。spark.sql.join.preferSortMergeJoin必须是 false4.x 默认 true这一条默认就把门堵上了。build 侧的sizeInBytes要小于spark.sql.autoBroadcastJoinThreshold乘以spark.sql.shuffle.partitions按默认值是 10MB 乘 200。build 侧的sizeInBytes乘spark.sql.shuffledHashJoinFactor默认 3还要不大于另一侧。后两条一条管绝对量一条管相对比例都是字节口径源码里没有行数参与这道判断想显式点名就挂SHUFFLE_HASHhint。AQE 那条转换路默认也是死的spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold默认 0还要求它不小于spark.sql.adaptive.advisoryPartitionSizeInBytes这个配置自身不给默认值源码里是回落到内部键spark.sql.adaptive.shuffle.targetPostShuffleInputSize的 64MB不动配置永远不会触发。从计划里认它看节点名ShuffledHashJoin后面跟着BuildLeft或BuildRight标出建表侧。选错的信号集中在内存这个 stage 的 task 峰值内存和 Spill Memory 一起抬高Shuffle Read 总量看着并不离谱就是单个分区装不下。它比 sort-merge 更怕倾斜同一个热点 keysort-merge 是慢shuffled hash 是建表建到 OOM。五. 退化成 BroadcastNestedLoopJoin 和笛卡尔积的是条件的形状这一类降级和大小无关是形状问题。ExtractEquiJoinKeys匹配不到等值 key等值那整段规则就直接跳过不管你的表多小。a.ts between b.start and b.end只有范围没有等号a.k b.k和a.k like b.pattern都是非等值比较。形状过关了还有类型这一关这一关在 4.x 是新添的。hashJoinSupported要求所有 join key 的类型满足UnsafeRowUtils.isBinaryStableSpark 4.0 起字符串列带 collation大小写不敏感或重音不敏感的字符串比较没法直接比二进制这类 key 会把 BroadcastHashJoin 和 ShuffledHashJoin 两条路一起跳过。跳的时候 driver 日志打一条 WARN原文是Hash based joins are not supported due to joining on keys that dont support binary equality后面跟着具体哪几个 key 和它的类型名。搜这条 WARN 是识别此类降级最快的入口。key 能哈希但没法排序也一样出事RowOrdering.isOrderable过不去sort-merge 也关掉于是只剩最后两条。inner 系走 CartesianProduct两侧一行配一行地乘出来其余 join 类型走 BroadcastNestedLoopJoin把较小的一侧广播出去逐行配对。指标上的特征很醒目。CartesianProduct 和 BroadcastNestedLoopJoin 都只有一个 metric名字是number of output rows。前者的输出量级是两侧行数相乘一千万行配一万行就是中间那一千亿行这个数字直接写在节点指标上不用猜。后者更隐蔽输出行数被条件过滤后看着不大但每个 task 都把整张广播表重扫一遍task 时间均匀地长而没有哪个特别慢这种要看 stage 总耗时和 Task State 里 REMAINING 的分布。单列的NOT IN子查询是个例外4.x 里由ExtractSingleColumnNullAwareAntiJoin兜住命中就直接生成BroadcastHashJoin ... LeftAnti广播右侧不查大小也不查阈值。六. 计划里写的那一行可能不是跑出来的那一行spark.sql.adaptive.enabled默认 true开着的时候物理规划定下来的那份计划文本只是草稿它在每个 shuffle 边界拿落地的真实数字重跑一遍选择。改判有三个方向。一侧的真实字节数掉到广播阈值以下sort-merge 换成 BroadcastHashJoin阈值spark.sql.adaptive.autoBroadcastJoinThreshold默认沿用非 adaptive 那份换完广播spark.sql.adaptive.localShuffleReader.enabled默认 true让 Spark 尽量本地读 shuffle 文件。所有 post-shuffle 分区都低于spark.sql.adaptive.maxShuffledHashJoinLocalMapThresholdsort-merge 换成 shuffled hash这道默认 0等于默认不走。一侧 map 统计里的非空分区比例低于spark.sql.adaptive.nonEmptyPartitionRatioForBroadcastJoininternal 键默认 0.2它就被判掉广播资格DynamicJoinSelection挂上 NO_BROADCAST_HASH 改走 shuffle join理由是空分区多时大量 task 一侧为空立刻返回shuffle join 反而快。关键在 AQE 只改判它看得见的那一侧。DynamicJoinSelection匹配的是携带MapOutputStatistics的 shuffle query stage从没经过 shuffle 的 leaf 扫描在 AQE 眼里仍是编译期那份估算。小表侧是直接从盘上读进来的维表时广播决策从头到尾用的都是估算值估歪了没人纠正一路开到BroadcastExchange让 driver 收全表。ANALYZE TABLE ... COMPUTE STATISTICS是全量收集带NOSCAN只收字节数不扫数据。表元数据里没有统计时spark.sql.statistics.fallBackToHdfs默认 false而它只对非分区的 Hive 表有意义分区表和其它数据源走另一套粗估。列级直方图靠spark.sql.statistics.histogram.enabled默认 false没开它过滤后的行数只能按启发式猜。想知道自己站在哪一把尺上有三个入口。EXPLAIN FORMATTED在执行前顶层是AdaptiveSparkPlan isFinalPlanfalse执行后取会切成 true同屏给出 Final Plan 和 Initial Plan 两段diff 它们就是改判的凭据。Spark UI 的 SQL 页节点挂的指标是跑出来的真数和计划文本里那行估算对照估得离不离谱一眼可见。spark.sql.planChangeLog.level设成 WARN 能把每条规则对计划的改动写进日志4.x 源码默认 TRACE 也就是不打配spark.sql.planChangeLog.rules只留 join 相关的规则。七. 判据和动手边界判错的时候按这个顺序读别一上来翻配置。先读计划里的 join 节点名落在哪一格就回去对上面那一节的条件。再读BroadcastExchange的data size、time to collect、time to build、time to broadcast四个指标前两个偏大就是赌错了后两个偏长是建表和分发吃力。最后 diffFinal Plan与Initial Plan确认运行时到底改没改判。配置能不能扭转看判据读的是哪一类东西。读大小和读比例的那部分能管阈值调大调小、设-1关广播、spark.sql.join.preferSortMergeJoin设 false 开 shuffled hash 的门、倾斜两道线往下压让拆分更激进、广播超时与内存闸配合 executor 余量逐步抬。读条件形状和读类型的那部分管不了非等值条件不会因为调阈值变成等值得给范围条件补一个等值锚点先按天或按桶键等值配对再在结果上过滤范围collation 让字符串 key 不能哈希就把 key 换成二进制可比的形式。left join 的左表不能广播是 join 语义决定的想广播只能改 SQL。聚合的热点在这套机制的管辖之外spark.sql.adaptive.skewJoin.enabled默认 true只管 join 一侧的 sort-merge 和 shuffled hash聚合要自己把热点打散成两段再聚。统计缺失也不是一次 SET 补得上的ANALYZE 一遍或者把过滤后的结果先落成带统计的中间表再 join。hint 站在两者中间BROADCAST、SHUFFLE_HASH、MERGE、SHUFFLE_REPLICATE_NL直接压过大小判断。它绕过的是估算绕不过约束类型和 join 语义不配合时 hint 照样被忽略用了 hint 一定要回来看计划有没有真的变。spark.sql.optimizer.excludedRules与spark.sql.adaptive.optimizer.excludedRules管的是优化器规则名4.x 里 join 算法的选择是一个 Strategy 而不是规则想关掉某条 join 路径靠的还是阈值和 hint 这两只手。