快速开始

在 Apache Spark 上运行 XTable 同步

师成师成· 更新于 2026-09-28· 阅读 11 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

xtable-spark-runtime 是一个运行时 jar,用于在 Apache Spark 集群上运行 Apache XTable™ (Incubating) 同步。与任何 XTable 同步一样,同步不会重写任何数据文件——它读取源表的元数据,并将目标格式的元数据写入到已有数据的旁边。

这个 jar 旨在直接放入你已经在运行的 Spark 作业中,因此你不需要额外的进程或额外的集群,就能让同一张表在不同格式之间保持互操作性。

获取 jar

该 jar 已发布到 Maven Central,坐标为 org.apache.xtable:xtable-spark-runtime_2.12:0.4.0-incubating。下载一次即可:

shell

curl -O https://repo1.maven.org/maven2/org/apache/xtable/xtable-spark-runtime_2.12/0.4.0-incubating/xtable-spark-runtime_2.12-0.4.0-incubating.jar

每个引擎依赖都是 provided 的,因此这个 jar 只有约 4 MB,并且会复用你的集群中已有的 Hudi、Iceberg 和 Delta 库。完整发布列表请参见 Downloads 页面。

在 Spark 作业中添加同步

这正是该 jar 的用途。你的作业已经以某种格式写入了一张表,你希望在这次写入完成后,该表能立即与其他格式互操作。

将该 jar 加入到你已经在使用的 spark-submit 命令中,与你的应用 jar 放在一起:

shell

$SPARK_HOME/bin/spark-submit \
  --jars xtable-spark-runtime_2.12-0.4.0-incubating.jar \
  --class com.example.OrdersJob \
  orders-job.jar

然后在写入完成之后调用 XTableSyncService。用 TableSyncSpec 描述表结构,并把会话的 Hadoop 配置传给它:

  • Scala
  • Java

OrdersJob.scala

import java.util.{Arrays => JArrays}
import org.apache.xtable.spark.{TableSyncSpec, XTableSyncService}

val basePath = "s3://example-warehouse/db/orders"

// the write your job already does
df.write.format("hudi").options(hudiOptions).mode("append").save(basePath)

// the one call you add
new XTableSyncService().sync(
  TableSyncSpec.builder()
    .key("orders")
    .basePath(basePath)
    .sourceFormat("HUDI")
    .targets(JArrays.asList("ICEBERG", "DELTA"))
    .build(),
  spark.sparkContext.hadoopConfiguration)

要针对这些类进行编译,请在构建中添加相同的 Maven 坐标,并将其作用域(scope)设为 provided。

同步以增量方式运行,并在目标的同步元数据中记录自身的水位线(watermark);当增量同步不安全时——例如首次运行——会自动回退为全量快照。因此,在每次写入之后调用同步都是安全的。

表通过路径来标识

运行时 jar 通过路径而非目录(catalog)来标识表,因此目前还没有与 RunSync 的 Iceberg 目录配置(-i)等价的选项。当表的数据文件不直接位于 basePath 下时,请将 dataPath 设置为其实际所在位置,因为各目标就是将元数据写入到该位置的。Iceberg 表的数据通常位于 <basePath>/data,但 Iceberg 并不强制要求这种布局——write.data.path 以及对象存储的布局可以把数据文件放在任意位置——因此请确认你的表实际写入的位置。

可运行示例

demo/spark-runtime 是一个完整的作业示例,可以进行双向同步并校验行数。

将同步作为独立作业运行

该 jar 还提供了一个 spark-submit 入口点,相当于此组件包中的 RunSync。请将该 jar 作为应用 jar 传入,而不是使用 --jars:

shell

$SPARK_HOME/bin/spark-submit \
  --class org.apache.xtable.spark.XTableSparkSync \
  --master 'local[*]' \
  xtable-spark-runtime_2.12-0.4.0-incubating.jar \
  --basepath /path/to/hudi_table \
  --sourceformat HUDI \
  --targets ICEBERG,DELTA

要在一次提交中同步多张表,请使用 --datasetconfig。它接受与 RunSync 相同的 YAML 配置,因此现有配置无需修改即可直接使用;而且与 RunSync 不同,配置文件本身也可以存放在云存储上:

my_config.yaml

sourceFormat: HUDI
targetFormats:
  - DELTA
  - ICEBERG
datasets:
  -
    tableBasePath: s3://tpc-ds-datasets/1GB/hudi/call_center
    tableDataPath: s3://tpc-ds-datasets/1GB/hudi/call_center/data
    tableName: call_center
    namespace: my.db
  -
    tableBasePath: s3://tpc-ds-datasets/1GB/hudi/catalog_sales
    tableName: catalog_sales
    partitionSpec: cs_sold_date_sk:VALUE
  -
    tableBasePath: s3://hudi/multi-partition-dataset
    tableName: multi_partition_dataset
    partitionSpec: time_millis:DAY:yyyy-MM-dd,type:VALUE

shell

$SPARK_HOME/bin/spark-submit \
  --class org.apache.xtable.spark.XTableSparkSync \
  xtable-spark-runtime_2.12-0.4.0-incubating.jar \
  --datasetconfig my_config.yaml

--datasetconfig 与 --basepath 互斥——两者只能传一个。

选项说明
--basepath源表的根路径。
--sourceformat源格式:HUDI、ICEBERG、DELTA、PAIMON 或 PARQUET。
--targets以逗号分隔的目标格式,例如 ICEBERG,DELTA。
--datasetconfig列出多张表的 YAML 配置文件路径。可以是本地路径或云存储路径。
--datapath数据文件的路径,当其与根路径不同时使用。
--tablename表名。默认取根路径的最后一段。
--namespace以点号分隔的表命名空间。
--partitionspecHudi 源的分区字段定义,例如 level:VALUE。
--usedeltakernel对 Delta 源或目标强制使用 Delta Kernel。
--help打印用法说明文本。

支持的格式

Paimon 和 Parquet 只能作为只读源;XTable 不会将这两种格式写为目标。

源 ↓ / 目标 →HudiIcebergDelta
Hudi–✅✅
Iceberg✅–✅
Delta✅✅–
Paimon✅✅✅
Parquet✅✅✅

Spark 版本支持

与 Hudi 和 Iceberg 之间的转换完全不需要 Spark,因此它们可以在下列任意 Spark 版本上运行。Delta 是唯一实现依赖 Spark 版本的引擎,jar 会自动选择正确的实现:

Spark 版本Hudi 和 IcebergDelta 实现
3.4.x✅Delta Standalone
3.5.x 及更高版本✅Delta Kernel,自动选择

Delta Standalone 无法在 Spark 3.5 上运行,因此在 3.5 及更高版本中,Delta 源或目标会自动通过 Delta Kernel 处理,无需额外传参。如果你也想在 Spark 3.4 上使用 Kernel,请传入 --usedeltakernel,或在 TableSyncSpec 上设置 .useDeltaKernel(true)。

下一步

  • 参阅快速入门,了解端到端的互操作性演练。
  • 参阅 Apache Spark,了解查询已同步表时各格式所需的选项。
  • 如果你更希望从源码构建项目,请参阅安装。

评论

登录后参与评论

正在加载评论…