记录合并器
Hudi 处理记录变更和流式数据的方式,在时间线条目顺序一节中已有简要介绍。为了向用户提供完善的流处理支持,Hudi 不遗余力地让存储引擎及其底层存储格式理解如何合并具有相同记录键的变更——这些变更可能在不同的时间以不同的顺序到达。随着移动应用和物联网的兴起,这类场景已成为常态而非例外。例如,社交应用可能会在用户重新连接到 Wi‑Fi 之后的数小时,才上传其发生时的用户事件。
合并模式
为应对这些挑战,Hudi 支持合并模式(merge mode),它定义了文件切片(file slice)中基础文件与日志文件的排序方式,以及该文件切片内具有相同记录键的不同记录如何被一致地合并,从而为快照查询、写入方和表服务产生相同的确定性结果。
合并模式是一个表级配置,用于以下代码路径:
- (写入) 在写入过程中读取输入数据时,合并同一记录键的多条变更记录。这是一个可选的优化,用于减少写入日志文件的记录数量,从而提升后续的查询和写入性能。
- (写入) 针对 COW 表,将最终的变更记录(部分更新/全量更新/删除)与存储中已有的记录进行合并。
- (压缩) 压缩服务会遵循合并模式,将日志文件中的所有变更记录与基础文件进行合并。
- (查询) 针对 MOR 表查询,在完成过滤和投影后,将日志文件中的变更记录与基础文件进行合并。
合并模式共有三种:COMMIT_TIME_ORDERING、EVENT_TIME_ORDERING 和 CUSTOM。默认合并模式会根据是否配置了排序字段自动推断得出。如果未指定排序字段(例如 hoodie.table.ordering.fields),合并模式默认为 COMMIT_TIME_ORDERING,即用传入批次中的新记录替换旧记录。如果指定了一个或多个排序字段,合并模式默认为 EVENT_TIME_ORDERING,即根据排序字段的值来比较记录,以处理乱序数据。
你也可以通过写入配置 hoodie.write.record.merge.mode 显式指定合并模式。当使用该配置创建或写入表时,它会以 hoodie.record.merge.mode 的形式持久化到表的配置文件(hoodie.properties)中。持久化之后,除非在写入配置中被显式覆盖,否则后续所有的读写操作都会使用该合并模式。
合并模式是为你的表选择和配置记录合并器(record merger)的方式。对于大多数用例,你只需设置合适的排序字段,即可选择 COMMIT_TIME_ORDERING 或 EVENT_TIME_ORDERING,无需任何额外配置或实现。这些模式提供了涵盖绝大多数场景的标准合并行为。对于需要自定义合并逻辑的高级用例,可以使用 CUSTOM 合并模式并实现 HoodieRecordMerger 接口。
note
表创建后不应更改合并模式,以避免在不同模式之间切换时压缩(compaction)产生不同的合并结果,从而导致行为不一致。
COMMIT_TIME_ORDERING
在此模式下,我们期望输入记录以严格顺序到达,即到达顺序与其在表上的增量提交顺序一致。合并时会选择属于最新写入的记录作为合并结果。用关系数据模型的术语来说,这提供了与时间线上可串行化写入一致的覆盖语义。

在上面的示例中,写入进程消费数据库变更日志,该日志预期按照逻辑序列号(lsn)的严格顺序排列,lsn 表示上游数据库中写入的先后顺序。
EVENT_TIME_ORDERING
这是默认的合并模式。虽然按提交时间排序提供了易于理解的标准行为,但它远远不够。提交时间与用户可能关心的数据实际顺序无关,而且在复杂的分布式系统中,输入的严格排序很难实现。采用事件时间排序时,合并会选择用户指定的排序字段或预合并字段上取值最高的记录作为合并结果。

