运维 Hudi

部署

师成师成· 更新于 2026-09-29· 阅读 30 分钟· 0 次阅读

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

本节提供大规模部署和运维 Hudi 表所需的全部帮助。具体而言,我们将涵盖以下方面。

部署

总而言之,Hudi 的部署不需要常驻服务器,也不会为你的数据湖带来额外的基础设施成本。事实上,Hudi 开创了利用现有基础设施构建事务性分布式存储层的模式,令人欣慰的是其他系统也在采用类似的方法。Hudi 的写入通过 Spark 作业(Hudi Streamer 或自定义 Spark 数据源作业)完成,按照 Apache Spark 的标准建议进行部署。对 Hudi 表的查询通过安装到 Apache Hive、Apache Spark 或 PrestoDB 中的库实现,因此无需额外的基础设施。

典型的 Hudi 数据摄取可以在两种模式下完成。在单次运行模式下,Hudi 摄取会读取下一批数据,将其写入 Hudi 表后退出。在持续模式下,Hudi 摄取作为长期运行的服务运行,循环执行摄取作业。

对于 Merge_On_Read 表,Hudi 摄取还需要负责合并(compact)增量文件。同样,压缩可以以异步模式执行,即让压缩与摄取并发运行,也可以以串行方式依次执行。

Hudi Streamer

Hudi Streamer 是一个独立工具,用于从 DFS、Kafka 和数据库变更日志(Changelog)等不同来源增量拉取上游变更,并将其摄取到 Hudi 表中。它以 Spark 应用的形式在两种模式下运行。

要在 Spark 中使用 Hudi Streamer,需要 hudi-utilities-slim-bundle 和 Hudi Spark bundle,只需在 spark-submit 命令中添加 --packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.1,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.1 即可。

请选择与你的 Spark 运行时相匹配的 Spark bundle —— 例如 hudi-spark3.3-bundle_2.12、hudi-spark3.4-bundle_2.12、hudi-spark3.5-bundle_2.12(Scala 2.12 或 2.13)、hudi-spark3.5-bundle_2.13、hudi-spark4.0-bundle_2.13 或 hudi-spark4.1-bundle_2.13。Spark 4.0 和 4.1 在运行时需要 Java 17 或更高版本;Spark 3.x 可运行在 Java 8 或更高版本上。

  • 单次运行模式(Run Once Mode):在此模式下,Hudi Streamer 执行一轮摄取,包括从上游源增量拉取事件并将其摄取到 Hudi 表中。清理旧文件版本、归档 hoodie 时间线等后台操作会作为该次运行的一部分自动执行。对于 Merge-On-Read 表,Compaction(压缩)也会以内联方式作为摄取的一部分运行,除非通过传入标志 --disable-compaction 来禁用。默认情况下,每次摄取运行都会内联执行 Compaction,可通过设置属性 hoodie.compact.inline.max.delta.commits 来更改该行为。你可以手动运行这个 Spark 应用,也可以使用任何 cron 定时任务或工作流编排工具(最常见的部署方式),例如 Apache Airflow 来启动该应用。有关运行该 Spark 应用的命令行选项,请参见此章节。

以下是一个示例调用,用于在单次运行模式下从 Kafka 主题读取数据,并在 YARN 集群上写入 Merge-On-Read 表类型:

