表服务

压缩

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

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

背景

Compaction(合并)是 Hudi 专门在 MOR(Merge-on-Read,读时合并)表上使用的一种表服务,用于定期将基于行的日志文件中的更新合并到相应的基于列式存储的基础文件中,从而生成基础文件的新版本。Compaction 不适用于 COW(Copy-on-Write,写时复制)表,只适用于 MOR 表。

为什么 MOR 表需要 compaction?

要理解 compaction 在 MOR 表中的重要性,首先需要了解 MOR 表的布局。在 Hudi 中,数据是以文件组为单位组织的。MOR 表中的每个文件组由一个基础文件(base file)和一个或多个日志文件(log file)组成。通常在写入时,插入操作的数据存储在基础文件中,而更新操作的数据则追加到日志文件中。

mor_table_file_layout 图:MOR 表文件布局,展示了包含基础数据文件和日志文件的不同文件组

在 compaction 过程中,来自日志文件的更新会与基础文件合并,形成基础文件的新版本,如下图所示。由于 MOR 是为写入优化设计的,因此在新写入时,索引标记完成后,Hudi 会将与每个文件组相关的记录以日志块(log block)的形式追加到日志文件中。写入过程中不会进行同步合并,从而降低了写放大(write amplification)并提升了写入延迟表现。相比之下,对 COW 表进行新写入时,Hudi 会将新写入的数据与原有的基础文件合并,产生基础文件的新版本,这会导致更高的写放大和更大的写入延迟。

mor_table_file_layout 图:对给定文件组执行 compaction

在服务读取查询(快照读)时,对于每个文件组,基础文件中的记录与其所有对应日志文件中的记录会被合并在一起后返回。因此,MOR 快照查询的读取延迟可能会比 COW 表更高,因为 COW 在读取时无需进行合并。Compaction 会定期将日志文件中的更新与基础文件合并,从而限制日志文件的增长,确保读取延迟不会飙升。

Compaction 架构

Compaction 分为两个步骤。

  • Compaction 调度:在这一步中,Hudi 会扫描分区并选择需要合并的文件切片(file slice),最终将 compaction 计划写入 Hudi 时间线(timeline)。
  • Compaction 执行:在这一步中,读取 compaction 计划并执行文件切片的合并。

Compaction 调度中的策略

compaction 的调度涉及两种策略:

  • 触发策略(Trigger Strategy):决定多久触发一次 compaction 调度。
  • 合并策略(Compaction Strategy):决定对哪些文件组执行合并。

Hudi 为这两种策略提供了多种选项,下面分别进行讨论。

增量调度

Hudi 支持增量式压缩(compaction)调度,这在分区数量庞大的表上能显著提升性能。增量式调度不会在每次压缩调度运行时扫描所有分区,而只处理自上次完成压缩以来发生变更的分区。

该功能默认通过 hoodie.table.services.incremental.enabled 启用。启用后,压缩调度将会:

  1. 通过分析上次完成压缩与当前调度即时(instant)之间时间窗口内的提交元数据,识别出自上次完成压缩以来被修改过的分区
  2. 纳入之前调度运行中被标记为遗漏的分区(例如,由于 IO 限制)
  3. 仅扫描并处理这些增量分区
  4. 如果找不到上次完成压缩的即时(例如,由于归档),或在获取增量分区时发生异常,则回退为扫描所有分区

对于分区数量很多的表,这种优化可以大幅降低调度开销,并提升作业的整体稳定性。

触发策略

配置名称默认值描述
hoodie.compact.inline.trigger.strategyNUM_COMMITS(可选)org.apache.hudi.table.action.compact.CompactionTriggerStrategy:控制何时调度压缩。 Config Param: INLINE_COMPACT_TRIGGER_STRATEGY
  • NUM_COMMITS:当上次完成压缩之后至少有 N 个增量提交时触发压缩。
  • NUM_COMMITS_AFTER_LAST_REQUEST:当上次完成或请求的压缩之后至少有 N 个增量提交时触发压缩。
  • TIME_ELAPSED:当距上次压缩已过去 N 秒后触发压缩。
  • NUM_AND_TIME:当上次完成压缩之后同时满足至少有 N 个增量提交且已过去 N 秒时触发压缩(两个条件都必须满足)。
  • NUM_OR_TIME:当上次完成压缩之后至少有 N 个增量提交或已过去 N 秒时触发压缩(满足任一条件即可)。

