This is an automated email from the ASF dual-hosted git repository. czy006 pushed a commit to branch jira/JSPT-2568 in repository https://gitbox.apache.org/repos/asf/amoro.git
commit 8cef40cc64150ee5ef9f8147f6cde6a9f215fb58 Author: ConradJam <[email protected]> AuthorDate: Wed Jun 17 13:00:16 2026 +0800 [JSPT-2540][feat] 小文件优化:新增Paimon Manifest合并脚本 --- amoro-producer/README.md | 141 ++++++- .../amoro/producer/PaimonCompactManifest.scala | 424 +++++++++++++++++++++ 2 files changed, 564 insertions(+), 1 deletion(-) diff --git a/amoro-producer/README.md b/amoro-producer/README.md index 57c88cfe9..d92ba1b05 100644 --- a/amoro-producer/README.md +++ b/amoro-producer/README.md @@ -21,6 +21,7 @@ - Snapshot 清理入口类为 `org.apache.amoro.producer.PaimonExpireSnapshots`,任务会通过 Spark SQL 执行 `CALL sys.expire_snapshots`。 - 孤儿文件清理入口类为 `org.apache.amoro.producer.PaimonRemoveOrphanFiles`,任务会通过 Spark SQL 执行 `CALL <catalog>.sys.remove_orphan_files`。 +- Manifest 合并入口类为 `org.apache.amoro.producer.PaimonCompactManifest`,任务会通过 Spark SQL 执行 `CALL <catalog>.sys.compact_manifest`。 ## 打包 @@ -73,6 +74,20 @@ amoro-producer/target/amoro-producer-0.9-SNAPSHOT-jar-with-dependencies.jar 孤儿文件清理的目标表参数组合规则与 Snapshot 清理保持一致。 +## Manifest 合并参数说明 + +| 参数 | 是否必填 | 默认值 | 说明 | +|------|----------|--------|------| +| `--catalogName` | 否 | `paimon` | Paimon catalog 名称。任务会使用该 catalog 扫表并调用 `compact_manifest`。 | +| `--databaseName` | 否 | 无 | 数据库名。仅传该参数时,任务会扫描并处理该库下全部非临时表。 | +| `--tableName` | 否 | 无 | 表名。与 `--databaseName` 同时传入时,只处理指定表;单独传入时必须使用 `db.table` 格式。 | +| `--procedureOptions` | 否 | 空字符串 | 传给 `compact_manifest` 的 `options` 字符串。传空字符串时不拼接 `options` 参数。 | +| `--retryTimes` | 否 | `3` | 单表 `compact_manifest` SQL 失败后的最大执行次数。 | +| `--continueOnTableFailure` | 否 | `true` | 单表 Manifest 合并失败后是否继续处理后续表。 | +| `--help` / `-h` | 否 | 无 | 打印帮助信息。 | + +Manifest 合并的目标表参数组合规则与 Snapshot 清理保持一致。 + ## Snapshot 清理 Spark Submit Demo 以下示例基于 YARN cluster 模式提交,主类替换为 @@ -331,10 +346,134 @@ catalog 名称不同,只需要调整 `--catalogName` 以及对应的 --continueOnTableFailure true ``` +## Manifest 合并 Spark Submit Demo + +Manifest 合并入口显式接收 `--catalogName`,示例中使用 `paimon`。如果生产环境 +catalog 名称不同,只需要调整 `--catalogName` 以及对应的 +`spark.sql.catalog.<catalog>` 配置。 + +### Manifest 按数据库合并 + +```bash +./spark-submit \ + --master yarn \ + --deploy-mode cluster \ + --driver-memory 8g \ + --driver-cores 4 \ + --executor-memory 4g \ + --executor-cores 4 \ + --conf spark.dynamicAllocation.enabled=true \ + --conf spark.dynamicAllocation.minExecutors=1 \ + --conf spark.dynamicAllocation.maxExecutors=2 \ + --conf spark.dynamicAllocation.initialExecutors=1 \ + --conf spark.dynamicAllocation.shuffleTracking.enabled=true \ + --conf spark.driver.memoryOverhead=2g \ + --conf spark.executor.memoryOverhead=2g \ + --conf spark.driver.bindAddress=0.0.0.0 \ + --conf spark.sql.shuffle.partitions=400 \ + --conf spark.sql.cbo.enabled=true \ + --conf spark.sql.files.maxPartitionBytes=512m \ + --conf spark.default.parallelism=200 \ + --conf spark.locality.wait=0s \ + --conf spark.sql.crossJoin.enabled=true \ + --conf spark.sql.adaptive.localShuffleReader.enabled=true \ + --conf spark.sql.fuse.unionAllOnJoin.enabled=true \ + --conf spark.sql.optimizer.runtime.bloomFilter.enabled=true \ + --conf spark.sql.optimizer.enableMergeScalarAggsInInnerJoin=true \ + --conf spark.sql.optimizer.pushdownAggregateBelowJoin=true \ + --conf spark.sql.optimizer.inferDistinctFromIntersect=true \ + --conf spark.sql.optimizer.groupSplitsByLocation=false \ + --conf spark.sql.adaptive.amend.join.selection.enabled=true \ + --conf spark.sql.mergeScalaSubquery.pullupAggFilter=true \ + --conf spark.sql.execution.optimizeExpand=true \ + --conf spark.sql.execution.optimizeExpand.ratio=5 \ + --conf spark.sql.legacy.ctePrecedencePolicy=LEGACY \ + --conf spark.sql.auto.reused.cte.enabled=true \ + --conf spark.sql.auto.clear.cte.cache.enabled=true \ + --conf spark.locality.wait.node=0s \ + --conf spark.sql.sources.ignoreDataLocality=true \ + --conf spark.driver.host=10.89.56.103 \ + --conf spark.sql.catalog.paimon=org.apache.paimon.spark.SparkCatalog \ + --conf spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions \ + --conf spark.sql.catalog.paimon.type=hive \ + --conf spark.sql.catalog.paimon.warehouse=hdfs://slcluster01/hive_warehouse \ + --conf spark.hadoop.net.topology.script.file.name=/dev/null \ + --conf spark.sql.defaultCatalog=paimon \ + --keytab /data/work/hadoop-shopline/etc/hadoop/sljdp.keytab \ + --jars /tmp/paimon-hive-connector-3.1-1.5-SNAPSHOT.jar,/tmp/paimon-spark-3.5_2.12-1.5-SNAPSHOT.jar \ + --principal [email protected] \ + --class org.apache.amoro.producer.PaimonCompactManifest \ + /data/work/frame/amoro-0.9-SNAPSHOT/plugin/producer/amoro-producer-0.9-SNAPSHOT-jar-with-dependencies.jar \ + --catalogName paimon \ + --databaseName sl_oki_test \ + --retryTimes 3 \ + --continueOnTableFailure true +``` + +### Manifest 按单表合并 + +```bash +./spark-submit \ + --master yarn \ + --deploy-mode cluster \ + --driver-memory 8g \ + --driver-cores 4 \ + --executor-memory 4g \ + --executor-cores 4 \ + --conf spark.dynamicAllocation.enabled=true \ + --conf spark.dynamicAllocation.minExecutors=1 \ + --conf spark.dynamicAllocation.maxExecutors=2 \ + --conf spark.dynamicAllocation.initialExecutors=1 \ + --conf spark.dynamicAllocation.shuffleTracking.enabled=true \ + --conf spark.driver.memoryOverhead=2g \ + --conf spark.executor.memoryOverhead=2g \ + --conf spark.driver.bindAddress=0.0.0.0 \ + --conf spark.sql.shuffle.partitions=400 \ + --conf spark.sql.cbo.enabled=true \ + --conf spark.sql.files.maxPartitionBytes=512m \ + --conf spark.default.parallelism=200 \ + --conf spark.locality.wait=0s \ + --conf spark.sql.crossJoin.enabled=true \ + --conf spark.sql.adaptive.localShuffleReader.enabled=true \ + --conf spark.sql.fuse.unionAllOnJoin.enabled=true \ + --conf spark.sql.optimizer.runtime.bloomFilter.enabled=true \ + --conf spark.sql.optimizer.enableMergeScalarAggsInInnerJoin=true \ + --conf spark.sql.optimizer.pushdownAggregateBelowJoin=true \ + --conf spark.sql.optimizer.inferDistinctFromIntersect=true \ + --conf spark.sql.optimizer.groupSplitsByLocation=false \ + --conf spark.sql.adaptive.amend.join.selection.enabled=true \ + --conf spark.sql.mergeScalaSubquery.pullupAggFilter=true \ + --conf spark.sql.execution.optimizeExpand=true \ + --conf spark.sql.execution.optimizeExpand.ratio=5 \ + --conf spark.sql.legacy.ctePrecedencePolicy=LEGACY \ + --conf spark.sql.auto.reused.cte.enabled=true \ + --conf spark.sql.auto.clear.cte.cache.enabled=true \ + --conf spark.locality.wait.node=0s \ + --conf spark.sql.sources.ignoreDataLocality=true \ + --conf spark.driver.host=10.89.56.103 \ + --conf spark.sql.catalog.paimon=org.apache.paimon.spark.SparkCatalog \ + --conf spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions \ + --conf spark.sql.catalog.paimon.type=hive \ + --conf spark.sql.catalog.paimon.warehouse=hdfs://slcluster01/hive_warehouse \ + --conf spark.hadoop.net.topology.script.file.name=/dev/null \ + --conf spark.sql.defaultCatalog=paimon \ + --keytab /data/work/hadoop-shopline/etc/hadoop/sljdp.keytab \ + --jars /tmp/paimon-hive-connector-3.1-1.5-SNAPSHOT.jar,/tmp/paimon-spark-3.5_2.12-1.5-SNAPSHOT.jar \ + --principal [email protected] \ + --class org.apache.amoro.producer.PaimonCompactManifest \ + /data/work/frame/amoro-0.9-SNAPSHOT/plugin/producer/amoro-producer-0.9-SNAPSHOT-jar-with-dependencies.jar \ + --catalogName paimon \ + --databaseName sl_oki_test \ + --tableName t_order \ + --retryTimes 3 \ + --continueOnTableFailure true +``` + ## 运行注意事项 Spark 运行环境必须提前配置好 Paimon catalog,否则 `SHOW TABLES` 或 -`CALL sys.expire_snapshots`、`CALL <catalog>.sys.remove_orphan_files` 会在运行期失败。 +`CALL sys.expire_snapshots`、`CALL <catalog>.sys.remove_orphan_files`、 +`CALL <catalog>.sys.compact_manifest` 会在运行期失败。 如果运行环境已经通过 Spark 安装目录或平台层提供 Paimon 相关 JAR,可以直接提交 `amoro-producer-0.9-SNAPSHOT.jar`;如果希望单 JAR 提交,使用 diff --git a/amoro-producer/src/main/scala/org/apache/amoro/producer/PaimonCompactManifest.scala b/amoro-producer/src/main/scala/org/apache/amoro/producer/PaimonCompactManifest.scala new file mode 100644 index 000000000..3d7e82a4b --- /dev/null +++ b/amoro-producer/src/main/scala/org/apache/amoro/producer/PaimonCompactManifest.scala @@ -0,0 +1,424 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.amoro.producer + +import java.util.Locale + +import scala.util.control.NonFatal + +import org.apache.spark.sql.{Row, SparkSession} + +object PaimonCompactManifest { + + private val DefaultCatalogName = "paimon" + private val DefaultRetryTimes = 3 + private val DefaultContinueOnTableFailure = true + + final private case class AppConfig( + catalogName: String, + databaseName: Option[String], + tableName: Option[String], + retryTimes: Int, + continueOnTableFailure: Boolean, + procedureOptions: String) + + final private case class TargetTable(raw: String, database: String, table: String) + + final private case class CompactResult(result: Boolean, costMs: Double) + + def main(args: Array[String]): Unit = { + val config = + try { + val parsed = parseArgs(args) + validateConfig(parsed) + parsed + } catch { + case e: IllegalArgumentException => + printUsage() + println(s"参数错误: ${e.getMessage}") + sys.exit(2) + } + + val spark = SparkSession.builder().getOrCreate() + try { + val targets = resolveTargetTables(spark, config) + val modeDesc = + if (targets.size > 1 || config.tableName.isEmpty) { + "整库 Manifest 合并" + } else { + "单表 Manifest 合并" + } + + if (targets.nonEmpty) { + if (targets.size > 1 || config.tableName.isEmpty) { + val databaseName = config.databaseName.getOrElse(targets.head.database) + println( + s"${modeDesc}启动,catalog=${config.catalogName}, 数据库 $databaseName 共发现 ${targets.size} 张表:") + targets.foreach(t => println(s" - ${t.raw}")) + } else { + println(s"${modeDesc}启动,catalog=${config.catalogName}, 目标表: ${targets.head.raw}") + } + } + + var successTables = 0 + var failedTables = 0 + val allStartNs = System.nanoTime() + + targets.foreach { table => + try { + compactOneTable(spark, config, table) + successTables += 1 + } catch { + case NonFatal(t) => + failedTables += 1 + println(s"[表 ${table.raw} Manifest 合并失败] ${t.getClass.getSimpleName}: ${t.getMessage}") + if (!config.continueOnTableFailure) { + throw t + } + } + } + + val allCostMs = (System.nanoTime() - allStartNs) / 1e6 + println( + f"[Manifest 合并汇总] catalog=${config.catalogName}, 模式=$modeDesc, " + + f"successTables=$successTables, failedTables=$failedTables, " + + f"总耗时=${allCostMs}%.3f ms, 表数=${targets.size}") + } finally { + spark.stop() + } + } + + private def resolveTargetTables(spark: SparkSession, config: AppConfig): Seq[TargetTable] = { + (config.databaseName, config.tableName) match { + case (Some(database), Some(table)) => + Seq(makeTarget( + requireNonEmpty(database, "--databaseName"), + requireNonEmpty(table, "--tableName"))) + + case (Some(database), None) => + listTables(spark, config.catalogName, requireNonEmpty(database, "--databaseName")) + + case (None, Some(table)) => + val trimTable = requireNonEmpty(table, "--tableName") + val parts = trimTable.split("\\.", -1) + if (parts.length != 2 || parts.exists(_.trim.isEmpty)) { + throw new IllegalArgumentException( + "只传入 --tableName 时,必须使用 db.table 格式,例如 --tableName db.table") + } + Seq(makeTarget(parts(0).trim, parts(1).trim)) + + case (None, None) => + throw new IllegalArgumentException("参数不足:请至少提供 --databaseName 或 --tableName") + } + } + + private def listTables( + spark: SparkSession, + catalogName: String, + database: String): Seq[TargetTable] = { + val showTablesSql = s"SHOW TABLES IN ${quoteIdent(catalogName)}.${quoteIdent(database)}" + println(s"listing tables: $showTablesSql") + + val rows = spark.sql(showTablesSql).collect() + rows.flatMap { row => + val tableNameOpt = fieldOpt(row, Set("tablename", "table")) + .orElse(valueAt(row, 1)) + .filter(_ != null) + .map(_.toString.trim) + + val isTemporary = fieldOpt(row, Set("istemporary", "temporary")).exists { + case b: java.lang.Boolean => b.booleanValue() + case s: String => s.equalsIgnoreCase("true") || s == "1" + case i: java.lang.Integer => i == 1 + case l: java.lang.Long => l == 1L + case _ => false + } + + tableNameOpt + .filter(name => name.nonEmpty && !isTemporary) + .map(name => makeTarget(database, name)) + }.toSeq + } + + private def compactOneTable( + spark: SparkSession, + config: AppConfig, + table: TargetTable): CompactResult = { + println(s"========== 表 ${table.raw} 开始 Manifest 合并 ==========") + val startNs = System.nanoTime() + val result = runOnceWithRetry(spark, config, table) + val costMs = (System.nanoTime() - startNs) / 1e6 + + println( + f"[表 ${table.raw} Manifest 合并结束] " + + f"耗时=${costMs}%.3f ms, result=$result") + + CompactResult(result, costMs) + } + + private def runOnceWithRetry( + spark: SparkSession, + config: AppConfig, + table: TargetTable): Boolean = { + var attempt = 1 + var lastError: Throwable = null + + while (attempt <= config.retryTimes) { + try { + println(s"[表 ${table.raw}] 第 $attempt 次执行 compact_manifest") + return runCompactManifest(spark, config, table) + } catch { + case NonFatal(t) => + lastError = t + println( + s"[表 ${table.raw}] 第 $attempt 次执行失败: ${t.getClass.getSimpleName}: ${t.getMessage}") + attempt += 1 + } + } + + throw lastError + } + + private def runCompactManifest( + spark: SparkSession, + config: AppConfig, + table: TargetTable): Boolean = { + val sqlText = buildCompactManifestSql(config, table) + println(s"executing compact_manifest, table = ${table.raw}") + println(sqlText) + + val resultRows = spark.sql(sqlText).collect() + val result = resultRows.headOption + .flatMap(row => fieldOpt(row, Set("result")).orElse(valueAt(row, 0))) + .exists(toBoolean) + + println(s"SUCCESS table = ${table.raw}, result = $result") + result + } + + private def buildCompactManifestSql(config: AppConfig, table: TargetTable): String = { + val optionsClause = + if (config.procedureOptions.trim.isEmpty) { + "" + } else { + s",\n options => ${sqlString(config.procedureOptions)}" + } + + s""" + |CALL ${quoteIdent(config.catalogName)}.sys.compact_manifest( + | table => ${sqlString(table.raw)}$optionsClause + |) + |""".stripMargin + } + + private def makeTarget(database: String, table: String): TargetTable = { + TargetTable(s"${database}.${table}", database, table) + } + + private def validateConfig(config: AppConfig): Unit = { + requireNonEmpty(config.catalogName, "--catalogName") + + (config.databaseName, config.tableName) match { + case (Some(database), Some(table)) => + requireNonEmpty(database, "--databaseName") + val trimTable = requireNonEmpty(table, "--tableName") + if (trimTable.contains(".")) { + throw new IllegalArgumentException( + "当 --databaseName 与 --tableName 同时提供时,--tableName 仅允许传入表名,不允许使用 db.table") + } + + case (Some(database), None) => + requireNonEmpty(database, "--databaseName") + + case (None, Some(table)) => + val trimTable = requireNonEmpty(table, "--tableName") + val parts = trimTable.split("\\.", -1) + if (parts.length != 2 || parts.exists(_.trim.isEmpty)) { + throw new IllegalArgumentException( + "只传入 --tableName 时,必须使用 db.table 格式,例如 --tableName db.table") + } + + case (None, None) => + throw new IllegalArgumentException("参数不足:请至少提供 --databaseName 或 --tableName") + } + + if (config.retryTimes <= 0) { + throw new IllegalArgumentException("--retryTimes 必须大于 0") + } + } + + private def parseArgs(args: Array[String]): AppConfig = { + var catalogName = DefaultCatalogName + var databaseName: Option[String] = None + var tableName: Option[String] = None + var retryTimes = DefaultRetryTimes + var continueOnTableFailure = DefaultContinueOnTableFailure + var procedureOptions = "" + + var i = 0 + while (i < args.length) { + args(i) match { + case "--help" | "-h" => + printUsage() + sys.exit(0) + + case "--catalogName" => + catalogName = readValue(args, i, "--catalogName") + i += 1 + + case "--databaseName" => + databaseName = Some(readValue(args, i, "--databaseName")) + i += 1 + + case "--tableName" => + tableName = Some(readValue(args, i, "--tableName")) + i += 1 + + case "--retryTimes" => + retryTimes = parseInt(readValue(args, i, "--retryTimes"), "--retryTimes") + i += 1 + + case "--continueOnTableFailure" => + continueOnTableFailure = + parseBoolean(readValue(args, i, "--continueOnTableFailure"), "--continueOnTableFailure") + i += 1 + + case "--procedureOptions" => + procedureOptions = readValue(args, i, "--procedureOptions") + i += 1 + + case other => + throw new IllegalArgumentException(s"未知参数: $other") + } + i += 1 + } + + AppConfig( + catalogName, + databaseName, + tableName, + retryTimes, + continueOnTableFailure, + procedureOptions) + } + + private def readValue(args: Array[String], index: Int, argName: String): String = { + if (index + 1 >= args.length || args(index + 1).startsWith("--")) { + throw new IllegalArgumentException(s"$argName 缺少参数值") + } + args(index + 1) + } + + private def parseInt(value: String, argName: String): Int = { + try { + value.toInt + } catch { + case _: NumberFormatException => + throw new IllegalArgumentException(s"$argName 必须是整数") + } + } + + private def parseBoolean(value: String, argName: String): Boolean = { + value.trim.toLowerCase(Locale.ROOT) match { + case "true" => true + case "false" => false + case _ => throw new IllegalArgumentException(s"$argName 必须是 true 或 false") + } + } + + private def requireNonEmpty(value: String, argName: String): String = { + val trimmed = value.trim + if (trimmed.isEmpty) { + throw new IllegalArgumentException(s"$argName 的值不能为空") + } + trimmed + } + + private def quoteIdent(name: String): String = { + "`" + name.replace("`", "``") + "`" + } + + private def sqlString(value: String): String = { + "'" + value.replace("'", "''") + "'" + } + + private def fieldOpt(row: Row, normalizedNames: Set[String]): Option[Any] = { + val schemaOpt = + try { + Option(row.schema) + } catch { + case NonFatal(_) => None + } + + schemaOpt.flatMap { schema => + schema.fields.zipWithIndex + .find { case (field, _) => normalizedNames.contains(normalizeFieldName(field.name)) } + .flatMap { case (_, index) => valueAt(row, index) } + } + } + + private def valueAt(row: Row, index: Int): Option[Any] = { + if (index >= 0 && index < row.size && !row.isNullAt(index)) { + Some(row.get(index)) + } else { + None + } + } + + private def normalizeFieldName(name: String): String = { + name.toLowerCase(Locale.ROOT).replace("_", "") + } + + private def toBoolean(value: Any): Boolean = { + value match { + case null => false + case b: java.lang.Boolean => b.booleanValue() + case s: String => s.trim.equalsIgnoreCase("true") || s.trim == "1" + case i: java.lang.Integer => i == 1 + case l: java.lang.Long => l == 1L + case other => other.toString.trim.equalsIgnoreCase("true") + } + } + + private def printUsage(): Unit = { + println( + s""" + |Usage: + | spark-submit --class org.apache.amoro.producer.PaimonCompactManifest <jar> [options] + | + |Required target options: + | --databaseName <db> 合并指定库下全部非临时表 Manifest + | --databaseName <db> --tableName <t> 只合并指定库下的单表 Manifest + | --tableName <db.table> 只合并全限定单表 Manifest + | + |Optional options: + | --catalogName <catalog> Paimon catalog 名称,默认: $DefaultCatalogName + | --procedureOptions <options> compact_manifest options,默认不传 + | --retryTimes <num> 单表失败最大执行次数,默认: $DefaultRetryTimes + | --continueOnTableFailure <true|false> + | 单表失败后是否继续后续表,默认: $DefaultContinueOnTableFailure + | --help, -h 打印帮助信息 + | + |Examples: + | --databaseName sl_oki_test + | --databaseName sl_oki_test --tableName t_order + | --tableName sl_oki_test.t_order + |""".stripMargin) + } +}