spark-submit \
 --packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.1,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.1 \
 --master yarn \
 --deploy-mode cluster \
 --num-executors 10 \
 --executor-memory 3g \
 --driver-memory 6g \
 --conf spark.driver.extraJavaOptions="-XX:+PrintGCApplicationStoppedTime -XX:+PrintGCApplicationConcurrentTime -XX:+PrintGCTimeStamps -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/tmp/varadarb_ds_driver.hprof" \
 --conf spark.executor.extraJavaOptions="-XX:+PrintGCApplicationStoppedTime -XX:+PrintGCApplicationConcurrentTime -XX:+PrintGCTimeStamps -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/tmp/varadarb_ds_executor.hprof" \
 --queue hadoop-platform-queue \
 --conf spark.scheduler.mode=FAIR \
 --conf spark.yarn.executor.memoryOverhead=1072 \
 --conf spark.yarn.driver.memoryOverhead=2048 \
 --conf spark.task.cpus=1 \
 --conf spark.executor.cores=1 \
 --conf spark.task.maxFailures=10 \
 --conf spark.memory.fraction=0.4 \
 --conf spark.rdd.compress=true \
 --conf spark.kryoserializer.buffer.max=200m \
 --conf spark.serializer=org.apache.spark.serializer.KryoSerializer \
 --conf spark.memory.storageFraction=0.1 \
 --conf spark.shuffle.service.enabled=true \
 --conf spark.sql.hive.convertMetastoreParquet=false \
 --conf spark.ui.port=5555 \
 --conf spark.driver.maxResultSize=3g \
 --conf spark.executor.heartbeatInterval=120s \
 --conf spark.network.timeout=600s \
 --conf spark.eventLog.overwrite=true \
 --conf spark.eventLog.enabled=true \
 --conf spark.eventLog.dir=hdfs:///user/spark/applicationHistory \
 --conf spark.yarn.max.executor.failures=10 \
 --conf spark.sql.catalogImplementation=hive \
 --conf spark.sql.shuffle.partitions=100 \
 --driver-class-path $HADOOP_CONF_DIR \
 --class org.apache.hudi.utilities.streamer.HoodieStreamer \
 --table-type MERGE_ON_READ \
 --source-class org.apache.hudi.utilities.sources.JsonKafkaSource \
 --source-ordering-field ts  \
 --target-base-path /user/hive/warehouse/stock_ticks_mor \
 --target-table stock_ticks_mor \
 --props /var/demo/config/kafka-source.properties \
 --schemaprovider-class org.apache.hudi.utilities.schema.FilebasedSchemaProvider
  • Continuous Mode:在此模式下,Hudi Streamer 会运行一个无限循环,每一轮都执行一次 Run Once Mode 中描述的摄取流程。数据摄取的频率可以通过配置项 --min-sync-interval-seconds 来控制。对于 Merge-On-Read 表,除非通过传入 --disable-compaction 标志禁用,否则 Compaction 会以异步方式与摄取并发运行。每次摄取运行都会异步触发一次 compaction 请求,该频率可以通过设置属性 hoodie.compact.inline.max.delta.commits 来调整。由于摄取和 compaction 运行在同一个 Spark context 中,因此可以使用 Hudi Streamer CLI 中的资源分配配置,例如 --delta-sync-scheduling-weight、--compact-scheduling-weight、--delta-sync-scheduling-minshare 和 --compact-scheduling-minshare,来控制执行器在摄取与 compaction 之间的分配。

以下是一个示例调用,用于以连续模式从 Kafka topic 读取数据,并写入 yarn 集群中的 Merge On Read 表类型:

spark-submit \
 --packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.1,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.1 \
 --master yarn \
 --deploy-mode cluster \
 --num-executors 10 \
 --executor-memory 3g \
 --driver-memory 6g \
 --conf spark.driver.extraJavaOptions="-XX:+PrintGCApplicationStoppedTime -XX:+PrintGCApplicationConcurrentTime -XX:+PrintGCTimeStamps -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/tmp/varadarb_ds_driver.hprof" \
 --conf spark.executor.extraJavaOptions="-XX:+PrintGCApplicationStoppedTime -XX:+PrintGCApplicationConcurrentTime -XX:+PrintGCTimeStamps -XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/tmp/varadarb_ds_executor.hprof" \
 --queue hadoop-platform-queue \
 --conf spark.scheduler.mode=FAIR \
 --conf spark.yarn.executor.memoryOverhead=1072 \
 --conf spark.yarn.driver.memoryOverhead=2048 \
 --conf spark.task.cpus=1 \
 --conf spark.executor.cores=1 \
 --conf spark.task.maxFailures=10 \
 --conf spark.memory.fraction=0.4 \
 --conf spark.rdd.compress=true \
 --conf spark.kryoserializer.buffer.max=200m \
 --conf spark.serializer=org.apache.spark.serializer.KryoSerializer \
 --conf spark.memory.storageFraction=0.1 \
 --conf spark.shuffle.service.enabled=true \
 --conf spark.sql.hive.convertMetastoreParquet=false \
 --conf spark.ui.port=5555 \
 --conf spark.driver.maxResultSize=3g \
 --conf spark.executor.heartbeatInterval=120s \
 --conf spark.network.timeout=600s \
 --conf spark.eventLog.overwrite=true \
 --conf spark.eventLog.enabled=true \
 --conf spark.eventLog.dir=hdfs:///user/spark/applicationHistory \
 --conf spark.yarn.max.executor.failures=10 \
 --conf spark.sql.catalogImplementation=hive \
 --conf spark.sql.shuffle.partitions=100 \
 --driver-class-path $HADOOP_CONF_DIR \
 --class org.apache.hudi.utilities.streamer.HoodieStreamer \
 --table-type MERGE_ON_READ \
 --source-class org.apache.hudi.utilities.sources.JsonKafkaSource \
 --source-ordering-field ts  \
 --target-base-path /user/hive/warehouse/stock_ticks_mor \
 --target-table stock_ticks_mor \
 --props /var/demo/config/kafka-source.properties \
 --schemaprovider-class org.apache.hudi.utilities.schema.FilebasedSchemaProvider \
 --continuous