压缩策略

配置名称默认值描述
hoodie.compaction.strategyorg.apache.hudi.table.action.compact.strategy.LogFileSizeBasedCompactionStrategy(可选)压缩策略决定每次压缩运行时选择哪些文件组进行压缩。默认情况下,Hudi 会选择累积未合并数据最多的日志文件。Config Param: COMPACTION_STRATEGY

可用策略(使用策略时请提供完整的包名):

  • LogFileNumBasedCompactionStrategy:根据日志文件总数对压缩任务进行排序,过滤出日志文件数超过阈值的文件组,并将压缩操作限制在配置的 IO 上限之内。
  • LogFileSizeBasedCompactionStrategy:根据日志文件总大小对压缩任务进行排序,过滤出日志文件大小超过阈值的文件组,并将压缩操作限制在配置的 IO 上限之内。
  • BoundedIOCompactionStrategy:一种 CompactionStrategy,它会考虑压缩所需的总 IO(读取 + 写入),并将压缩列表限制在配置的 IO 上限以内。
  • BoundedPartitionAwareCompactionStrategy:该策略确保即使表中创建了更新的分区,最后 N 个分区仍会被选中。lastNPartitions 定义为 currentDate 之前的 N 个分区。假设 currentDay = 2018/01/01,表中存在超出 currentDay 的 2018/02/02 和 2018/03/03 分区,该策略将选择以下分区进行压缩:(2018/01/01, allPartitionsInRange[(2018/01/01 - lastNPartitions) 到 2018/01/01), 2018/02/02, 2018/03/03)
  • DayBasedCompactionStrategy:该策略按照 Hive 分区创建时间的倒序排列压缩任务。它有助于优先压缩最新分区中的数据,然后再处理较旧的分区,总量受限于允许的 Total_IO。
  • UnBoundedCompactionStrategy:UnBoundedCompactionStrategy 不会改变排序,也不会过滤任何压缩任务。它是一个直通策略,会压缩所有存在日志文件的基础文件。这通常意味着压缩时不做任何智能决策。
  • UnBoundedPartitionAwareCompactionStrategy:UnBoundedPartitionAwareCompactionStrategy 是一个自定义的 UnBounded 策略。它会过滤出所有可由 {@link BoundedPartitionAwareCompactionStrategy} 压缩的分区,并返回结果。这样做是为了避免长时间运行的 UnBoundedPartitionAwareCompactionStrategy 覆盖正在较短时间内运行的 BoundedPartitionAwareCompactionStrategy 所选中的分区。本质上,它是 BoundedPartitionAwareCompactionStrategy 所选分区的反向选择。

note

更多详情请参考 高级配置。

元数据表压缩触发策略

自 Hudi 1.2.0 起,元数据表(MDT)支持与数据表相同的压缩触发策略集合,并额外支持基于时间的选项。

配置名称默认值描述
hoodie.metadata.compact.trigger.strategyNUM_COMMITSMDT(元数据表)压缩的触发策略。接受与 hoodie.compact.inline.trigger.strategy 相同的取值:NUM_COMMITS、NUM_COMMITS_AFTER_LAST_REQUEST、TIME_ELAPSED、NUM_AND_TIME、NUM_OR_TIME。
hoodie.metadata.compact.max.delta.commits10在上一次 MDT 压缩之后,经过多少次增量提交(delta commit)才会调度一次新的压缩(适用于基于 NUM_COMMITS 的策略)。
hoodie.metadata.compact.max.delta.seconds7200在上一次 MDT 压缩之后,经过多少秒才会调度一次新的压缩。仅对 TIME_ELAPSED、NUM_AND_TIME 和 NUM_OR_TIME 策略生效。

触发 Compaction 的方式

内联(Inline)

默认情况下,压缩是异步执行的。

