在 Apache Spark 上运行 XTable 同步
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:VALUEshell
$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 | 以点号分隔的表命名空间。 |
--partitionspec | Hudi 源的分区字段定义,例如 level:VALUE。 |
--usedeltakernel | 对 Delta 源或目标强制使用 Delta Kernel。 |
--help | 打印用法说明文本。 |
支持的格式
Paimon 和 Parquet 只能作为只读源;XTable 不会将这两种格式写为目标。
| 源 ↓ / 目标 → | Hudi | Iceberg | Delta |
|---|---|---|---|
| Hudi | – | ✅ | ✅ |
| Iceberg | ✅ | – | ✅ |
| Delta | ✅ | ✅ | – |
| Paimon | ✅ | ✅ | ✅ |
| Parquet | ✅ | ✅ | ✅ |
Spark 版本支持
与 Hudi 和 Iceberg 之间的转换完全不需要 Spark,因此它们可以在下列任意 Spark 版本上运行。Delta 是唯一实现依赖 Spark 版本的引擎,jar 会自动选择正确的实现:
| Spark 版本 | Hudi 和 Iceberg | Delta 实现 |
|---|---|---|
| 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,了解查询已同步表时各格式所需的选项。
- 如果你更希望从源码构建项目,请参阅安装。
评论
登录后参与评论
KnowForge