Spark 数据源写入作业

如批量写入中所述,你可以使用 Spark 数据源将数据写入 Hudi 表。该机制允许你以 Hudi 格式写入任意 Spark DataFrame。Hudi Spark 数据源还支持通过 Spark Streaming 将流式数据源写入 Hudi 表。对于 Merge On Read 表类型,默认启用内联压缩(inline compaction),该操作会在每次写入任务结束后执行。压缩频率可以通过设置属性 hoodie.compact.inline.max.delta.commits 来调整。

以下是使用 Spark 数据源的调用示例:

inputDF.write()
       .format("org.apache.hudi")
       .options(clientOpts) // any of the Hudi client opts can be passed in as well
       .option("hoodie.datasource.write.recordkey.field", "_row_key")
       .option("hoodie.datasource.write.partitionpath.field", "partition")
       .option("hoodie.table.ordering.fields", "timestamp")
       .option("hoodie.table.name", tableName)
       .mode(SaveMode.Append)
       .save(basePath);

升级

新的 Hudi 版本会在发布页面上列出,其中包含详细说明,列出所有变更以及每个版本的亮点。归根结底,Hudi 是一个存储系统,随之而来的是一系列责任,我们会认真对待。

一般准则:

  • 我们力求保持所有变更向后兼容(即新代码可以读取旧的数据/时间线文件),当无法做到时,我们会通过 CLI 提供升级/降级工具。
  • 我们无法始终保证向前兼容(即旧代码能够读取更高版本写入的数据/时间线文件)。这通常是常态,因为否则无法构建任何新功能。不过,此类重大变更默认会被关闭,以便顺利过渡到新版本。经过几个版本之后,一旦足够多的用户认为该功能在生产环境中稳定,我们会在后续版本中更改默认值。
  • 始终先升级查询端的 bundle(mr-bundle、presto-bundle、spark-bundle),然后再升级写入端(Hudi Streamer、使用 datasource 的 Spark 作业)。这样通常能获得最佳体验,而且通过回滚或推进写入端代码来修复问题也很简单(通常你对写入端拥有更多控制权)。
  • 对于体量大、功能丰富的版本,我们建议逐步迁移,先在预发布环境中测试并运行你自己的测试。升级 Hudi 与升级任何数据库系统并无二致。

请注意,发布说明可能会针对具体情况覆盖此信息并给出特定指引。

升级到 1.0.0

1.0.0 是一个包含重大格式变更的主要版本。为确保迁移过程顺利,建议按照以下步骤操作:

  1. 完全停止 0.x 中的所有异步表服务。
  2. 将写入端升级到 1.x,使用表版本(tv)6,禁用 autoUpgrade 和元数据(这不会自动升级任何内容);0.x 的读取端将继续正常工作;写入端同时也可以作为读取端,并将继续读取 tv=6。
    a. 将 hoodie.write.auto.upgrade 设置为 false。
    b. 将 hoodie.metadata.enable 设置为 false。
    upgrade1.0-1
  3. 将表服务升级到 1.x,使用 tv=6,并恢复其运行。
    upgrade1.0-2
  4. 将所有剩余的读取端升级到 1.x,使用 tv=6。
    upgrade1.0-3
  5. 使用 tv=8 重新部署写入端;表服务和读取端会自适应并动态获取 tv=8。
  6. 当所有读取端和写入端都升级到 1.x 后,我们就可以在 1.x 表上启用任何新功能了,包括元数据。
    upgrade1.0-4

在升级过程中,元数据表不会被更新,因此会落后于数据表。需要注意的是,只有当写入者升级到 tv=8 时,元数据表才会被更新。因此,在滚动升级期间,即使读取者也应保持元数据表功能处于禁用状态,直到所有写入者都升级到 tv=8 为止。

:::caution
大多数情况都会由自动升级过程无缝处理,但仍然存在一些限制。在继续进行迁移之前,请先通读升级与降级过程的限制说明。详情请参阅 RFC-78。
:::

降级