如果你非常看重记录的写入延迟,那么你很可能在使用 Merge-on-Read 表。Merge-on-Read 表采用列式(如 Parquet)与行式(如 Avro)文件格式相结合的方式来存储数据。更新会先记录到增量文件(delta file)中,之后再通过压缩生成新的列式文件版本。为了降低写入延迟,异步压缩(Async Compaction)是默认配置。

如果你更看重新提交(commit)的即时读取性能,或者希望不必单独管理压缩作业以简化运维,那么你可能希望使用同步内联压缩,也就是说,在写入一次提交的同时,也会由同一个作业完成压缩。

对于这种部署模式,请为 Spark Datasource 和 Spark SQL 写入器设置 hoodie.compact.inline = true。对于 Hudi Streamer 的 sync once 模式,可以通过传入标志 --disable-compaction(即禁用异步压缩)来实现内联压缩。此外,在 Hudi Streamer 中,当写入(ingestion)与压缩运行在同一个 Spark context 中时,你可以使用 Hudi Streamer CLI 中的资源分配配置,例如 --delta-sync-scheduling-weight、--compact-scheduling-weight、--delta-sync-scheduling-minshare 和 --compact-scheduling-minshare,来控制写入与压缩之间的执行器(executor)资源分配。

异步与离线压缩模型

这里有几种触发压缩的方式。

在同一进程内异步执行

在 Hudi Streamer 连续模式、Flink 以及 Spark Streaming 等流式摄入写入模型中,异步 compaction 默认处于启用状态,会在正常摄入的同时并行运行,而不会阻塞常规写入。

Spark Structured Streaming

compaction 会在流作业内部被异步调度并执行。以下是用 Java 编写的示例代码片段:

import org.apache.hudi.DataSourceWriteOptions;
import org.apache.hudi.HoodieDataSourceHelpers;
import org.apache.hudi.config.HoodieCompactionConfig;
import org.apache.hudi.config.HoodieWriteConfig;

import org.apache.spark.sql.streaming.OutputMode;
import org.apache.spark.sql.streaming.ProcessingTime;
 DataStreamWriter<Row> writer = streamingInput.writeStream().format("org.apache.hudi")
        .option("hoodie.datasource.write.operation", operationType)
        .option("hoodie.datasource.write.table.type", tableType)
        .option("hoodie.datasource.write.recordkey.field", "_row_key")
        .option("hoodie.datasource.write.partitionpath.field", "partition")
        .option("hoodie.table.ordering.fields", "timestamp")
        .option("hoodie.compact.inline.max.delta.commits", "10")
        .option("hoodie.datasource.compaction.async.enable", "true")
        .option("hoodie.table.name", tableName).option("checkpointLocation", checkpointLocation)
        .outputMode(OutputMode.Append());
 writer.trigger(new ProcessingTime(30000)).start(tablePath);
Hudi Streamer 连续模式

Hudi Streamer 提供连续摄取模式,在该模式下,单个长期运行的 Spark 应用程序会持续从上游数据源摄取数据并写入 Hudi 表。在此模式下,Hudi 支持管理异步压缩(compaction)。以下是在连续模式下运行并启用异步压缩的示例片段:

spark-submit --packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.0,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.0 \
--class org.apache.hudi.utilities.streamer.HoodieStreamer \
--table-type MERGE_ON_READ \
--target-base-path <hudi_base_path> \
--target-table <hudi_table> \
--source-class org.apache.hudi.utilities.sources.JsonDFSSource \
--source-ordering-field ts \
--props /path/to/source.properties \
--continous

由独立进程调度与执行

对于某些表服务运行时间较长的使用场景,用户可以选择将压缩的两个步骤(调度与执行)完全放在一个独立的进程中离线完成,而不是让常规写入阻塞。这样常规写入方无需关心这些压缩步骤,用户也可以根据需要为压缩作业分配更多资源。

note

这种模式需要为所有作业(包括常规写入作业和离线压缩作业)配置锁提供者。

内联调度、异步执行

在这种模式下,Spark Datasource 写入器或 Flink 作业可以只在内联调度压缩(这会将压缩计划序列化到 timeline 中,但不会执行它)。随后,可以由 HudiCompactor 或 HoodieFlinkCompactor 等独立工具定期执行压缩计划。

