Apache Spark SQL 性能调优完全指南:缓存、分区、连接策略与自适应查询执行

Apache Spark SQL 性能调优完全指南:缓存、分区、连接策略与自适应查询执行 Apache Spark SQL 性能调优完全指南缓存、分区、连接策略与自适应查询执行【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark本文是 Apache Spark SQL 性能调优的实战指南围绕 DataFrame/SQL 工作负载的几大类调优手段展开数据缓存、分区调整、统计信息利用、聚合优化、连接策略选择、子计划合并、自适应查询执行AQE以及存储分区连接SPJ。文中所有配置项均以当前仓库 docs/sql-performance-tuning.md 为准并结合源码实现与测试用例给出底层原理说明。读完本文你将掌握每项优化技术的适用场景、关键配置参数的默认值与含义以及如何通过EXPLAIN、Web UI 等手段验证优化效果。Caching Data 缓存数据Spark SQL 可以使用内存列式格式缓存表通过spark.catalog.cacheTable(tableName)或dataFrame.cache()触发。缓存后Spark SQL 只扫描查询所需的列并自动为每一列选择压缩编解码器从而最小化内存占用和 GC 压力。移除缓存使用spark.catalog.uncacheTable(tableName)或dataFrame.unpersist()。检查缓存状态spark.catalog.isCached(tableName)判断指定表或视图是否已缓存读取任意Dataset的storageLevel属性未缓存时返回StorageLevel.NONE应用运行期间所有持久化对象的整体视图包括通过Dataset.cache()直接缓存的数据可查看 Web UI 的 Storage 页该页在 action 物化数据后展示每个持久化关系的存储级别、大小和分区数。Spark 支持两种缓存格式默认缓存格式标准的内存列式缓存默认使用Arrow 缓存格式基于 Apache Arrow 的缓存可改善列式工作负载的读取性能并支持与 Arrow 生态的互操作详见 Arrow Cache Format 文档。内存缓存相关配置可通过spark.conf.set或 SQL 的SET keyvalue命令设置属性名默认值含义引入版本spark.sql.inMemoryColumnarStorage.compressedtrue为 true 时Spark SQL 基于数据统计信息自动为每一列选择压缩编解码器1.0.1spark.sql.inMemoryColumnarStorage.batchSize10000控制列式缓存的分批大小。较大的 batchSize 可提升内存利用率和压缩效果但缓存数据时存在 OOM 风险1.1.1源码印证列式缓存相关配置定义在 sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala 中IN_MEMORY_TABLE_STORAGE_LEVEL等配置项与上述两个参数共同控制 InMemoryTableScanExec 的执行行为。Tuning Partitions 分区调优读取文件型数据源Parquet、JSON、ORC时分区数量的合理性直接影响并行度与任务效率。相关配置如下属性名默认值含义引入版本spark.sql.files.maxPartitionBytes134217728128 MB读取文件时打包到单个分区的最大字节数。仅对文件型数据源Parquet、JSON、ORC生效2.0.0spark.sql.files.openCostInBytes41943044 MB打开文件的估算开销以相同时间内可扫描的字节数衡量。用于把多个小文件合并进一个分区。建议高估该值这样小文件分区会比大文件分区先被调度更快。仅对文件型数据源生效2.0.0spark.sql.files.minPartitionNum默认并行度拆分的文件分区数的建议非保证最小值。未设置时默认取spark.sql.leafNodeDefaultParallelism的值。仅对文件型数据源生效3.1.0spark.sql.files.maxPartitionNum无拆分的文件分区数的建议非保证最大值。设置后若初始分区数超过该值Spark 会重新缩放每个分区使分区数接近该值。仅对文件型数据源生效3.5.0spark.sql.shuffle.partitions200对连接或聚合进行 shuffle 时使用的分区数1.1.0spark.sql.sources.parallelPartitionDiscovery.threshold32启用作业输入路径并行列举的阈值。输入路径数大于该阈值时Spark 使用分布式作业列举文件否则回退到顺序列举。仅对文件型数据源Parquet、ORC、JSON生效1.5.0spark.sql.sources.parallelPartitionDiscovery.parallelism10000作业输入路径的最大列举并行度。输入路径数超过该值时会被限流到该值。仅对文件型数据源生效2.1.1源码印证上述分区参数定义于 sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala#L3098-L3143。其中spark.sql.files.maxPartitionBytes默认值与parquet.block.size对齐128MBspark.sql.files.openCostInBytes标记为.internal()minPartitionNum与maxPartitionNum都校验必须为正整数。实际读取时FileSourceScanExec 会依据这些参数把文件拆分成接近目标大小的分区。Coalesce Hints 合并提示Coalesce hints 允许 Spark SQL 用户像 Dataset API 中的coalesce、repartition和repartitionByRange一样控制输出文件数量既可用于性能调优也可用于减少输出文件数。各 hint 的参数规则COALESCE只有一个分区数参数REPARTITION可有分区数、列、两者皆有或两者皆无作为参数REPARTITION_BY_RANGE必须有列名分区数可选REBALANCE可有初始分区数、列、两者皆有或两者皆无作为参数REBALANCE_BY_SIZE需要一个建议分区大小参数可选地后跟列。SELECT /* COALESCE(3) */ * FROM t; SELECT /* REPARTITION(3) */ * FROM t; SELECT /* REPARTITION(c) */ * FROM t; SELECT /* REPARTITION(3, c) */ * FROM t; SELECT /* REPARTITION */ * FROM t; SELECT /* REPARTITION_BY_RANGE(c) */ * FROM t; SELECT /* REPARTITION_BY_RANGE(3, c) */ * FROM t; SELECT /* REBALANCE */ * FROM t; SELECT /* REBALANCE(3) */ * FROM t; SELECT /* REBALANCE(c) */ * FROM t; SELECT /* REBALANCE(3, c) */ * FROM t; SELECT /* REBALANCE_BY_SIZE(134217728) */ * FROM t; SELECT /* REBALANCE_BY_SIZE(134217728, c) */ * FROM t; SELECT /* REBALANCE_BY_SIZE(128m) */ * FROM t; SELECT /* REBALANCE_BY_SIZE(128m, c) */ * FROM t;更多细节参见 Partitioning Hints 文档。其中REBALANCE与REBALANCE_BY_SIZE在开启 AQE 时还会与spark.sql.adaptive.optimizeSkewsInRebalancePartitions.enabled配合对倾斜分区进行拆分见下文 AQE 章节。Leveraging Statistics 利用统计信息Apache Spark 在众多候选执行计划中选优的能力取决于它对执行计划中每个节点read、filter、join 等输出行数的估计。这些估计基于通过以下途径提供给 Spark 的统计信息数据源统计Spark 直接从底层数据源读取的统计信息如 Parquet 文件元数据中的行数、min/max 值由数据源自身维护Catalog 统计Spark 从 catalog如 Hive Metastore读取的统计信息每次执行ANALYZE TABLE时收集或更新运行时统计Spark 在查询运行过程中自行计算的统计信息属于 自适应查询执行框架 的一部分。统计信息缺失或不准确会妨碍 Spark 选择最优计划可能导致查询性能下降。因此建议检查 Spark 可用的统计信息以及查询规划/执行阶段的估计值数据对象统计用DESCRIBE EXTENDED查看表或列的统计信息查询计划估计用EXPLAIN COST或DataFrame.explain(modecost)查看优化后查询计划中的代价估计运行时统计查询运行过程中可在 SQL UI 的 Details 区域查看在计划中寻找Statistics(..., isRuntimetrue)标记。Optimizing the Aggregate 优化聚合Adaptive Partial Aggregation 自适应部分聚合分组聚合通常分两个阶段执行shuffle 之前的部分聚合partial aggregation和 shuffle 之后的最终聚合final aggregation。部分聚合只有在确实减少行数时才有价值当分组键接近唯一时聚合 map 会膨胀到与输入差不多大小甚至发生溢出却几乎按原样输出消费的行数此时部分聚合得不偿失。启用自适应部分聚合后hash 聚合会在运行时测量压缩比compaction ratio——已处理行数除以聚合 map 中持有的键数。若部分聚合折叠的行数不足以抵消其开销则停止填充聚合 map将剩余行作为单行部分聚合缓冲直接透传给最终聚合合并。透传激活后map 会被冻结其输出始终排在透传行之前与冻结 map 中键冲突的行被排在队列里待 map 排空后才冲刷因此 map 中已存在键的重复行仍会合并到该键累积行之后保证first/last等对顺序敏感的聚合与从不透传的运行结果一致。压缩比会周期性评估也会在聚合 map 即将溢出前再次评估——此时溢出会被完全跳过。两次评估都使用当前 map 周期的累计行数与键数且一旦触发透传在当前任务剩余部分不会撤销因此偏斜的前缀会向任一方向影响结果有利前缀掩盖不利后缀较好的前缀压缩比可能掩盖后续大量不同的键使聚合持续到溢出才重新记账。若不想等到溢出可调大minCompaction让累计压缩比在较弱的后段趋势下更早触发阈值更高的阈值也会在其他输入上更激进地透传不同键前缀提前触发透传不同的键较多时可能过早触发透传并一直保持。可通过调大minRows推迟周期性评估的启动来避免过早锁定任务的透传决策。属性名默认值含义引入版本spark.sql.execution.aggregate.adaptivePartialAggregation.enabledfalse为 true 时hash 聚合在运行时观察到部分聚合未将行数减少到值得的程度会自适应地绕过 shuffle 前的部分聚合。仅适用于带分组键的 hash 聚合4.4.0spark.sql.execution.aggregate.adaptivePartialAggregation.minRows100000周期性压缩比评估之间的行数。设为 0 则禁用周期性评估但 map 即将溢出时仍可能评估。较大的值会推迟周期性评估因此触发透传时冻结 map 往往持有更多行冻结 map 在输出排空前驻留内存较大值会抬高这一瞬时内存峰值4.4.0spark.sql.execution.aggregate.adaptivePartialAggregation.minCompaction1.05保持部分聚合所需的最小压缩比。压缩比 10 表示部分聚合将十行折叠为一个键当某次评估发现压缩比低于该值时部分聚合在剩余输入上被绕过。更大的值绕过得更激进4.4.0源码印证该特性在 sql/core/src/main/scala/org/apache/spark/sql/execution/aggregate/HashAggregateExec.scala#L201-L213 中接入HashAggregateExec读取conf.adaptivePartialAggregationEnabled、adaptivePartialAggregationMinRows与adaptivePartialAggregationMinCompaction并在其聚合逻辑中实现压缩比采样、map 冻结与透传缓冲。三个配置项定义于 sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala#L4410-L4444均标注NOT_APPLICABLE绑定策略即不支持通过 SQLSET动态切换需在启动前通过配置指定。Optimizing the Join Strategy 优化连接策略Automatically Broadcasting Joins 自动广播连接属性名默认值含义引入版本spark.sql.autoBroadcastJoinThreshold1048576010 MB执行连接时广播到所有 worker 节点的表的最大字节数。设为 -1 可禁用广播1.1.0spark.sql.broadcastTimeout300广播连接中广播等待的超时时间秒1.3.0Join Strategy Hints 连接策略提示连接策略提示BROADCAST、MERGE、SHUFFLE_HASH和SHUFFLE_REPLICATE_NL指示 Spark 在将指定关系与其他关系连接时使用被提示的策略。例如对表t1使用BROADCAST提示时即使统计信息显示t1的大小超过spark.sql.autoBroadcastJoinThresholdSpark 仍会优先选择以t1为 build 侧的广播连接具体为 broadcast hash join 或 broadcast nested loop join取决于是否存在等值连接键。当连接两侧指定了不同的策略提示时Spark 的优先级为BROADCASTMERGESHUFFLE_HASHSHUFFLE_REPLICATE_NL。当两侧都指定BROADCAST或都指定SHUFFLE_HASH时Spark 根据连接类型和关系大小选择 build 侧。注意并不保证Spark 一定采用提示指定的策略因为特定策略可能不支持所有连接类型。各语言与 SQL 中的示例spark.table(src).join(spark.table(records).hint(broadcast), key).show()spark.table(src).join(spark.table(records).hint(broadcast), key).show()spark.table(src).join(spark.table(records).hint(broadcast), key).show();src - sql(SELECT * FROM src) records - sql(SELECT * FROM records) head(join(src, hint(records, broadcast), src$key records$key))-- We accept BROADCAST, BROADCASTJOIN and MAPJOIN for broadcast hint SELECT /* BROADCAST(r) */ * FROM src s JOIN records r ON s.key r.key更多细节参见 Join Hints 文档。Merging Subplans 合并子计划Spark 会合并返回单行且读取相同输入的子计划使输入只被扫描一次而不是每个子计划各扫描一次。合并的候选对象是非相关的确定性标量子查询和不带GROUP BY的非分组聚合。合并后的子计划只执行一次并输出单个 struct每个原始位置从该 struct 中读取自己需要的字段。该优化默认开启。例如以下查询的两个子查询都扫描store_salesSELECT (SELECT min(ss_net_paid) FROM store_sales), (SELECT max(ss_net_paid) FROM store_sales)它们被合并为一个同时计算min和max的聚合store_sales只被读取一次。在EXPLAIN输出中合并后的子计划表现为输出单列名为mergedValue的子查询共享它的位置显示为ReusedSubquery。两个子计划按节点逐一对齐才能合并Project列表求并集Aggregate必须具有相同的分组并使用相同的聚合实现因此min不会与collect_list合并Filter必须具有相同的条件Join必须具有相同的类型、条件和提示叶子节点必须读取相同的输入。对于行的内容依赖所读取列的 V1 文件关系仅当两侧子计划读取该关系的相同列时才可合并csv、json、xml其解析器根据所需 schema 判定什么算损坏记录以及任何在开启spark.sql.files.ignoreCorruptFiles作为读取选项或通过配置的情况下读取的文件关系——此时只有一侧读取的列发生读取失败会与整个文件的其余行一起被吞掉。spark.sql.files.ignoreMissingFiles也计入考量但原因是其中一个谓词替两者作答。仅WHERE条件不同的子计划也可以合并方法是将每侧的条件变成布尔列并给每侧的聚合表达式加上FILTER (WHERE ...)子句由下列配置控制。当规则执行时查询中仍包含WITH子句未被内联的会被跳过。在 DataSource V2 读取路径上对于声明了SCAN_MERGING表能力的源叶子读取相同输入的要求被放宽两个仅投影列不同的叶子可合并为读取这些列并集的单次扫描。内置文件格式中 Parquet、ORC、text 和 Avro 声明了该能力格式只有在被移出spark.sql.sources.useV1SourceList后才走 V2 读取路径。当spark.sql.files.ignoreCorruptFiles为 true 时文件表会放弃该能力因为仅另一子计划投影的列发生读取失败时会被连同该文件其余行一起吞掉spark.sql.files.ignoreMissingFiles为 true 时也会放弃以匹配文件读取器使用的严格性谓词。当两个子计划中只有一个带 filter 时合并总是有益的因为未过滤侧反正要读全部数据。这种情况默认开启除非 filter 必须跨越Join才能到达聚合这需要下面的 through-join 配置。当两侧都带 filter对称情形时合并后的扫描过滤条件变成OR(f1, f2)其选择性低于任一原始 filter因此可能读取更多数据——例如 filter 裁剪分区或 Parquet row group 时。这正是对称情形默认关闭的原因。不过对于同一张表上计算多个不同过滤聚合这类常见分析形态仍然值得开启该优化SELECT (SELECT avg(ss_net_paid) FROM store_sales WHERE ss_quantity BETWEEN 1 AND 20), (SELECT avg(ss_net_paid) FROM store_sales WHERE ss_quantity BETWEEN 21 AND 40)在 TPC-DS 基准测试中开启对称 filter 传播使q9和q28提速约 3.5 倍同时开启经过 join 的传播后q88约提速 7 倍、q90约提速 2 倍测量数据见 SPARK-40193 与 SPARK-56677。收益取决于表当差异 filter 位于数据源无法裁剪的列上时收益最大在重度分区或文件裁剪的表上加宽的 filter 会丢失裁剪能力风险最高。在生产环境启用前务必在自己的工作负载上验证。属性名默认值含义引入版本spark.sql.optimizer.mergeSubplans.filterPropagation.enabledtrue为 true 时仅 filter 条件不同的子计划可通过将 filter 传播到外层非分组聚合来合并。这是下面三个配置的总开关它为 false 时三者均不生效。filter 条件相同的子计划不受此配置影响始终可合并4.2.0spark.sql.optimizer.mergeSubplans.filterPropagation.symmetricFilterPropagation.enabledfalse为 true 时两侧都带 filter 条件的非分组聚合子计划也可合并。默认关闭因为合并后的 filter 被加宽为OR(f1, f2)可能比两个原始 filter 读取更多数据尤其在重度分区或文件裁剪的表上4.2.0spark.sql.optimizer.mergeSubplans.filterPropagation.throughJoin.enabledfalse为 true 时filter 条件也可跨Join节点传播使仅 filter 条件不同且共享同一 join 的子计划可合并。为 false 时不跨 join 传播任何 filter即使只有一侧带 filter。filter 只能从 join 的保留侧传播LEFT OUTER/LEFT SEMI/LEFT ANTI的左、RIGHT OUTER的右、INNER/CROSS的任一侧。FULL OUTER连接永远不适用。filter 不同的子计划通常两侧都有 filter因此该配置通常与symmetricFilterPropagation.enabled一起开启4.2.0spark.sql.optimizer.mergeSubplans.filterPropagation.dsv2SymmetricFilterPropagation.enabledfalse为 true 时两个下推了相同严格强制 filter 但携带不同 best-effort扫描后filter 的 DataSource V2 扫描即使在symmetricFilterPropagation.enabled为 false 时也可合并。此情形下加宽不会改变扫描必须返回的行集合因为严格 filter 会原样重新下推外层Filter会在扫描之上重新检查其余条件。仅适用于通过SCAN_MERGING表能力选择加入扫描合并的 V2 源。对文件源而言严格强制 filter 即分区 filter因此该配置允许同一分区、不同数据 filter的两个扫描合并4.3.0源码印证子计划合并规则实现于 sql/core/src/main/scala/org/apache/spark/sql/execution/planmerging/MergeSubplans.scala#L151规则类型为Rule[LogicalPlan]配套的 PlanMerger.scala 负责实际的对齐与合并逻辑。相关行为有完整的测试覆盖例如 MergeSubplansSuite.scala、DSv2PlanMergingSuite.scala 与 FileSourceV2PlanMergingSuite.scala。完全关闭子计划合并可将规则加入spark.sql.optimizer.excludedRulesspark.sql.optimizer.excludedRulesorg.apache.spark.sql.execution.planmerging.MergeSubplans必须使用这个准确名称该规则在 Spark 4.2 之前叫MergeScalarSubqueries且 Spark 4.3 之前位于不同的包spark.sql.optimizer.excludedRules中未知的名称会被静默忽略因此从旧版本迁移过来的旧名称并不会关闭该规则。旧名称参见 SQL migration guide。Adaptive Query Execution 自适应查询执行自适应查询执行AQE是 Spark SQL 中的一种优化技术利用运行时统计信息选择最高效的查询执行计划自 Apache Spark 3.2.0 起默认开启。可通过spark.sql.adaptive.enabled作为总开关开启或关闭 AQE。属性名默认值含义引入版本spark.sql.adaptive.enabledtrue为 true 时启用自适应查询执行基于准确的运行时统计信息在查询执行中途重新优化查询计划1.6.0Coalescing Post Shuffle Partitions 合并 shuffle 后分区当spark.sql.adaptive.enabled和spark.sql.adaptive.coalescePartitions.enabled都为 true 时该特性基于 map 输出统计信息合并 shuffle 后的分区。它简化了运行查询时 shuffle 分区数的调优无需为数据集设置精确的分区数只需通过spark.sql.adaptive.coalescePartitions.initialPartitionNum设置足够大的初始 shuffle 分区数Spark 就能在运行时选出合适的分区数。属性名默认值含义引入版本spark.sql.adaptive.coalescePartitions.enabledtrue为 true 且spark.sql.adaptive.enabled为 true 时Spark 根据目标大小由spark.sql.adaptive.advisoryPartitionSizeInBytes指定合并连续的 shuffle 分区避免过多小任务3.0.0spark.sql.adaptive.coalescePartitions.parallelismFirsttrue为 true 时合并连续 shuffle 分区时忽略spark.sql.adaptive.advisoryPartitionSizeInBytes默认 64MB指定的目标大小只尊重spark.sql.adaptive.coalescePartitions.minPartitionSize默认 1MB指定的最小分区大小以最大化并行度。这是为了避免启用 AQE 时的性能回退。建议在繁忙集群上将其设为 false以提高资源利用效率避免大量小任务3.2.0spark.sql.adaptive.coalescePartitions.minPartitionSize1MB合并后 shuffle 分区的最小大小。当合并分区时目标大小被忽略默认情形时该参数有用3.2.0spark.sql.adaptive.coalescePartitions.maxReducerPartitionsPerTaskInt.MaxValue可合并进单个任务的连续 reducer 分区最大数量。它独立于建议分区大小限制 reducer 分区的扇入4.3.0spark.sql.adaptive.coalescePartitions.initialPartitionNum无合并前 shuffle 分区的初始数量。未设置时等于spark.sql.shuffle.partitions。仅在spark.sql.adaptive.enabled与spark.sql.adaptive.coalescePartitions.enabled同时为 true 时生效3.0.0spark.sql.adaptive.advisoryPartitionSizeInBytes64 MB自适应优化期间 shuffle 分区的建议字节大小spark.sql.adaptive.enabled为 true 时。Spark 合并小 shuffle 分区或拆分倾斜 shuffle 分区时生效3.0.0Splitting skewed shuffle partitions 拆分倾斜的 shuffle 分区属性名默认值含义引入版本spark.sql.adaptive.optimizeSkewsInRebalancePartitions.enabledtrue为 true 且spark.sql.adaptive.enabled为 true 时Spark 优化 RebalancePartitions 中倾斜的 shuffle 分区按目标大小spark.sql.adaptive.advisoryPartitionSizeInBytes指定将其拆分为更小的分区避免数据倾斜3.2.0spark.sql.adaptive.rebalancePartitionsSmallPartitionFactor0.2拆分过程中若分区大小小于该因子乘以spark.sql.adaptive.advisoryPartitionSizeInBytes则该分区会被合并3.3.0Converting sort-merge join to broadcast join 将 sort-merge join 转换为 broadcast join当连接任一侧的运行时统计信息小于自适应广播连接阈值时AQE 将 sort-merge join 转换为 broadcast hash join。这不如一开始就规划 broadcast hash join 高效但仍优于继续 sort-merge join——可以避免对两侧排序并在本地读取 shuffle 文件以节省网络流量前提是spark.sql.adaptive.localShuffleReader.enabled为 true。属性名默认值含义引入版本spark.sql.adaptive.autoBroadcastJoinThreshold无连接时广播到所有 worker 节点的表的最大字节数。设为 -1 可禁用广播。默认值与spark.sql.autoBroadcastJoinThreshold相同。注意此配置仅在自适应框架中使用3.2.0spark.sql.adaptive.localShuffleReader.enabledtrue为 true 且spark.sql.adaptive.enabled为 true 时在 shuffle 分区不再需要时例如 sort-merge join 转换为 broadcast-hash join 后Spark 尝试使用本地 shuffle reader 读取 shuffle 数据3.0.0Converting sort-merge join to shuffled hash join 将 sort-merge join 转换为 shuffled hash join当所有 post shuffle 分区都小于spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold配置的阈值时AQE 将 sort-merge join 转换为 shuffled hash join。属性名默认值含义引入版本spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold0允许构建本地 hash map 的每分区最大字节数。若该值不小于spark.sql.adaptive.advisoryPartitionSizeInBytes且所有分区大小都不超过该配置则无论spark.sql.join.preferSortMergeJoin的值如何连接选择都倾向使用 shuffled hash join 而非 sort merge join3.2.0spark.sql.adaptive.convertSortMergeJoinToShuffledHashJoin.enabledtrue为 true 时自适应执行期间当 build 侧物化的每分区大小都在spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold之内时同时要求spark.sql.adaptive.advisoryPartitionSizeInBytes不大于它Spark 将 sort-merge join 转换为 shuffled hash join。这是转换的总开关。默认只穿透 join 自身所需的 sort 到达直接输入 shuffle设置spark.sql.adaptive.convertSortMergeJoinToShuffledHashJoin.lookThroughOperators.enabled为 true 可同时穿透 join 与输入 shuffle 之间的非 shuffle 算子如 aggregate、project、filter、window4.3.0spark.sql.adaptive.convertSortMergeJoinToShuffledHashJoin.lookThroughOperators.enabledfalse为 true 时sort-merge join 到 shuffled hash join 的转换额外穿透 join 与输入 shuffle 之间的非 shuffle 算子如 aggregate、project、filter、window而不只是 join 自身所需的 sort。仅在convertSortMergeJoinToShuffledHashJoin.enabled为 true 时生效4.3.0spark.sql.adaptive.convertSortMergeJoinToShuffledHashJoin.minWideningFactor1.0sort-merge-join 到 shuffled-hash-join 转换约束 build 侧 shuffled hash map 大小时应用的行加宽因子下限。该因子用 join 与 shuffle 之间算子的估算单行大小增长来缩放输入 shuffle 字节较大的下限更保守在统计信息可能低估 build 大小时使转换更不可能发生。必须为正数4.3.0spark.sql.adaptive.costEvaluator.countLocalSort.enabledspark.sql.adaptive.convertSortMergeJoinToShuffledHashJoin.lookThroughOperators.enabled的值为 true 时默认的 AQE 代价评估器还会将本地排序数量作为 shuffle 数量之下的低优先级决胜项计入因此 shuffle 数相同的计划中本地排序更少者被优先选择。例如只有转换不会在计划其他地方引入额外排序时sort-merge join 才被 shuffled hash join 替换。默认随 look-through 转换一起开启4.3.0Optimizing Skew Join 优化倾斜连接数据倾斜会严重降低连接查询的性能。该特性通过在 sort-merge join 中动态处理倾斜将倾斜任务拆分为必要时复制大小大致均匀的任务。它在spark.sql.adaptive.enabled与spark.sql.adaptive.skewJoin.enabled同时开启时生效。属性名默认值含义引入版本spark.sql.adaptive.skewJoin.enabledtrue为 true 且spark.sql.adaptive.enabled为 true 时Spark 动态处理 sort-merge join 中的倾斜拆分必要时复制倾斜分区3.0.0spark.sql.adaptive.skewJoin.skewedPartitionFactor5.0分区大小大于该因子乘以分区大小中位数且同时大于spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes时判定为倾斜分区3.0.0spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes256MB分区字节大小大于该阈值且同时大于spark.sql.adaptive.skewJoin.skewedPartitionFactor乘以分区大小中位数时判定为倾斜分区。理想情况下该配置应设置得比spark.sql.adaptive.advisoryPartitionSizeInBytes大3.0.0spark.sql.adaptive.forceOptimizeSkewedJoinfalse为 true 时强制启用 OptimizeSkewedJoin——即使引入额外 shuffle 也要优化倾斜连接以避免落后任务3.3.0Advanced Customization 高级定制可以通过提供自定义代价评估器类或排除 AQE 优化器规则控制 AQE 的细节。属性名默认值含义引入版本spark.sql.adaptive.optimizer.excludedRules无配置自适应优化器中要禁用的规则列表按规则名指定并用逗号分隔。优化器会记录确实被排除的规则3.1.0spark.sql.adaptive.customCostEvaluatorClass无用于自适应执行的自定义代价评估器类。未设置时 Spark 默认使用自己的SimpleCostEvaluator3.2.0Storage Partition Join 存储分区连接存储分区连接SPJ是 Spark SQL 中的一种优化技术利用现有存储布局避免 shuffle 阶段。它是 Bucket Join 概念的推广——Bucket Join 只适用于 bucketed分桶 表而 SPJ 可适用于按 FunctionCatalog 中注册的函数分区的表。存储分区连接目前支持兼容的 V2 DataSource。以下 SQL 属性在不同连接查询中以各种优化方式启用存储分区连接属性名默认值含义引入版本spark.sql.sources.v2.bucketing.enabledtrue为 true 时尝试利用兼容 V2 数据源报告的划分信息消除 shuffle3.3.0spark.sql.sources.v2.bucketing.pushPartValues.enabledtrue启用后若连接一侧相对于另一侧缺少分区值则尝试消除 shuffle。此配置要求spark.sql.sources.v2.bucketing.enabled为 true3.4.0spark.sql.requireAllClusterKeysForCoPartitiontrue为 true 时存储分区连接要求每个 join 或 MERGE 键都被某个分区键覆盖而非按位置匹配分区键才能消除 shuffle。当分区键只覆盖部分 join 或 MERGE 键时可设为false以消除 shuffle但代价是较粗的存储分区可能带来数据倾斜和并行度下降3.3.0spark.sql.sources.v2.bucketing.partiallyClusteredDistribution.enabledfalse为 true 且连接不是 full outer join 时在避免 shuffle 的同时启用倾斜优化来处理数据量大的分区。将根据表统计信息选择一侧作为大表该侧的拆分采用部分聚类partially-clustered另一侧的拆分被分组并复制以匹配。此配置要求spark.sql.sources.v2.bucketing.enabled和spark.sql.sources.v2.bucketing.pushPartValues.enabled都为 true3.4.0spark.sql.sources.v2.bucketing.allowKeysSubsetOfPartitionKeys.enabledfalse启用后若 join 或 MERGE 条件未包含全部分区列也尝试避免 shuffle。此配置要求spark.sql.sources.v2.bucketing.enabled为 true4.0.0spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabledfalse启用后若分区 transform 兼容但不完全相同也尝试避免 shuffle。此配置要求spark.sql.sources.v2.bucketing.enabled与pushPartValues.enabled都为 true且partiallyClusteredDistribution.enabled为 false4.0.0spark.sql.sources.v2.bucketing.shuffle.enabledfalse启用后通过识别另一侧 V2 数据源报告的分区信息尝试避免连接一侧的 shuffle4.0.0如果执行了存储分区连接查询计划中 join 之前将不会出现 Exchange 节点。下面的示例使用 Iceberg——一个支持存储分区连接的 Spark V2 DataSourceCREATE TABLE prod.db.target (id INT, salary INT, dep STRING) USING iceberg PARTITIONED BY (dep, bucket(8, id)) CREATE TABLE prod.db.source (id INT, salary INT, dep STRING) USING iceberg PARTITIONED BY (dep, bucket(8, id)) EXPLAIN SELECT * FROM target t INNER JOIN source s ON t.dep s.dep AND t.id s.id -- Plan without Storage Partition Join Physical Plan * Project (12) - * SortMergeJoin Inner (11) :- * Sort (5) : - Exchange (4) // DATA SHUFFLE : - * Filter (3) : - * ColumnarToRow (2) : - BatchScan (1) - * Sort (10) - Exchange (9) // DATA SHUFFLE - * Filter (8) - * ColumnarToRow (7) - BatchScan (6) SET spark.sql.sources.v2.bucketing.enabled true SET spark.sql.iceberg.planning.preserve-data-grouping true SET spark.sql.sources.v2.bucketing.pushPartValues.enabled true -- Plan with Storage Partition Join Physical Plan * Project (10) - * SortMergeJoin Inner (9) :- * Sort (4) : - * Filter (3) : - * ColumnarToRow (2) : - BatchScan (1) - * Sort (8) - * Filter (7) - * ColumnarToRow (6) - BatchScan (5)对比两个计划可以发现启用 SPJ 后join 之前的ExchangeDATA SHUFFLE节点被消除两侧的数据直接以存储布局的天然划分参与连接显著节省了 shuffle 与网络开销。对于倾斜连接可以考虑启用spark.sql.sources.v2.bucketing.partiallyClusteredDistribution.enabled并在自己的工作负载上测量效果。该选项会复制连接一侧的分区可能增加读取的数据量上面的示例保持其默认值false。总结Spark SQL 的性能调优是一个从静态配置到动态自适应的渐进体系列式缓存与分区参数解决数据摆放与并行度问题统计信息与EXPLAIN COST帮助理解优化器的决策依据聚合、连接策略与子计划合并让静态优化器做出更聪明的选择而 AQE 与存储分区连接则将决策时机推进到运行时用真实数据分布指导计划调整。实践建议先用EXPLAIN与 Web UI 定位瓶颈再针对性调整上述配置并在生产环境前于自有工作负载上逐一验证尤其对默认关闭的symmetricFilterPropagation、partiallyClusteredDistribution等激进优化。本文全部配置的权威定义见 docs/sql-performance-tuning.md源码实现位于 sql/catalyst 与 sql/core 两个模块。【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考