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 key=value命令设置:
| 属性名 | 默认值 | 含义 | 引入版本 |
|---|---|---|---|
spark.sql.inMemoryColumnarStorage.compressed | true | 为 true 时,Spark SQL 基于数据统计信息自动为每一列选择压缩编解码器 | 1.0.1 |
spark.sql.inMemoryColumnarStorage.batchSize | 10000 | 控制列式缓存的分批大小。较大的 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.maxPartitionBytes | 134217728(128 MB) | 读取文件时打包到单个分区的最大字节数。仅对文件型数据源(Parquet、JSON、ORC)生效 | 2.0.0 |
spark.sql.files.openCostInBytes | 4194304(4 MB) | 打开文件的估算开销,以"相同时间内可扫描的字节数"衡量。用于把多个小文件合并进一个分区。建议高估该值,这样小文件分区会比大文件分区(先被调度)更快。仅对文件型数据源生效 | 2.0.0 |
spark.sql.files.minPartitionNum | 默认并行度 | 拆分的文件分区数的建议(非保证)最小值。未设置时默认取spark.sql.leafNodeDefaultParallelism的值。仅对文件型数据源生效 | 3.1.0 |
spark.sql.files.maxPartitionNum | 无 | 拆分的文件分区数的建议(非保证)最大值。设置后,若初始分区数超过该值,Spark 会重新缩放每个分区,使分区数接近该值。仅对文件型数据源生效 | 3.5.0 |
spark.sql.shuffle.partitions | 200 | 对连接或聚合进行 shuffle 时使用的分区数 | 1.1.0 |
spark.sql.sources.parallelPartitionDiscovery.threshold | 32 | 启用作业输入路径并行列举的阈值。输入路径数大于该阈值时,Spark 使用分布式作业列举文件,否则回退到顺序列举。仅对文件型数据源(Parquet、ORC、JSON)生效 | 1.5.0 |
spark.sql.sources.parallelPartitionDiscovery.parallelism | 10000 | 作业输入路径的最大列举并行度。输入路径数超过该值时会被限流到该值。仅对文件型数据源生效 | 2.1.1 |
源码印证:上述分区参数定义于 sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala#L3098-L3143。其中spark.sql.files.maxPartitionBytes默认值与parquet.block.size对齐(128MB);spark.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(mode="cost")查看优化后查询计划中的代价估计; - 运行时统计:查询运行过程中可在 SQL UI 的 "Details" 区域查看,在计划中寻找
Statistics(..., isRuntime=true)标记。
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.enabled | false | 为 true 时,hash 聚合在运行时观察到部分聚合未将行数减少到值得的程度,会自适应地绕过 shuffle 前的部分聚合。仅适用于带分组键的 hash 聚合 | 4.4.0 |
spark.sql.execution.aggregate.adaptivePartialAggregation.minRows | 100000 | 周期性压缩比评估之间的行数。设为 0 则禁用周期性评估(但 map 即将溢出时仍可能评估)。较大的值会推迟周期性评估,因此触发透传时冻结 map 往往持有更多行;冻结 map 在输出排空前驻留内存,较大值会抬高这一瞬时内存峰值 | 4.4.0 |
spark.sql.execution.aggregate.adaptivePartialAggregation.minCompaction | 1.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.autoBroadcastJoinThreshold | 10485760(10 MB) | 执行连接时广播到所有 worker 节点的表的最大字节数。设为 -1 可禁用广播 | 1.1.0 |
spark.sql.broadcastTimeout | 300 | 广播连接中广播等待的超时时间(秒) | 1.3.0 |
Join Strategy Hints 连接策略提示
连接策略提示BROADCAST、MERGE、SHUFFLE_HASH和SHUFFLE_REPLICATE_NL指示 Spark 在将指定关系与其他关系连接时使用被提示的策略。例如对表t1使用BROADCAST提示时,即使统计信息显示t1的大小超过spark.sql.autoBroadcastJoinThreshold,Spark 仍会优先选择以t1为 build 侧的广播连接(具体为 broadcast hash join 或 broadcast nested loop join,取决于是否存在等值连接键)。
当连接两侧指定了不同的策略提示时,Spark 的优先级为:BROADCAST>MERGE>SHUFFLE_HASH>SHUFFLE_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_sales:
SELECT (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.enabled | true | 为 true 时,仅 filter 条件不同的子计划可通过将 filter 传播到外层非分组聚合来合并。这是下面三个配置的总开关:它为 false 时三者均不生效。filter 条件相同的子计划不受此配置影响,始终可合并 | 4.2.0 |
spark.sql.optimizer.mergeSubplans.filterPropagation.symmetricFilterPropagation.enabled | false | 为 true 时,两侧都带 filter 条件的非分组聚合子计划也可合并。默认关闭,因为合并后的 filter 被加宽为OR(f1, f2),可能比两个原始 filter 读取更多数据,尤其在重度分区或文件裁剪的表上 | 4.2.0 |
spark.sql.optimizer.mergeSubplans.filterPropagation.throughJoin.enabled | false | 为 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.0 |
spark.sql.optimizer.mergeSubplans.filterPropagation.dsv2SymmetricFilterPropagation.enabled | false | 为 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.excludedRules:
spark.sql.optimizer.excludedRules=org.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.enabled | true | 为 true 时启用自适应查询执行,基于准确的运行时统计信息在查询执行中途重新优化查询计划 | 1.6.0 |
Coalescing 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.enabled | true | 为 true 且spark.sql.adaptive.enabled为 true 时,Spark 根据目标大小(由spark.sql.adaptive.advisoryPartitionSizeInBytes指定)合并连续的 shuffle 分区,避免过多小任务 | 3.0.0 |
spark.sql.adaptive.coalescePartitions.parallelismFirst | true | 为 true 时,合并连续 shuffle 分区时忽略spark.sql.adaptive.advisoryPartitionSizeInBytes(默认 64MB)指定的目标大小,只尊重spark.sql.adaptive.coalescePartitions.minPartitionSize(默认 1MB)指定的最小分区大小,以最大化并行度。这是为了避免启用 AQE 时的性能回退。建议在繁忙集群上将其设为 false,以提高资源利用效率(避免大量小任务) | 3.2.0 |
spark.sql.adaptive.coalescePartitions.minPartitionSize | 1MB | 合并后 shuffle 分区的最小大小。当合并分区时目标大小被忽略(默认情形)时该参数有用 | 3.2.0 |
spark.sql.adaptive.coalescePartitions.maxReducerPartitionsPerTask | Int.MaxValue | 可合并进单个任务的连续 reducer 分区最大数量。它独立于建议分区大小限制 reducer 分区的扇入 | 4.3.0 |
spark.sql.adaptive.coalescePartitions.initialPartitionNum | 无 | 合并前 shuffle 分区的初始数量。未设置时等于spark.sql.shuffle.partitions。仅在spark.sql.adaptive.enabled与spark.sql.adaptive.coalescePartitions.enabled同时为 true 时生效 | 3.0.0 |
spark.sql.adaptive.advisoryPartitionSizeInBytes | 64 MB | 自适应优化期间 shuffle 分区的建议字节大小(spark.sql.adaptive.enabled为 true 时)。Spark 合并小 shuffle 分区或拆分倾斜 shuffle 分区时生效 | 3.0.0 |
Splitting skewed shuffle partitions 拆分倾斜的 shuffle 分区
| 属性名 | 默认值 | 含义 | 引入版本 |
|---|---|---|---|
spark.sql.adaptive.optimizeSkewsInRebalancePartitions.enabled | true | 为 true 且spark.sql.adaptive.enabled为 true 时,Spark 优化 RebalancePartitions 中倾斜的 shuffle 分区,按目标大小(spark.sql.adaptive.advisoryPartitionSizeInBytes指定)将其拆分为更小的分区,避免数据倾斜 | 3.2.0 |
spark.sql.adaptive.rebalancePartitionsSmallPartitionFactor | 0.2 | 拆分过程中,若分区大小小于该因子乘以spark.sql.adaptive.advisoryPartitionSizeInBytes,则该分区会被合并 | 3.3.0 |
Converting 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.0 |
spark.sql.adaptive.localShuffleReader.enabled | true | 为 true 且spark.sql.adaptive.enabled为 true 时,在 shuffle 分区不再需要时(例如 sort-merge join 转换为 broadcast-hash join 后),Spark 尝试使用本地 shuffle reader 读取 shuffle 数据 | 3.0.0 |
Converting 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.maxShuffledHashJoinLocalMapThreshold | 0 | 允许构建本地 hash map 的每分区最大字节数。若该值不小于spark.sql.adaptive.advisoryPartitionSizeInBytes且所有分区大小都不超过该配置,则无论spark.sql.join.preferSortMergeJoin的值如何,连接选择都倾向使用 shuffled hash join 而非 sort merge join | 3.2.0 |
spark.sql.adaptive.convertSortMergeJoinToShuffledHashJoin.enabled | true | 为 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、window) | 4.3.0 |
spark.sql.adaptive.convertSortMergeJoinToShuffledHashJoin.lookThroughOperators.enabled | false | 为 true 时,sort-merge join 到 shuffled hash join 的转换额外穿透 join 与输入 shuffle 之间的非 shuffle 算子(如 aggregate、project、filter、window),而不只是 join 自身所需的 sort。仅在convertSortMergeJoinToShuffledHashJoin.enabled为 true 时生效 | 4.3.0 |
spark.sql.adaptive.convertSortMergeJoinToShuffledHashJoin.minWideningFactor | 1.0 | sort-merge-join 到 shuffled-hash-join 转换约束 build 侧 shuffled hash map 大小时应用的行加宽因子下限。该因子用 join 与 shuffle 之间算子的估算单行大小增长来缩放输入 shuffle 字节;较大的下限更保守,在统计信息可能低估 build 大小时使转换更不可能发生。必须为正数 | 4.3.0 |
spark.sql.adaptive.costEvaluator.countLocalSort.enabled | spark.sql.adaptive.convertSortMergeJoinToShuffledHashJoin.lookThroughOperators.enabled的值 | 为 true 时,默认的 AQE 代价评估器还会将本地排序数量作为 shuffle 数量之下的低优先级决胜项计入,因此 shuffle 数相同的计划中,本地排序更少者被优先选择。例如,只有转换不会在计划其他地方引入额外排序时,sort-merge join 才被 shuffled hash join 替换。默认随 look-through 转换一起开启 | 4.3.0 |
Optimizing Skew Join 优化倾斜连接
数据倾斜会严重降低连接查询的性能。该特性通过在 sort-merge join 中动态处理倾斜:将倾斜任务拆分为(必要时复制)大小大致均匀的任务。它在spark.sql.adaptive.enabled与spark.sql.adaptive.skewJoin.enabled同时开启时生效。
| 属性名 | 默认值 | 含义 | 引入版本 |
|---|---|---|---|
spark.sql.adaptive.skewJoin.enabled | true | 为 true 且spark.sql.adaptive.enabled为 true 时,Spark 动态处理 sort-merge join 中的倾斜,拆分(必要时复制)倾斜分区 | 3.0.0 |
spark.sql.adaptive.skewJoin.skewedPartitionFactor | 5.0 | 分区大小大于该因子乘以分区大小中位数,且同时大于spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes时,判定为倾斜分区 | 3.0.0 |
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes | 256MB | 分区字节大小大于该阈值,且同时大于spark.sql.adaptive.skewJoin.skewedPartitionFactor乘以分区大小中位数时,判定为倾斜分区。理想情况下该配置应设置得比spark.sql.adaptive.advisoryPartitionSizeInBytes大 | 3.0.0 |
spark.sql.adaptive.forceOptimizeSkewedJoin | false | 为 true 时强制启用 OptimizeSkewedJoin——即使引入额外 shuffle 也要优化倾斜连接以避免落后任务 | 3.3.0 |
Advanced Customization 高级定制
可以通过提供自定义代价评估器类或排除 AQE 优化器规则,控制 AQE 的细节。
| 属性名 | 默认值 | 含义 | 引入版本 |
|---|---|---|---|
spark.sql.adaptive.optimizer.excludedRules | 无 | 配置自适应优化器中要禁用的规则列表,按规则名指定并用逗号分隔。优化器会记录确实被排除的规则 | 3.1.0 |
spark.sql.adaptive.customCostEvaluatorClass | 无 | 用于自适应执行的自定义代价评估器类。未设置时 Spark 默认使用自己的SimpleCostEvaluator | 3.2.0 |
Storage Partition Join 存储分区连接
存储分区连接(SPJ)是 Spark SQL 中的一种优化技术,利用现有存储布局避免 shuffle 阶段。
它是 Bucket Join 概念的推广——Bucket Join 只适用于 bucketed(分桶) 表,而 SPJ 可适用于按 FunctionCatalog 中注册的函数分区的表。存储分区连接目前支持兼容的 V2 DataSource。
以下 SQL 属性在不同连接查询中以各种优化方式启用存储分区连接:
| 属性名 | 默认值 | 含义 | 引入版本 |
|---|---|---|---|
spark.sql.sources.v2.bucketing.enabled | true | 为 true 时,尝试利用兼容 V2 数据源报告的划分信息消除 shuffle | 3.3.0 |
spark.sql.sources.v2.bucketing.pushPartValues.enabled | true | 启用后,若连接一侧相对于另一侧缺少分区值,则尝试消除 shuffle。此配置要求spark.sql.sources.v2.bucketing.enabled为 true | 3.4.0 |
spark.sql.requireAllClusterKeysForCoPartition | true | 为 true 时,存储分区连接要求每个 join 或 MERGE 键都被某个分区键覆盖(而非按位置匹配分区键)才能消除 shuffle。当分区键只覆盖部分 join 或 MERGE 键时,可设为false以消除 shuffle,但代价是较粗的存储分区可能带来数据倾斜和并行度下降 | 3.3.0 |
spark.sql.sources.v2.bucketing.partiallyClusteredDistribution.enabled | false | 为 true 且连接不是 full outer join 时,在避免 shuffle 的同时启用倾斜优化来处理数据量大的分区。将根据表统计信息选择一侧作为大表,该侧的拆分采用部分聚类(partially-clustered);另一侧的拆分被分组并复制以匹配。此配置要求spark.sql.sources.v2.bucketing.enabled和spark.sql.sources.v2.bucketing.pushPartValues.enabled都为 true | 3.4.0 |
spark.sql.sources.v2.bucketing.allowKeysSubsetOfPartitionKeys.enabled | false | 启用后,若 join 或 MERGE 条件未包含全部分区列,也尝试避免 shuffle。此配置要求spark.sql.sources.v2.bucketing.enabled为 true | 4.0.0 |
spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled | false | 启用后,若分区 transform 兼容但不完全相同,也尝试避免 shuffle。此配置要求spark.sql.sources.v2.bucketing.enabled与pushPartValues.enabled都为 true,且partiallyClusteredDistribution.enabled为 false | 4.0.0 |
spark.sql.sources.v2.bucketing.shuffle.enabled | false | 启用后,通过识别另一侧 V2 数据源报告的分区信息,尝试避免连接一侧的 shuffle | 4.0.0 |
如果执行了存储分区连接,查询计划中 join 之前将不会出现 Exchange 节点。
下面的示例使用 Iceberg——一个支持存储分区连接的 Spark V2 DataSource:
CREATE 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 之前的Exchange(DATA 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),仅供参考