note

如果启用了元数据表(metadata table),这种模式可能需要配置锁提供者。

Hudi Compactor 工具

Hudi 提供了一个独立工具,用于异步执行特定的压缩操作。下面是一个示例,你可以在部署指南中了解更多信息。compactor 工具支持压缩的调度与执行。

示例:

spark-submit --packages org.apache.hudi:hudi-utilities-slim-bundle_2.12:1.2.0,org.apache.hudi:hudi-spark3.5-bundle_2.12:1.2.0 \
--class org.apache.hudi.utilities.HoodieCompactor \
--base-path <base_path> \
--table-name <table_name> \
--schema-file <schema_file> \
--instant-time <compaction_instant>

注意,instant-time 参数对于 Hudi Compactor Utility 现在是可选的。如果在不指定 --instant-time 的情况下使用该工具,spark-submit 将执行 Hudi 时间线上最早调度的 compaction。

可用选项

HoodieCompactor 工具支持以下重试和超时选项(在 scheduleAndExecute 模式下生效):

选项名称短标志默认值说明
--retry-last-failed-job-rcfalse设置为 true 时,会检查并回滚最近一次失败的 compaction 计划并重新执行,而不是直接规划新的 compaction 作业。这有助于从之前的失败中恢复。
--job-max-processing-time-ms-jt0将 compaction 作业判定为失败之前的最大处理时间(毫秒)。如果超过该时间且作业仍未完成,Hudi 将视为作业失败并重新启动它(需与 --retry-last-failed-job 一起使用)。值为 0 或负数时将禁用超时检查。

note

这些重试选项仅在使用 --mode scheduleAndExecute 时生效。--retry-last-failed-job 选项需要将 --job-max-processing-time-ms 设置为正值,才能检测滞留的 inflight 瞬时操作(instant)。

Hudi CLI

Hudi CLI 是另一种异步执行特定 compaction 的方式。以下是一个示例,你可以在部署指南中了解更多信息。

示例:

hudi:trips->compaction run --tableName <table_name> --parallelism <parallelism> --compactionInstant <InstantTime>
...

Flink 离线 Compaction

离线 Compaction 需要在命令行提交 Flink 任务。程序入口如下:hudi-flink-bundle_2.11-0.9.0-SNAPSHOT.jar : org.apache.hudi.sink.compact.HoodieFlinkCompactor

# Command line
./bin/flink run -c org.apache.hudi.sink.compact.HoodieFlinkCompactor lib/hudi-flink-bundle_2.11-0.9.0.jar --path hdfs://xxx:9000/table

选项

选项名称默认值描述
--pathn/a **(必填)**目标表在 Hudi 中的存储路径
--compaction-max-memory100(可选)压缩(compaction)过程中日志数据索引映射的大小,默认为 100 MB。如果内存充足,可以调大该参数
--schedulefalse(可选)是否执行调度压缩计划的操作。当写入流程仍在写入时,开启该参数存在数据丢失的风险。因此,开启该参数时必须确保当前没有任何写入任务正在向该表写入数据
--seqLIFO(可选)压缩任务的执行顺序。默认从最新的压缩计划开始执行。LIFO:从最新的计划开始执行。FIFO:从最旧的计划开始执行。
--servicefalse(可选)是否启动一个监控服务,按配置的时间间隔检查并调度新的压缩任务。
--min-compaction-interval-seconds600(s)(可选)服务模式下的检查间隔,默认为 10 分钟。
--retry0(可选)压缩操作的重试次数。仅在单次运行模式(非服务模式)下有效。默认为 0(不重试)。
--retry-last-failed-jobfalse(可选)当进行中的即时(inflight instant)超过最大处理时间时,检查并重试上次失败的压缩作业。仅在单次运行模式下有效。需要将 --job-max-processing-time-ms 设置为正值。
--job-max-processing-time-ms0(可选)在将压缩作业判定为失败之前的最大处理时间(毫秒)。与 --retry-last-failed-job 搭配使用。默认值 0 表示不进行超时检查。

note

