使用 Delta Lake UniForm 将表转换为 Hudi 格式:完整实战指南 📅 发布时间:2026/9/16 12:08:46 👁 浏览次数: 使用 Delta Lake UniForm 将表转换为 Hudi 格式完整实战指南【免费下载链接】deltaAn open-source storage framework that enables building a Lakehouse architecture with compute engines including Spark, PrestoDB, Flink, Trino, and Hive and APIs项目地址: https://gitcode.com/GitHub_Trending/del/delta本文以 Delta Lake 仓库中hudi/模块为主线讲解如何通过 UniFormUniversal Format机制让同一份 Parquet 数据文件同时具备 Delta 与 Hudi 两套元数据从而被 Hudi 生态的 Spark 版本直接读取。你将掌握如何用 spark-sql 创建启用 Hudi UniForm 的 Delta 表并写入数据、如何在低版本 Spark 中启动 Hudi Spark Shell 读取该表、底层异步元数据转换的执行原理以及 schema 转换、类型支持、时间旅行与清理归档等进阶行为。背景UniForm 与 Hudi 的双格式数据Delta Lake、Apache Hudi 与 Apache Iceberg 的底层数据文件都是 Parquet差异只在于各自的元数据层。Delta UniForm 正是利用这一点在 Delta 提交之后异步地为同一份 Parquet 数据文件生成 Iceberg / Hudi 的元数据让 Iceberg、Hudi 客户端可以把 Delta 表当作自己的表来读取。仓库中的 hudi/README.md 给出了从创建到读取的最短路径而hudi/src/main/scala/org/apache/spark/sql/delta/hudi/目录下的四个 Scala 文件HudiConverter、HudiConversionTransaction、HudiSchemaUtils、HudiTransactionUtils则是这套能力在 Spark 侧的完整实现配合 ConvertToHudiSuite.scala 与 write_uniform_hudi.py 集成测试可以验证整个读写闭环。需要特别说明的是UniForm 从 Hudi 视角是只读的官方文档docs/src/content/docs/delta-uniform.mdx明确提示外部 Hudi 客户端只能读不能写任何非 Delta 的外部写入者都可能破坏 Delta 表并导致数据丢失。一、创建启用 Hudi UniForm 的 Delta 表1.1 准备 Spark 环境与依赖 JarUniForm 的 Hudi 支持是一个独立模块需要把对应的 assembly Jar 放到 Spark 的 classpath 上。hudi/README.md给出的启动命令如下spark-sql --packages io.delta:delta-spark_2.12:3.2.0-SNAPSHOT \ --jars delta-hudi-assembly_2.12-3.2.0-SNAPSHOT.jar \ --conf spark.sql.extensionsio.delta.sql.DeltaSparkSessionExtension \ --conf spark.sql.catalog.spark_catalogorg.apache.spark.sql.delta.catalog.DeltaCatalog要点说明--packages拉取 Delta Spark 连接器本身--jars指定本地的delta-hudi-assembly产物在该仓库中由hudi/模块构建而来。spark.sql.extensions与spark.sql.catalog.spark_catalog是 Delta 接入 Spark 的标准配置分别注册 SQL 扩展与 Catalog 实现。官方文档delta-uniform.mdx也提供了通过 Maven 坐标直接引用的方式--packages io.delta:io.delta:delta-hudi_2.12:version。写入侧要求 Delta Lake 3.2 及以上版本且 UniForm 的 Hudi 支持目前仍处于preview预览状态。1.2 建表并写入数据进入 spark-sql 交互界面后创建一张启用 Hudi UniForm 的 Delta 表CREATE TABLE delta_table_with_hudi (col1 INT) USING DELTA TBLPROPERTIES(delta.universalFormat.enabledFormats hudi) LOCATION /tmp/delta-table-with-hudi;然后插入记录INSERT INTO delta_table_with_hudi VALUES (1);这里起关键作用的是表属性delta.universalFormat.enabledFormats hudi。它等价于源码中的DeltaConfigs.UNIVERSAL_FORMAT_ENABLED_FORMATS由 UniversalFormat.scala 中的hudiEnabled(metadata)判定当前快照是否启用了 Hudi 转换。除了在建表时一次性声明也可以在表创建后再开启——集成测试 write_uniform_hudi.py 验证了这种场景CREATE TABLE delta_table_with_hudi_1 (col1 INT, col2 STRING) USING DELTA; INSERT INTO delta_table_with_hudi_1 VALUES (1, a), (2, b); ALTER TABLE delta_table_with_hudi_1 SET TBLPROPERTIES(delta.universalFormat.enabledFormats hudi);注意启用 UniForm 需要表已开启列映射column mapping。在建表时 Delta 会自动设置而对存量表执行上述ALTER TABLE前需确认列映射已启用。1.3 建表过程中的 Hudi 元数据初始化当hudiEnabled判定通过后转换逻辑会走 HudiTransactionUtils.scala 的loadTableMetaClient如果目标路径下还没有 Hudi 表就调用initializeHudiTable初始化一个全新的 Hudi 表。从源码可以看到初始化时固定写入的关键属性表类型HoodieTableType.COPY_ON_WRITE写时复制时间语义setCommitTimezone(HoodieTimelineTimeZone.UTC)即 Hudi 提交时间统一使用 UTC这保证了 Delta 提交时间戳与 Hudi instant 时间可对齐分区风格setHiveStylePartitioningEnable(true)与读取侧hoodie.datasource.write.hive_style_partitioningtrue对应KeyGenerator 选择无分区用NonpartitionedKeyGenerator单分区用SimpleKeyGenerator多分区用CustomKeyGenerator关闭 populateMetaFieldssetPopulateMetaFields(false)避免在 Parquet 文件中写入 Hudi 的_hoodie_*元数据列。这也解释了为什么 Hudi 侧读取时要显式开启元数据表Hudi 的元数据timeline、commit 等全部以文件形式存放在同一路径下而非依赖 Parquet 内的内联字段。二、用 Hudi 读取这张表2.1 版本约束Spark 3.4 及以下Hudi 目前不支持 Spark 3.5.x因此读取侧必须启动一个 Spark 3.4 或更早版本的 Spark Shell。按照 Hudi 官方的 Quick Start Guide 启动 spark-shell 后用如下 Scala 代码读取val df spark.read.format(hudi) .option(hoodie.metadata.enable, true) .load(/tmp/delta-table-with-hudi)两个关键点format(hudi)告诉 Spark 使用 Hudi 数据源hoodie.metadata.enable true启用 Hudi 元数据表。由于 UniForm 生成的 Hudi 表在初始化时关闭了 populateMetaFields读取必须依赖 Hudi 的文件系统元数据timeline metadata table来确定有效文件列表。集成测试 write_uniform_hudi.py 中Hudi 读取会话还需要额外的 Kryo 序列化配置SparkSession.builder \ .appName(delta-uniform-hudi-reader) \ .config(spark.sql.extensions, org.apache.spark.sql.hudi.HoodieSparkSessionExtension) \ .config(spark.sql.catalog.spark_catalog, org.apache.spark.sql.hudi.catalog.HoodieCatalog) \ .config(spark.serializer, org.apache.spark.serializer.KryoSerializer) \ .config(spark.kryo.registrator, org.apache.spark.HoodieSparkKryoRegistrar) \ .getOrCreate()2.2 数据一致性验证集成测试对每一张测试表都会同时用 Delta 和 Hudi 两种数据源读取并通过assertDataFrameEqual(df_delta, df_hudi)断言两份 DataFrame 完全一致覆盖了以下场景建表后立即转换_0表包含 BIGINT、BOOLEAN、DATE、DOUBLE、FLOAT、INT、STRING、TIMESTAMP、BINARY、DECIMAL、STRUCT、ARRAY、MAP 等全类型建表后再通过ALTER TABLE开启 UniForm_1表执行DELETE后的数据_2表顶层 schema 演进ALTER TABLE ADD COLUMN col3 INT FIRST_3表嵌套字段演进ALTER TABLE ... ADD COLUMN col1.field3 INT AFTER field1_4表时间旅行读取_5表同时开启了delta.columnMapping.modename。三、底层原理异步 Hudi 元数据转换3.1 Post-commit 钩子与异步线程每一次 Delta 提交完成后HudiConverterHook.scala 作为 post-commit hook 被触发。它首先校验两点txn.committedVersion postCommitSnapshot.version避免对同一批 action 重复转换以及快照元数据中hudiEnabled为真。随后根据配置决定同步还是异步执行DELTA_UNIFORM_HUDI_SYNC_CONVERT_ENABLED异步模式默认调用converter.enqueueSnapshotForConversion把快照放入待转换队列同步模式直接调用converter.convertSnapshot阻塞完成转换。异步的核心实现在 HudiConverter.scala 的enqueueSnapshotForConversion中转换工作跑在一个名为async-hudi-converter的守护线程里队列中始终只保留一个待转换快照——如果转换尚未完成又有新提交旧快照会被新快照替换同时记录delta.hudi.conversion.async.backlog事件。这样既能以异步方式消除对 Delta 写入延迟的影响又不会让待转换积压无限增长。3.2 增量转换只处理上次转换以来的提交convertSnapshot的增量设计值得关注。转换前先从 Hudi meta client 读取最近一次已转换的 Delta 版本loadLastDeltaVersionConverted逻辑是读取 Hudi timeline 上最后一个已完成 commit 的 extra metadata其中保存了HudiConverter.DELTA_VERSION_PROPERTYkey 为delta-version与DELTA_TIMESTAMP_PROPERTYkey 为delta-timestamp两个属性分别记录对应的 Delta commit 版本号与毫秒时间戳。增量转换有三种路径最近转换版本 当前事务读取快照版本直接复用事务的 read snapshot避免重新加载距离上次转换的提交数在阈值内由配置hudi.maxPendingCommits控制源码为DeltaSQLConf.HUDI_MAX_COMMITS_TO_CONVERT调用log.getSnapshotAt(version)拿到上次转换时的快照然后通过DeltaFileProviderUtils.getCommitsInVersionRange只读取该版本之后到当前版本之间的 Delta JSON 文件按批batch转换为 Hudi action距离过远或 commit 文件已过期DeltaFileNotFoundException退化为全量状态重建直接枚举当前快照的allFiles进行转换。这种设计把每次转换的开销控制在上次转换以来的增量上并避免 driver 因加载过长历史而 OOM。3.3 Hudi 提交事务与清理归档每个批次的 action 通过 HudiConversionTransaction.scala 提交为一次 Hudi 事务setCommitFileUpdates把 Delta 的AddFile转成 Hudi 的WriteStatus含文件路径、分区、记录数、字节数把RemoveFile按分区归类为被替换的文件组commit()使用HoodieJavaWriteClient以REPLACE_COMMIT_ACTION提交并把delta-version/delta-timestamp写入 Hudi commit 的 extra metadata供后续增量转换与版本对应关系查询使用提交后还会执行清理cleaning与归档archival按KEEP_LATEST_BY_HOURS策略默认保留 7×24 小时内的 commit标记过期文件手动触发HoodieTimelineArchiver归档旧 commit从而控制 Hudi timeline 与元数据表的规模增长。单元测试ConvertToHudiSuite中专门有validate Hudi timeline archival and cleaning用例使用ManualClock模拟 12 天前开始、连续 20 次提交的场景验证 Hudi 侧 active/archived timeline 的 commit 数量符合预期verifyNumHudiCommits则断言activeCommits archivedCommits与预期一致。3.4 Schema 转换Delta 类型到 Hudi/Avro 类型Hudi 以 Avro schema 描述表结构。HudiSchemaUtils.scala 的convertDeltaSchemaToHudiSchema递归地把 Delta 的StructType转为 Avro schemaDelta 类型Avro 类型StringTypeSTRINGLongTypeLONGIntegerTypeINTFloatType/DoubleTypeFLOAT/DOUBLEDecimalType(p, s)BYTES decimal logical typeBooleanTypeBOOLEANBinaryTypeBYTESDateTypeINT date logical typeTimestampTypeLONG timestamp-micros logical typeStructType/ArrayType/MapTypeAvro record / array / map递归nullable 字段包装为 union测试ConvertToHudiSuite对 ARRAY、ARRAY 、嵌套 ARRAY、MAP、多层嵌套 STRUCT 等类型均有专门用例验证转换后的 Hudi 表 schema 与 Delta 原始 schema 完全一致verifyFilesAndSchemaMatch同时比对文件列表与 schema。四、约束与注意事项结合源码与官方文档delta-uniform.mdx使用 UniForm Hudi 时有以下约束需要留意Deletion Vectors 与 UniForm Hudi 互斥。ConvertToHudiSuite中的两个用例明确验证在建表属性中同时声明delta.enableDeletionVectorstrue与delta.universalFormat.enabledFormatshudi会抛出DeltaUnsupportedOperationException在已启用 Hudi UniForm 的表上再开启 Deletion Vectors 同样被拒绝。读取侧版本限制。Hudi 客户端尚不支持 Spark 3.5.x读取时需使用 Spark 3.4 及以下版本当前仓库环境下写入侧要求 Delta 3.2。只读互操作。UniForm 生成的 Hudi 元数据只服务于读取外部 Hudi 引擎对表执行写入、清理等操作可能破坏 Delta 表必须避免。首次启用后的异步等待。首次启用 UniForm 时异步元数据生成任务需要一定时间完成之后外部客户端才能读到完整数据可借助 Hudi 表元数据中delta-version/delta-timestamp属性核对 Delta 与 Hudi 的版本对应关系。版本对齐语义。Delta 与 Hudi 的 commit时间戳是对齐的都基于同一 Delta 提交时间Hudi 侧固定使用 UTC但版本号并不一一对应——高提交频率下多个 Delta commit 可能被打包进同一个 Hudi commit这也正是hudi.maxPendingCommits阈值存在的意义。五、快速验证清单如果你想在当前仓库环境下复现整条链路可以按如下顺序操作构建或获取delta-hudi-assembly产物按 1.1 节命令启动 spark-sqlDelta 3.2执行 1.2 节的建表与INSERT或对已存在的 Delta 表执行ALTER TABLE ... SET TBLPROPERTIES等待片刻让异步转换完成观察 Hudi 路径下出现.hoodie目录与 timeline 文件启动 Spark 3.4 及以下的 Hudi spark-shell用 2.1 节的 Scala 代码读取同一路径比对数据若需自动化验证可参考 write_uniform_hudi.py 的完整流程全类型、删除、schema 演进、时间旅行以及 ConvertToHudiSuite.scala 中verifyFilesAndSchemaMatch的文件与 schema 双维度断言方法。至此你已完成Delta 写入 → Hudi 读取的 UniForm 全流程并理解了其背后异步转换、增量提交、schema 映射与清理归档的实现细节。【免费下载链接】deltaAn open-source storage framework that enables building a Lakehouse architecture with compute engines including Spark, PrestoDB, Flink, Trino, and Hive and APIs项目地址: https://gitcode.com/GitHub_Trending/del/delta创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考