在上面的示例中,两个微服务在不同时间产生关于订单的变更记录,这些记录可能乱序到达。如颜色编码所示,如果仅仅按提交时间顺序进行合并,可能导致表中出现应用层面的不一致状态,例如已取消的订单被重新创建,或已支付的订单退回到刚创建的状态并再次等待支付。事件时间排序通过忽略延迟到达的旧状态变更,避免订单状态在时间上“回退”,从而解决了这一问题。结合非阻塞并发控制,这种方式为高效且正确地处理此类数据流提供了非常强大的手段。
CUSTOM
tip
**对大多数用户而言:**内置的 COMMIT_TIME_ORDERING 和 EVENT_TIME_ORDERING 合并模式应当已经足够。只有在需要标准模式无法实现的特殊合并逻辑时,才使用 CUSTOM 模式。
在某些情况下,可能需要更强的控制和定制能力。延续上面的例子,两个微服务可能分别更新两组不同的列——order_info 和 payment_info,以及订单状态。此时的合并逻辑不仅需要解析出正确的状态,还需要将处于 created 状态的记录中的 order_info 合并到处于 canceled 状态的记录中,后者已填充了因支付失败原因而对应的 payment_info 字段。这种对账机制为下游消费提供了一个简单的反规范化数据模型,查询(例如欺诈检测)可以直接跨 order_info 和 payment_info 过滤字段,而无需在每次访问时执行代价高昂的自连接。
要实现自定义合并逻辑,需要实现 HoodieRecordMerger 接口。Hudi 在标准记录合并器 API 之上支持编写跨语言的自定义记录合并器,该 API 支持完全合并和部分合并。HoodieRecordMerger 接口使用 BufferedRecord 类,通过直接操作引擎原生记录格式(无需转换为 Avro),从而提供更好的性能和一致性。
BufferedRecord 类封装了记录数据以及记录键、排序值、schema 标识符和 HoodieOperation 等关键信息。RecordContext 为跨不同引擎(Spark、Flink 等)访问和操作记录提供了统一接口,使自定义合并器与引擎无关。
以下从高层次概述了 Java API。该接口接收包装在 BufferedRecord 实例中的较旧/较新记录,并输出合并后的 BufferedRecord。记录合并器通过 hoodie.write.record.merge.strategy.id 写配置进行配置,其值为一个 UUID,写入器会将其持久化到表配置中,并且预期由下面的 getMergingStrategy() 方法返回。通过这一机制,Hudi 可以在不同的语言和引擎运行时中自动推断出表格应使用的记录合并器。
interface HoodieRecordMerger {
<T> BufferedRecord<T> merge(BufferedRecord<T> older, BufferedRecord<T> newer,
RecordContext<T> recordContext, TypedProperties props) throws IOException {
// Merges full records. Returns a non-null BufferedRecord.
// If the result represents a deletion, set HoodieOperation.DELETE on the returned record.
// Ordering values must always be set if there are ordering fields for the table.
...
}
<T> BufferedRecord<T> partialMerge(BufferedRecord<T> older, BufferedRecord<T> newer,
HoodieSchema readerSchema, RecordContext<T> recordContext,
TypedProperties props) throws IOException {
// Merges records which can contain partial updates.
// Returns a non-null BufferedRecord with only changed fields included.
// If the result represents a deletion, set HoodieOperation.DELETE on the returned record.
// Ordering values must always be set if there are ordering fields for the table.
...
}
boolean isProjectionCompatible() {
// Whether this merger can work on a projection of the table schema rather than every column.
// Defaults to false, which makes MOR reads fetch all columns. See "Projection compatibility" below.
...
}
String[] getMandatoryFieldsForMerging(HoodieSchema dataSchema, HoodieTableConfig cfg,
TypedProperties properties) {
// The columns merge() and partialMerge() read, beyond those the query itself requests.
// Only consulted when isProjectionCompatible() returns true.
// Defaults to the record key field plus the ordering fields.
...
}
HoodieRecordType getRecordType() {...}
String getMergingStrategy() {...}
}实现指南
在实现 HoodieRecordMerger 接口时,请遵循以下准则以确保结果的一致性:
- 返回非空记录:
merge()和partialMerge()方法都必须返回非空的BufferedRecord。只要可能,返回的记录都应包含合并后的数据,即使该数据代表一次删除。这样可以使后续的合并操作能够引用数据的前值。不过,如果数据不可用或不需要,返回 data 为 null 的BufferedRecord也是可以接受的(例如使用BufferedRecords.createDelete()时)。 - 处理删除操作:如果合并结果应当删除与该记录键匹配的行,请在返回的
BufferedRecord上通过setHoodieOperation(HoodieOperation.DELETE)将HoodieOperation设置为DELETE。 - 保留排序值:如果表配置了任何排序字段,务必在结果中设置排序值。这样可以确保后续的合并操作能够正确引用这些值。
- 使用 RecordContext:
RecordContext参数提供了与引擎无关的方法,用于访问字段值、提取记录键以及操作记录。请使用getValue()、getRecordKey()和mergeWithEngineRecord()等方法来处理底层数据。 - 结合律:
merge()方法应当满足结合律:对于同一条记录的任意三个版本 A、B、C,merge(a, merge(b, c))应当与merge(merge(a, b), c)得到相同的结果。 - 声明合并逻辑读取的每一列:如果你使 merger 兼容投影,
getMandatoryFieldsForMerging()必须列出merge()和partialMerge()所触及、且查询本身并未请求的每一列。详见下文。
有关实现的更多细节,请参阅 RFC 101。
投影兼容性
在 Merge-on-Read 表上,读取方更希望只获取查询所要求的列。但对于自定义 merger,它无法安全地做到这一点,因为它无从得知你的合并逻辑会读取哪些列。有两个方法可以解决这个问题,它们成对工作:
isProjectionCompatible()—— 默认值为false,此时会告诉读取方回退到读取完整的表 schema。在这种模式下合并始终是正确的,但即使查询只选择了两列,仍会读取每个日志块的每一列,因此这是较慢的选项。COMMIT_TIME_ORDERING和EVENT_TIME_ORDERING是投影兼容的;CUSTOMmerger 在你明确声明之前并不兼容。getMandatoryFieldsForMerging()—— 仅在isProjectionCompatible()返回true后才会被查询。读取方随后会读取查询所请求的列以及该方法指定的列。其默认值为记录键字段和排序字段,这正是只比较排序值的 merger 所需要的。
这两者必须一致,而问题恰恰出在这里:如果 isProjectionCompatible() 返回 true,而你的合并逻辑读取了某个既未被查询请求、也未出现在 getMandatoryFieldsForMerging() 中的列,那么该列在你的合并器收到的记录中是不存在的。接下来会发生什么取决于你的代码——读取缺失字段通常会抛出 NullPointerException,而一个能容忍该字段为 null 的合并器则会基于不完整的数据做出决策,从而悄悄返回错误结果。无论哪种情况,其表现都取决于查询,只会出现在恰好没有选择该列的查询上,因此在测试中很容易被遗漏。
因此,如果你的合并器通过比较某个 priority 列来决定胜出者,要么将其声明出来:
@Override
public boolean isProjectionCompatible() {
return true;
}
@Override
public String[] getMandatoryFieldsForMerging(HoodieSchema dataSchema, HoodieTableConfig cfg,
TypedProperties properties) {
// Keep the defaults - record key and ordering fields - and add the column merge() reads, so the
// reader includes it even when a query does not project it.
LinkedHashSet<String> fields = new LinkedHashSet<>(
Arrays.asList(HoodieRecordMerger.super.getMandatoryFieldsForMerging(dataSchema, cfg, properties)));
fields.add("priority");
return fields.toArray(new String[0]);
}保持 isProjectionCompatible() 的默认值 false,接受全字段读取。两种做法都是正确的;只有折中的做法——声明为投影兼容却不声明字段——是不正确的。
合并模式配置
记录合并模式,以及可选的记录合并策略 ID 和自定义合并实现类,可通过以下配置指定。
配置名称默认值说明
hoodie.write.record.merge.modeEVENT_TIME_ORDERING(设置了排序字段时)
COMMIT_TIME_ORDERING(未设置排序字段时)决定具有相同记录键的不同记录的合并逻辑。有效值:(1) COMMIT_TIME_ORDERING:使用提交时间合并记录,即较晚提交的记录覆盖相同键的较早记录。(2) EVENT_TIME_ORDERING:使用事件时间作为排序依据合并记录,即事件时间较大的记录覆盖相同键上事件时间较小的记录,与提交时间无关。事件时间或排序字段需要由用户指定。当配置了排序字段时,这是默认值。(3) CUSTOM:使用用户指定的自定义合并逻辑。Config Param: RECORD_MERGE_MODESince Version: 1.0.0
hoodie.write.record.merge.strategy.id无(可选)记录合并策略的 ID。Hudi 会从 hoodie.write.record.merge.custom.implementation.classes 中选取合并策略 ID 相同的 HoodieRecordMerger 实现。使用自定义合并逻辑时,需要同时指定此配置和 hoodie.write.record.merge.custom.implementation.classes。Config Param: RECORD_MERGE_STRATEGY_IDSince Version: 0.13.0Alternative: hoodie.datasource.write.record.merger.strategy(已弃用)
hoodie.write.record.merge.custom.implementation.classes无(可选)由 HoodieRecordMerger 实现组成的列表,Hudi 根据所用引擎从中构建合并策略。Hudi 会从该列表中选择第一个满足以下条件的实现:(1) 与 hoodie.write.record.merge.strategy.id 中指定的合并策略 ID 相同(如果提供了该配置);(2) 与执行引擎兼容(例如,Spark 使用 SPARK merger,Flink 使用 FLINK merger,Java/Hive 使用 AVRO)。列表中的顺序很重要——请将首选实现放在最前面。特定于引擎的实现(SPARK、FLINK)效率更高,因为它们避免了 Avro 序列化/反序列化的开销。Config Param: RECORD_MERGE_IMPL_CLASSESSince Version: 0.13.0Alternative: hoodie.datasource.write.record.merger.impls(已弃用)
记录 Payload(已弃用)
caution
弃用公告: 从 1.1.0 版本开始,基于 payload 的记录合并方式已被弃用。这种方式与 Avro 格式的记录紧密耦合,因此与原生查询引擎格式(例如 Spark InternalRow)的兼容性较差,维护难度也更大。我们强烈建议迁移到合并模式,它能为现代湖仓架构提供更好的灵活性、性能和可维护性。
现有的基于 payload 的配置将通过向后兼容继续可用,但我们建议用户迁移其实现方式。详情请参阅 RFC 97。
Record payload 是一种较早期的抽象/API,用于实现类似的记录级合并能力。虽然 record payload 曾非常有用且广受欢迎,但它存在一些缺点:由于需要将引擎原生记录格式转换为 Apache Avro 进行合并,性能较低;同时缺乏跨语言支持。正如下面将要介绍的,Hudi 针对不同用例提供了对多种 payload 的开箱即用支持。为了保证现有 payload 实现的向后兼容性,Hudi 内部实现了从 record merger API 到 payload API 的回退机制。
OverwriteWithLatestAvroPayload(已弃用)
hoodie.datasource.write.payload.class=org.apache.hudi.common.model.OverwriteWithLatestAvroPayload这是默认的记录载荷实现。它通过选取 precombine 键值最大的记录(即对 precombine 键的值调用 .compareTo() 来确定大小)来打破并列,合并时则直接选取最新记录。这提供了 latest-write-wins(后写入者获胜)式的语义。
DefaultHoodieRecordPayload(已弃用)
hoodie.datasource.write.payload.class=org.apache.hudi.common.model.DefaultHoodieRecordPayloadOverwriteWithLatestAvroPayload 依据排序字段进行预合并,并在合并阶段选取最新记录,而 DefaultHoodieRecordPayload 在预合并和合并两个阶段都遵循排序字段。让我们通过一个例子来理解两者的区别:
假设排序字段为 ts,记录主键为 id,其 schema 如下:
{
[
{"name":"id","type":"string"},
{"name":"ts","type":"long"},
{"name":"name","type":"string"},
{"name":"price","type":"string"}
]
}存储中的当前记录:
id ts name price
1 2 name_2 price_2传入记录:
id ts name price
1 1 name_1 price_1使用 OverwriteWithLatestAvroPayload(最新写入优先)合并后的结果数据:
id ts name price
1 1 name_1 price_1使用 DefaultHoodieRecordPayload 合并后的结果数据(始终遵循排序字段):
id ts name price
1 2 name_2 price_2EventTimeAvroPayload(已弃用)
hoodie.datasource.write.payload.class=org.apache.hudi.common.model.EventTimeAvroPayload这是基于 Flink 写入的默认记录载荷(payload)。某些用例需要按事件时间合并记录,因此事件时间在这里充当排序字段。对于乱序到达的数据,该载荷尤其有用。针对这类用例,用户需要设置 payload 事件时间字段 配置。
OverwriteNonDefaultsWithLatestAvroPayload(已弃用)
hoodie.datasource.write.payload.class=org.apache.hudi.common.model.OverwriteNonDefaultsWithLatestAvroPayload此 payload 与 OverwriteWithLatestAvroPayload 非常相似,仅在合并记录时略有不同。在预合并(precombining)阶段,与 OverwriteWithLatestAvroPayload 一样,它会基于排序字段为同一个 key 选取最新的记录。在合并阶段,它只会用新记录中与该字段默认值不相等的指定字段去覆盖存储中的现有记录。
PartialUpdateAvroPayload(已废弃)
hoodie.datasource.write.payload.class=org.apache.hudi.common.model.PartialUpdateAvroPayload该 payload 支持部分更新。通常,一旦合并步骤确定了要选取哪条记录,存储中的记录就会被完全替换为所解析出的记录。但在某些情况下,需求是仅更新某些字段,而不是替换整条记录,这称为部分更新。PartialUpdateAvroPayload 为此类用例提供了开箱即用的支持。为了说明这一点,我们来看一个简单的例子:
假设排序字段为 ts,记录键为 id,schema 如下:
{
[
{"name":"id","type":"string"},
{"name":"ts","type":"long"},
{"name":"name","type":"string"},
{"name":"price","type":"string"}
]
}存储中的当前记录:
id ts name price
1 2 name_1 null传入记录:
id ts name price
1 1 null price_1使用 PartialUpdateAvroPayload 合并后的结果数据:
id ts name price
1 2 name_1 price_1Record Payload 配置(已弃用)
可以使用以下配置指定 Payload 类。如需更高级的配置,请参阅配置页面。
基于 Spark 的配置:
配置名称默认值说明
hoodie.datasource.write.payload.classorg.apache.hudi.common.model.OverwriteWithLatestAvroPayload (可选)所使用的 Payload 类。如果你想在插入/更新时实现自己的合并逻辑,可以覆盖该配置。这将使为 PRECOMBINE_FIELD_OPT_VAL 设置的任何值失效
Config Param: WRITE_PAYLOAD_CLASS_NAME
基于 Flink 的配置:
配置名称默认值说明
payload.classorg.apache.hudi.common.model.EventTimeAvroPayload (可选)所使用的 Payload 类。如果你想在插入/更新时实现自己的合并逻辑,可以覆盖该配置。这将使为该选项设置的任何值失效
Config Param: PAYLOAD_CLASS_NAME
此外还有不少其他实现。开发人员可能有兴趣查看 HoodieRecordPayload 接口的类层次结构。例如,MySqlDebeziumAvroPayload 和 PostgresDebeziumAvroPayload 支持无缝应用通过 Debezium 为 MySQL 和 PostgresDB 捕获的变更;AWSDmsAvroPayload 支持将通过 Amazon Database Migration Service 捕获的变更应用到 S3。完整配置请参阅配置页面,如果你想实现自己的自定义 Payload,请查看 FAQ。
博客
评论
登录后参与评论
KnowForge