重试相关选项(--retry、--retry-last-failed-job、--job-max-processing-time-ms)仅在单次运行模式下生效,在服务模式下无效。服务模式通过其持续监控循环具备隐式的重试语义。如果启用了 --retry-last-failed-job 但 --job-max-processing-time-ms 未设置为正值,则会记录一条警告。

日志压缩(Log Compaction)

日志压缩是针对 Merge-on-Read 表的一种轻量压缩。它不会将日志文件合并为新的基础文件,而是在同一文件组内将若干小的日志块拼接成一个更大的块,因此对于频繁收到小规模更新的文件组,可以在不必重写其基础文件的情况下保持高效。读取时会跳过已经拼接过的日志块,从而需要合并的块更少。

拼接的块是追加写入的,不会替换任何内容:被取代的块会保留在磁盘上,直到下一次完整的压缩与清理。因此,日志压缩是以额外的存储空间换取读取时更少的合并块。如果合并输出超过了日志块大小,单次运行也可能产生一个以上的块。

日志压缩在时间线上以 logcompaction 动作进行调度,处于 requested 和 inflight 状态。完成时它会以 deltacommit 提交,因此不存在可查找的已完成 logcompaction 实例。

配置名称 默认值 描述
hoodie.log.compaction.inline false(可选) 设置为 true 时,每次写入后都会触发日志压缩服务。这种方式在运维上更简单,但会在写入路径上增加额外的延迟。
Config Param: INLINE_LOG_COMPACT
Since Version: 0.13.0
hoodie.log.compaction.blocks.threshold 5(可选) 当文件切片至少包含这么多日志文件或至少这么多日志块时,即可调度日志压缩。适用于所有调度尝试,无论是内联触发还是以编程方式触发。
Config Param: LOG_COMPACTION_BLOCKS_THRESHOLD
Since Version: 0.13.0

note

hoodie.log.compaction.inline 是在数据表上调度日志压缩的唯一内置方式。与压缩(compaction)不同,数据表没有异步日志压缩服务,也没有 SQL 存储过程、Hudi CLI 命令或独立的实用工具。它也未通过 Flink 选项暴露,因此 Flink 无法为数据表调度日志压缩。可以通过写入客户端的 scheduleLogCompaction 和 logCompact 方法以编程方式进行调度。

元数据表运行其自身的日志压缩,由另一组独立的配置控制。它由维护元数据表的写入者执行。在 Flink 上无需额外配置:压缩管道会直接调度并执行元数据表的日志压缩。

配置名称 默认值 描述
hoodie.metadata.log.compaction.enable false(可选) 为元数据表启用日志压缩。
Config Param: ENABLE_LOG_COMPACTION_ON_METADATA_TABLE
Since Version: 0.14.0

hoodie.metadata.log.compaction.blocks.threshold 5(可选)日志块数量超过该值时,会为元数据表调度日志压缩。

Config Param: LOG_COMPACT_BLOCKS_THRESHOLD
Since Version: 0.14.0

caution

hoodie.metadata.table.service.manager.actions 的支持动作中虽然列出了 logcompaction,但它无法用于将元数据表的日志压缩移出写入端。设置它只会让写入端停止内联执行日志压缩,却不会把该任务交给任何地方:Hudi 没有任何分发路径可以将日志压缩委派给表服务管理器,而且实际上也并未提供表服务管理器。这样,元数据表上就会不断堆积待处理的 logcompaction 实例,并且根据下面的说明,元数据表的压缩也将不再被调度。

note

在元数据表的日志压缩处于待处理状态期间,元数据表的主压缩不会被调度,因为记录级索引等元数据分区依赖于处理时序。参见 HUDI-7533。

caution

hoodie.log.compaction.enable 也出现在配置参考中,但它并不是一个应在你的表上设置的开关。Hudi 会将其内部应用于元数据表自身的写入配置,其值派生自 hoodie.metadata.log.compaction.enable。在数据表上设置它不会产生任何效果:数据表应使用 hoodie.log.compaction.inline,元数据表则应使用 hoodie.metadata.log.compaction.enable。

有关该表服务背后的设计,参见 RFC-48。

博客

Apache Hudi Compaction Standalone HoodieCompactor Utility

评论

登录后参与评论

正在加载评论…