使用新版本 Hudi 时,升级是自动完成的,而降级则是一个手动步骤。我们需要使用 Hudi CLI 将表从较高版本降级到较低版本。下面以一个示例说明:我们先使用 0.12.0 创建一个表,将其升级到 0.13.0,然后通过 Hudi CLI 对其进行降级。

以 Hudi 0.11.0 版本启动 spark shell。

spark-shell \
  --packages org.apache.hudi:hudi-spark3.2-bundle_2.12:0.11.0 \
  --conf 'spark.serializer=org.apache.spark.serializer.KryoSerializer' \
  --conf 'spark.sql.catalog.spark_catalog=org.apache.spark.sql.hudi.catalog.HoodieCatalog' \
  --conf 'spark.sql.extensions=org.apache.spark.sql.hudi.HoodieSparkSessionExtension'

使用下面的 Scala 脚本创建一个 Hudi 表。

import org.apache.hudi.QuickstartUtils._
import scala.collection.JavaConversions._
import org.apache.spark.sql.SaveMode._
import org.apache.hudi.DataSourceReadOptions._
import org.apache.hudi.DataSourceWriteOptions._
import org.apache.hudi.config.HoodieWriteConfig._
import org.apache.hudi.common.model.HoodieRecord
import org.apache.hudi.common.table.timeline.HoodieTimeline
import org.apache.hudi.common.fs.FSUtils
import org.apache.hudi.HoodieDataSourceHelpers

val dataGen = new DataGenerator
val tableType = MOR_TABLE_TYPE_OPT_VAL
val basePath = "file:///tmp/hudi_table"
val tableName = "hudi_table"

val inserts = convertToStringList(dataGen.generateInserts(100)).toList
val insertDf = spark.read.json(spark.sparkContext.parallelize(inserts, 2))
insertDf.write.format("hudi").
        options(getQuickstartWriteConfigs).
        option("hoodie.table.ordering.fields", "ts").
        option("hoodie.datasource.write.recordkey.field", "uuid").
        option("hoodie.datasource.write.partitionpath.field", "partitionpath").
        option("hoodie.table.name", tableName).
        option("hoodie.datasource.write.operation", "insert").
        mode(Append).
        save(basePath)

你将在 hoodie.properties 中看到一个 table version 条目,其中说明表版本为 4。

bash$ cat /tmp/hudi_table/.hoodie/hoodie.properties | grep hoodie.table.version
hoodie.table.version=4

使用 0.13.0 版本启动新的 Spark shell,并使用上面的脚本向同一张表追加数据。请注意,升级会随新版本自动完成。

spark-shell \
  --packages org.apache.hudi:hudi-spark3.2-bundle_2.12:0.13.1 \
  --conf 'spark.serializer=org.apache.spark.serializer.KryoSerializer' \
  --conf 'spark.sql.catalog.spark_catalog=org.apache.spark.sql.hudi.catalog.HoodieCatalog' \
  --conf 'spark.sql.extensions=org.apache.spark.sql.hudi.HoodieSparkSessionExtension'

升级后,表版本会更新为 5。

bash$ cat /tmp/hudi_table/.hoodie/hoodie.properties | grep hoodie.table.version
hoodie.table.version=5

让我们尝试将表降级回版本 4。降级时我们需要使用 Hudi CLI 并执行 downgrade 命令。有关降级的更多详细信息,请参阅此处的文档。

connect --path /tmp/hudi_table
downgrade table --toVersion 4

降级后,表版本将更新为 4。

bash$ cat /tmp/hudi_table/.hoodie/hoodie.properties | grep hoodie.table.version
hoodie.table.version=4

迁移

目前,迁移到 Hudi 有两种方式:

  • 将较新的分区转换为 Hudi:该模式适用于大型事件表(例如:点击流、广告曝光),这类表通常也只接收最近几天的写入。你可以将最近 N 个分区转换为 Hudi,然后像操作 Hudi 表一样继续进行写入。Hudi 查询侧的代码能够正确处理 Hudi 与非 Hudi 的数据分区。
  • 完全转换为 Hudi:如果你当前每天会对表进行几次批量/全量加载(例如数据库数据摄入),则适合采用该模式。完全转换为 Hudi 只需一次性操作(类似于运行一次现有的作业),即可将所有数据迁移到 Hudi 格式,并为后续写入提供增量更新的能力。

更多详情请参阅详细的迁移指南。未来,我们将支持对现有表进行无缝的零拷贝引导(bootstrap),并完整支持所有更新(upsert)与增量查询能力。

评论

登录后参与评论

正在加载评论…