辅助优化规则
Kyuubi 开箱即用地提供了 SQL 扩展。由于与 Apache Spark 的版本兼容性,目前我们支持 Apache Spark branch-3.1 及更高版本。同时不必担心,Kyuubi 未来将会支持新的 Apache Spark 版本。得益于自适应查询执行框架(AQE),Kyuubi 能够完成这些优化。
特性
自动合并小文件
小文件一直是 Apache Spark 长期存在的问题。Kyuubi 可以通过增加一次额外的 shuffle 来合并小文件。目前,Kyuubi 支持处理数据源表(datasource table)和 Hive 表中的小文件,同时也支持优化动态分区写入。例如,对于一个常见的写入查询
INSERT INTO TABLE $table1 SELECT * FROM $table2,Kyuubi 会在写入之前引入一次额外的 shuffle,从而消除小文件。在 Join 之前插入 shuffle 节点,使 AQE 的 OptimizeSkewedJoin 生效
在当前实现中,Apache Spark 只能针对标准 join(standard join)优化倾斜 join(skewed join),这意味着一个 join 必须具备两个 sort 和 shuffle 节点。然而在复杂场景下,这个假设很容易被打破。Kyuubi 可以通过在 join 之前增加一个额外的 shuffle 节点来保证 join 是标准的,从而使 OptimizeSkewedJoin 能够更好地工作。
AQE 中的 stage 级别配置隔离
众所周知,
spark.sql.adaptive.advisoryPartitionSizeInBytes是 Apache Spark AQE 中的一个关键配置。它控制着 shuffle 期间每个 task 应该处理多大的数据量,因此我们通常会使用 64MB 或更小的值来保证足够的并行度。然而在一般情况下,我们希望文件足够大,比如 256MB 或 512MB。Kyuubi 可以通过配置隔离来解决这一冲突,使得 staging 阶段的分区数据量较小,而最终的分区数据量较大。
使用方法
| Kyuubi Spark SQL 扩展 | 支持的 Spark 版本 | 可用起始版本 | 终止支持(EOL) | 包含于二进制发布包 | Maven profile |
|---|---|---|---|---|---|
| kyuubi-extension-spark-3-1 | 3.1.x | 1.3.0-incubating | 1.8.0 | 1.3.0-incubating | spark-3.1 |
| kyuubi-extension-spark-3-2 | 3.2.x | 1.4.0-incubating | N/A | 1.4.0-incubating | spark-3.2 |
| kyuubi-extension-spark-3-3 | 3.3.x | 1.6.0-incubating | N/A | 1.6.0-incubating | spark-3.3 |
| kyuubi-extension-spark-3-4 | 3.4.x | 1.8.0 | N/A | 1.8.0 | spark-3.4 |
| kyuubi-extension-spark-3-5 | 3.5.x | 1.8.0 | N/A | 1.9.0 | spark-3.5 |
查看上述矩阵,确认你使用的 Spark 版本是否受支持,并找到对应的 Kyuubi Spark SQL 扩展 jar
获取 Kyuubi Spark SQL 扩展 jar
- 每个 Kyuubi 二进制发布包仅包含一个默认版本的 Kyuubi Spark SQL Extension jar,如果你需要该版本,可以在
$KYUUBI_HOME/extension下找到。 - 所有受支持版本的 Kyuubi Spark SQL Extension jar 都将发布到 Maven 中央仓库
- 如果你愿意,也可以自行编译 Kyuubi Spark SQL Extension jar,在编译命令中激活相应的 Maven profile,例如使用
-Pspark-3.5编译时,可在extensions/spark/kyuubi-extension-spark-3-5/target下获取适用于 Spark 3.5 的 Kyuubi Spark SQL Extension jar。
- 每个 Kyuubi 二进制发布包仅包含一个默认版本的 Kyuubi Spark SQL Extension jar,如果你需要该版本,可以在
将 Kyuubi Spark SQL extension jar
kyuubi-extension-spark-*.jar放入$SPARK_HOME/jars中启用
KyuubiSparkSQLExtension,即在$SPARK_HOME/conf/spark-defaults.conf中添加一项配置,spark.sql.extensions=org.apache.kyuubi.sql.KyuubiSparkSQLExtension
现在,你就可以畅享 Kyuubi SQL Extension 了。
附加配置
Kyuubi 提供了一些配置,让这些功能更易于使用。
| 名称 | 默认值 | 描述 | 引入版本 |
|---|---|---|---|
| spark.sql.optimizer.insertRepartitionBeforeWrite.enabled | true | 在查询计划的顶部添加 repartition 节点。一种合并小文件的方法。 | 1.2.0 |
| spark.sql.optimizer.forceShuffleBeforeJoin.enabled | false | 确保在 shuffled join(shj 和 smj)之前存在 shuffle 节点,使 AQE OptimizeSkewedJoin 生效(复杂场景 join、多表 join)。 | 1.2.0 |
| spark.sql.optimizer.finalStageConfigIsolation.enabled | false | 如果为 true,最后阶段可以使用与之前阶段不同的配置。最后阶段配置项的前缀应为 spark.sql.finalStage.。例如,原始的 spark 配置为:spark.sql.adaptive.advisoryPartitionSizeInBytes,则最后阶段的配置应为:spark.sql.finalStage.adaptive.advisoryPartitionSizeInBytes。 | 1.2.0 |
| spark.sql.optimizer.insertZorderBeforeWriting.enabled | true | 为 true 时,我们将根据目标表属性决定是否插入 zorder。关键属性包括:1) kyuubi.zorder.enabled:如果该属性为 true,我们会在写入数据前插入 zorder。2) kyuubi.zorder.cols:以逗号分隔的字符串,我们将按这些列进行 zorder。 | 1.4.0 |
| spark.sql.optimizer.zorderGlobalSort.enabled | true | 为 true 时,我们使用 zorder 进行全局排序。注意,如果 zorder 列的基数较低,可能会导致数据倾斜问题。为 false 时,我们仅使用 zorder 进行局部排序。 | 1.4.0 |
| spark.sql.watchdog.maxPartitions | none | 设置 Spark 扫描数据源时的最大分区数。通过指定该配置启用 maxPartition 策略。添加 maxPartitions 策略以避免扫描分区表时扫描过多分区,该策略为可选,配合已定义的配置使用。 | 1.4.0 |
| spark.sql.watchdog.maxFileSize | none | 设置 Spark 扫描数据源时文件的最大字节大小。通过指定该配置启用 maxFileSize 策略。添加 maxFileSize 策略以避免扫描文件的总大小过大,该策略为可选,配合已定义的配置使用。 | 1.8.0 |
| spark.sql.optimizer.dropIgnoreNonExistent | false | 为 true 时,如果 DROP DATABASE/TABLE/VIEW/FUNCTION/PARTITION 指定了不存在的数据库/表/视图/函数/分区,则不会报错。 | 1.5.0 |
| spark.sql.optimizer.rebalanceBeforeZorder.enabled | false | 为 true 时,我们在 zorder 之前执行 rebalance,以防止数据倾斜。注意,如果写入是动态分区的,我们将使用分区列进行 rebalance。注意,该配置仅对 Spark 3.3.x 生效。 | 1.6.0 |
| spark.sql.optimizer.rebalanceZorderColumns.enabled | false | 当该配置和 spark.sql.optimizer.rebalanceBeforeZorder.enabled 均为 true 时,我们在 Z-Order 之前执行 rebalance。如果是动态分区写入,rebalance 表达式将同时包含分区列和 Z-Order 列。注意,该配置仅对 Spark 3.3.x 生效。 | 1.6.0 |
| spark.sql.optimizer.twoPhaseRebalanceBeforeZorder.enabled | false | 当该配置和 spark.sql.optimizer.rebalanceBeforeZorder.enabled 均为 true 时,我们在动态分区写入的 Z-Order 之前执行两阶段 rebalance。第一阶段使用动态分区列进行 rebalance;第二阶段使用动态分区列和 Z-Order 列进行 rebalance。注意,该配置仅对 Spark 3.3.x 生效。 | 1.6.0 |
| spark.sql.optimizer.zorderUsingOriginalOrdering.enabled | false | 当该配置和 spark.sql.optimizer.rebalanceBeforeZorder.enabled 均为 true 时,我们按原始顺序(即字典序)进行排序。注意,该配置仅对 Spark 3.3.x 生效。 | 1.6.0 |
| spark.sql.optimizer.inferRebalanceAndSortOrders.enabled | false | 为 true 时,从原始查询中推断 rebalance 和排序顺序的列,例如 join 的连接键。这可以避免压缩率下降。 | 1.7.0 |
| spark.sql.optimizer.inferRebalanceAndSortOrdersMaxColumns | 3 | 推断列的最大数量。 | 1.7.0 |
| spark.sql.optimizer.insertRepartitionBeforeWriteIfNoShuffle.enabled | false | 为 true 时,即使原始计划中没有 shuffle,也添加 repartition。 | 1.7.0 |
| spark.sql.optimizer.finalStageConfigIsolationWriteOnly.enabled | true | 为 true 时,仅为写入启用最后阶段的配置隔离。 | 1.7.0 |
| spark.sql.finalWriteStage.eagerlyKillExecutors.enabled | false | 为 true 时,在运行最终写入阶段之前主动终止多余的 executor。 | 1.8.0 |
| spark.sql.finalWriteStage.skipKillingExecutorsForTableCache | true | 为 true 时,如果计划中包含表缓存,则跳过终止 executor。 | 1.8.0 |
| spark.sql.finalWriteStage.retainExecutorsFactor | 1.2 | 如果 目标 executor 数 × factor < 活跃 executor 数,且 目标 executor 数 × factor > 最小 executor 数,则注入终止 executor 或注入自定义资源档案。 | 1.8.0 |
| spark.sql.finalWriteStage.resourceIsolation.enabled | false | 为 true 时,使用自定义 RDD 资源档案实现最终写入阶段的资源隔离。 | 1.8.0 |
| spark.sql.finalWriteStageExecutorCores | fallback spark.executor.cores | 指定最终写入阶段的 executor 核心数请求,将传递给 RDD 资源档案。 | 1.8.0 |
| spark.sql.finalWriteStageExecutorMemory | fallback spark.executor.memory | 指定最终写入阶段的 executor 堆内内存请求,将传递给 RDD 资源档案。 | 1.8.0 |
| spark.sql.finalWriteStageExecutorMemoryOverhead | fallback spark.executor.memoryOverhead | 指定最终写入阶段的 executor 内存开销请求,将传递给 RDD 资源档案。 | 1.8.0 |
| spark.sql.finalWriteStageExecutorOffHeapMemory | NONE | 指定最终写入阶段的 executor 堆外内存请求,将传递给 RDD 资源档案。 | 1.8.0 |
| spark.sql.execution.scriptTransformation.enabled | true | 为 false 时,不允许使用脚本转换。 | 1.9.0 |
评论
登录后参与评论
KnowForge