性能调优
Spark 提供了多种技术来调整 DataFrame 或 SQL 工作负载的性能。从广义上讲,这些技术包括缓存数据、更改数据集的分区方式、选择最佳连接策略,以及向优化器提供可用于构建更高效执行计划的额外信息。
缓存数据
Spark SQL 可以通过调用 spark.catalog.cacheTable("tableName") 或 dataFrame.cache() 使用内存列式格式缓存表。随后,Spark SQL 将仅扫描所需的列,并自动调整压缩以最小化内存使用和 GC 压力。您可以调用 spark.catalog.uncacheTable("tableName") 或 dataFrame.unpersist() 从内存中删除该表。
内存缓存的配置可以通过 spark.conf.set 或在 SQL 中运行 SET key=value 命令来完成。
| 属性名称 | 默认值 | 含义 | 起始版本 |
|---|---|---|---|
spark.sql.inMemoryColumnarStorage.compressed |
true | 设置为 true 时,Spark SQL 将根据数据统计信息自动为每一列选择压缩编解码器。 | 1.0.1 |
spark.sql.inMemoryColumnarStorage.batchSize |
10000 | 控制列式缓存的批次大小。较大的批次大小可以提高内存利用率和压缩效果,但在缓存数据时有导致内存溢出 (OOM) 的风险。 | 1.1.1 |
调整分区
| 属性名称 | 默认值 | 含义 | 起始版本 |
|---|---|---|---|
spark.sql.files.maxPartitionBytes |
134217728 (128 MB) | 读取文件时放入单个分区的最大字节数。此配置仅在使用基于文件的数据源(如 Parquet、JSON 和 ORC)时有效。 | 2.0.0 |
spark.sql.files.openCostInBytes |
4194304 (4 MB) | 打开文件的估计成本,以在相同时间内可以扫描的字节数来衡量。这用于将多个文件放入一个分区时。最好进行高估,这样小文件的分区将比大文件的分区更快(因为小文件分区会被优先调度)。此配置仅在使用基于文件的数据源(如 Parquet、JSON 和 ORC)时有效。 | 2.0.0 |
spark.sql.files.minPartitionNum |
默认并行度 | 建议(非保证)的文件拆分分区的最小数量。如果未设置,则默认值为 `spark.sql.leafNodeDefaultParallelism`。此配置仅在使用基于文件的数据源(如 Parquet、JSON 和 ORC)时有效。 | 3.1.0 |
spark.sql.files.maxPartitionNum |
None | 建议(非保证)的文件拆分分区的最大数量。如果设置了此值,当初始分区数超过该值时,Spark 将重新调整每个分区,使分区数量接近此值。此配置仅在使用基于文件的数据源(如 Parquet、JSON 和 ORC)时有效。 | 3.5.0 |
spark.sql.shuffle.partitions |
200 | 配置执行连接或聚合进行 Shuffle 时使用的分区数量。 | 1.1.0 |
spark.sql.sources.parallelPartitionDiscovery.threshold |
32 | 配置启用作业输入路径并行列表的阈值。如果输入路径的数量大于此阈值,Spark 将使用 Spark 分布式作业来列出文件。否则,它将回退到顺序列表。此配置仅在使用基于文件的数据源(如 Parquet、ORC 和 JSON)时有效。 | 1.5.0 |
spark.sql.sources.parallelPartitionDiscovery.parallelism |
10000 | 配置作业输入路径的最大列表并行度。如果输入路径的数量大于此值,它将被限制为使用此值。此配置仅在使用基于文件的数据源(如 Parquet、ORC 和 JSON)时有效。 | 2.1.1 |
Coalesce 提示
Coalesce 提示允许 Spark SQL 用户像 Dataset API 中的 coalesce、repartition 和 repartitionByRange 一样控制输出文件的数量,它们可用于性能调优和减少输出文件的数量。“COALESCE”提示仅包含分区数量作为参数。“REPARTITION”提示包含分区数量、列,或者两者皆有/皆无作为参数。“REPARTITION_BY_RANGE”提示必须包含列名,分区数量是可选的。“REBALANCE”提示包含初始分区数量、列,或者两者皆有/皆无作为参数。
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;
更多详细信息,请参考 分区提示 (Partitioning Hints) 文档。
利用统计信息
Apache Spark 在众多可能的选项中选择最佳执行计划的能力,部分取决于它对执行计划中每个节点(读取、过滤、连接等)将输出多少行的估计。这些估计又基于通过多种方式提供给 Spark 的统计信息:
- 数据源:Spark 直接从底层数据源读取的统计信息,例如 Parquet 文件元数据中的计数和最小值/最大值。这些统计信息由底层数据源维护。
- 目录 (Catalog):Spark 从目录(如 Hive Metastore)读取的统计信息。每当运行
ANALYZE TABLE时,这些统计信息就会被收集或更新。 - 运行时:Spark 在查询运行时自行计算的统计信息。这是 自适应查询执行框架 的一部分。
缺失或不准确的统计信息会阻碍 Spark 选择最佳计划的能力,并可能导致较差的查询性能。因此,检查 Spark 可用的统计信息以及它在查询规划和执行期间所做的估计会很有帮助。
- 数据对象统计信息:您可以使用
DESCRIBE EXTENDED检查表或列的统计信息。 - 查询计划估计:您可以通过
EXPLAIN COST或DataFrame.explain(mode="cost")在优化后的查询计划中检查 Spark 的成本估计。 - 运行时统计信息:您可以在查询运行时在 SQL UI 的“Details”部分检查这些统计信息。在计划中查找
Statistics(..., isRuntime=true)。
优化连接策略
自动广播连接
| 属性名称 | 默认值 | 含义 | 起始版本 |
|---|---|---|---|
spark.sql.autoBroadcastJoinThreshold |
10485760 (10 MB) | 配置在执行连接时将广播到所有工作节点(worker node)的表的最大字节数。将此值设置为 -1 可禁用广播。 | 1.1.0 |
spark.sql.broadcastTimeout |
300 |
广播连接中广播等待时间的超时(以秒为单位)。 |
1.3.0 |
连接策略提示
连接策略提示,即 BROADCAST、MERGE、SHUFFLE_HASH 和 SHUFFLE_REPLICATE_NL,指示 Spark 在将指定的某个关系与另一个关系连接时使用该提示策略。例如,当在表“t1”上使用 BROADCAST 提示时,即使统计信息建议的表“t1”的大小超过了配置 spark.sql.autoBroadcastJoinThreshold,Spark 也会优先使用以“t1”为构建侧(build side)的广播连接(取决于是否存在等值连接键,可能是广播哈希连接或广播嵌套循环连接)。
当在连接的两侧指定了不同的连接策略提示时,Spark 的优先级顺序为:BROADCAST > MERGE > SHUFFLE_HASH > SHUFFLE_REPLICATE_NL。当两侧都指定了 BROADCAST 提示或 SHUFFLE_HASH 提示时,Spark 将根据连接类型和关系的大小来选择构建侧。
请注意,不能保证 Spark 会选择提示中指定的连接策略,因为特定策略可能不支持所有连接类型。
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) 文档。
自适应查询执行 (AQE)
自适应查询执行 (AQE) 是 Spark SQL 中的一种优化技术,它利用运行时统计信息来选择最高效的查询执行计划,自 Apache Spark 3.2.0 起默认启用。Spark SQL 可以通过 spark.sql.adaptive.enabled 作为总开关来开启或关闭 AQE。
| 属性名称 | 默认值 | 含义 | 起始版本 |
|---|---|---|---|
spark.sql.adaptive.enabled |
true | 当为 true 时,启用自适应查询执行,它会根据准确的运行时统计信息在查询执行过程中重新优化查询计划。 | 1.6.0 |
合并 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 时,Spark 在合并连续 Shuffle 分区时会忽略 spark.sql.adaptive.advisoryPartitionSizeInBytes(默认 64MB)指定的目标大小,而仅遵守 spark.sql.adaptive.coalescePartitions.minPartitionSize(默认 1MB)指定的最小分区大小,以最大化并行度。这是为了避免启用自适应查询执行时出现性能退化。建议在繁忙的集群上将此配置设置为 false,以使资源利用更高效(避免产生大量微小任务)。 |
3.2.0 |
spark.sql.adaptive.coalescePartitions.minPartitionSize |
1MB | 合并后 Shuffle 分区的最小大小。当在分区合并期间忽略目标大小时(默认情况),此配置很有用。 | 3.2.0 |
spark.sql.adaptive.coalescePartitions.initialPartitionNum |
(无) | 合并前的初始 Shuffle 分区数。如果未设置,则等于 spark.sql.shuffle.partitions。此配置仅在 spark.sql.adaptive.enabled 和 spark.sql.adaptive.coalescePartitions.enabled 均启用时有效。 |
3.0.0 |
spark.sql.adaptive.advisoryPartitionSizeInBytes |
64 MB | 自适应优化期间(当 spark.sql.adaptive.enabled 为 true 时)Shuffle 分区的建议字节大小。它在 Spark 合并小的 Shuffle 分区或拆分倾斜的 Shuffle 分区时生效。 |
3.0.0 |
拆分倾斜的 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 |
将排序合并连接转换为广播连接
当任一连接侧的运行时统计信息小于自适应广播哈希连接阈值时,AQE 会将排序合并连接转换为广播哈希连接。这不如直接规划广播哈希连接有效,但比继续使用排序合并连接要好,因为我们可以避免对连接两侧进行排序,并从本地读取 Shuffle 文件以节省网络流量(前提是 spark.sql.adaptive.localShuffleReader.enabled 为 true)。
| 属性名称 | 默认值 | 含义 | 起始版本 |
|---|---|---|---|
spark.sql.adaptive.autoBroadcastJoinThreshold |
(无) | 配置执行连接时将广播到所有工作节点(worker node)的表的最大字节数。将此值设置为 -1 可禁用广播。默认值与 spark.sql.autoBroadcastJoinThreshold 相同。请注意,此配置仅在自适应框架中使用。 |
3.2.0 |
spark.sql.adaptive.localShuffleReader.enabled |
true | 当为 true 且 spark.sql.adaptive.enabled 为 true 时,Spark 尝试在不需要 Shuffle 分区(例如,将排序合并连接转换为广播哈希连接后)时使用本地 Shuffle 读取器来读取 Shuffle 数据。 |
3.0.0 |
将排序合并连接转换为 Shuffle 哈希连接
当所有 Shuffle 后分区都小于 spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold 中配置的阈值时,AQE 会将排序合并连接转换为 Shuffle 哈希连接。
| 属性名称 | 默认值 | 含义 | 起始版本 |
|---|---|---|---|
spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold |
0 | 配置允许构建本地哈希映射的每个分区的最大字节数。如果此值不小于 spark.sql.adaptive.advisoryPartitionSizeInBytes 且所有分区的大小都不大于此配置值,则连接选择将倾向于使用 Shuffle 哈希连接,而不是排序合并连接,无论 spark.sql.join.preferSortMergeJoin 的值如何。 |
3.2.0 |
优化倾斜连接
数据倾斜会严重降低连接查询的性能。此功能通过将倾斜任务拆分(并在需要时复制)为大小大致相等的任务,动态处理排序合并连接中的倾斜。当 spark.sql.adaptive.enabled 和 spark.sql.adaptive.skewJoin.enabled 配置均启用时,此功能生效。
| 属性名称 | 默认值 | 含义 | 起始版本 |
|---|---|---|---|
spark.sql.adaptive.skewJoin.enabled |
true | 当为 true 且 spark.sql.adaptive.enabled 为 true 时,Spark 通过拆分(并在需要时复制)倾斜分区来动态处理排序合并连接中的倾斜。 |
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。这是一个用于优化倾斜连接以避免延迟任务(straggler tasks)的自适应规则,即使它引入了额外的 Shuffle 也在所不惜。 | 3.3.0 |
高级自定义
您可以通过提供自定义的成本评估器类或排除 AQE 优化器规则来控制 AQE 的具体工作细节。
| 属性名称 | 默认值 | 含义 | 起始版本 |
|---|---|---|---|
spark.sql.adaptive.optimizer.excludedRules |
(无) | 配置自适应优化器中要禁用的规则列表,规则按名称指定并用逗号分隔。优化器将记录实际已被排除的规则。 | 3.1.0 |
spark.sql.adaptive.customCostEvaluatorClass |
(无) | 用于自适应执行的自定义成本评估器类。如果未设置,Spark 默认将使用其内置的 SimpleCostEvaluator。 |
3.2.0 |
存储分区连接 (Storage Partition Join)
存储分区连接 (Storage Partition Join, SPJ) 是 Spark SQL 中的一种优化技术,它利用现有的存储布局来避免 Shuffle 阶段。
这是桶连接 (Bucket Joins) 概念的推广。桶连接仅适用于 分桶 (bucketed) 表,而存储分区连接适用于由 FunctionCatalog 中注册的函数进行分区的表。存储分区连接目前支持兼容的 V2 数据源。
以下 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 时,要求连接或 MERGE 键与分区键相同且顺序一致,以消除 Shuffle。因此,在这种情况下设置为 false 以消除 Shuffle。 | 3.4.0 |
spark.sql.sources.v2.bucketing.partiallyClusteredDistribution.enabled |
false | 当为 true 且连接不是全外连接 (full outer join) 时,启用倾斜优化,以在避免 Shuffle 时处理包含大量数据的分区。将根据表统计信息选择一侧作为大表,该侧的拆分将是部分聚类的。另一侧的拆分将被分组并复制以进行匹配。此配置要求 spark.sql.sources.v2.bucketing.enabled 和 spark.sql.sources.v2.bucketing.pushPartValues.enabled 均为 true。 |
3.4.0 |
spark.sql.sources.v2.bucketing.allowJoinKeysSubsetOfPartitionKeys.enabled |
false | 启用后,如果连接或 MERGE 条件未包含所有分区列,尝试避免 Shuffle。此配置要求 spark.sql.sources.v2.bucketing.enabled 和 spark.sql.sources.v2.bucketing.pushPartValues.enabled 均为 true,并且 spark.sql.requireAllClusterKeysForCoPartition 为 false。 |
4.0.0 |
spark.sql.sources.v2.bucketing.allowCompatibleTransforms.enabled |
false | 启用后,如果分区转换兼容但不完全相同,尝试避免 Shuffle。此配置要求 spark.sql.sources.v2.bucketing.enabled 和 spark.sql.sources.v2.bucketing.pushPartValues.enabled 均为 true。 |
4.0.0 |
spark.sql.sources.v2.bucketing.shuffle.enabled |
false | 启用后,通过识别 V2 数据源在另一侧报告的分区信息,尝试在连接的一侧避免 Shuffle。 | 4.0.0 |
如果执行了存储分区连接,查询计划将不会包含连接前的 Exchange 节点。
以下示例使用了 Iceberg (https://iceberg.org.cn/docs/latest/spark-getting-started/),这是一个支持存储分区连接的 Spark V2 数据源。
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'
SET 'spark.sql.requireAllClusterKeysForCoPartition' 'false'
SET 'spark.sql.sources.v2.bucketing.partiallyClusteredDistribution.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)