数据湖架构演进(一):Hudi 的增量数据内核
Hive 通过 Schema、分区、Metastore 和 SQL,将 HDFS 上的文件组织成可管理、可查询的数仓表,奠定了大数据离线数仓的基础。但随着原始日志、事件流和机器学习数据快速增长,企业开始需要更低成本的共享存储,并让多个计算引擎访问同一份数据,数据湖架构由此逐渐形成。
早期数据湖解决了开放存储、弹性扩展和多引擎访问问题,却没有完整解决事务提交、更新删除、并发写入和版本管理。特别是面对数据库 CDC 和高频 Upsert,传统的分区覆盖与批量重写成本越来越高。Hudi 由此在文件之上补回增量存储内核:通过 Timeline 管理提交与可见性,通过 Index 和 File Group 定位记录,并利用 COW、MOR 与 Table Services 平衡读写成本。
本文将从数据湖的演进背景出发,分析 Hudi 的整体架构和完整读写流程,再深入 Timeline、File Group、File Slice、Index、MOR、并发控制与后台维护等核心机制,最后结合 CDC 场景讨论这些设计如何在生产系统中协同工作。Iceberg、Paimon 与 Lance 将在后续文章中分别展开。
本文以 Apache Hudi 1.2.0、Table Version 9 为主要分析对象。旧版本在 Timeline 布局、Record-Level Index 和并发控制能力上可能存在差异。
从 Hive 数仓到 Lakehouse:数据湖为什么需要数据库内核
从传统数据仓库到 Hadoop/Hive 数仓
在 Hadoop 出现之前,企业通常使用关系型数据库或专用 MPP 数据库建设数据仓库。业务数据经过 ETL 清洗和建模后写入数据库,由数据库统一负责存储、事务、索引、查询优化和 SQL 分析。这种架构适合结构化数据与 BI 报表,但面对快速增长的日志、点击流和多媒体数据时,存储成本与横向扩展压力逐渐增大。
Hadoop 将数据从数据库管理的存储中释放出来:HDFS 使用普通服务器提供分布式存储,MapReduce 提供并行计算,Hive 则在文件之上增加 Schema、分区、Metastore 和 SQL,形成面向大规模离线分析的 Hive 数仓。

Hive 数仓降低了海量数据的存储和计算成本,但仍以批处理、目录和分区为核心。随着数据类型、处理时效和计算引擎进一步增加,数据架构开始继续向数据湖演进。
数据湖:共享存储上的多类型数据
随着日志、事件流、文档、图片和机器学习样本不断增长,企业开始先保存原始数据,再由不同团队按需加工和消费。数据湖由此形成了“低成本存储、开放格式、按需计算”的基本架构。
数据湖最初可以建立在 HDFS 上,云对象存储则进一步实现了存算分离,使 Spark、Flink、Trino 和 Ray 能够共享同一份底层数据。

开放存储降低了成本和平台绑定,却没有提供完整的表级语义:对象存储知道文件是否存在,Parquet 知道单个文件如何编码,但它们都不知道哪些文件共同构成当前表。
因此,数据湖解决了“数据如何存储和共享”,却没有天然解决“某一时刻哪一份数据才是一张正确的表”。
数据湖与数据仓库的两层架构为何仍然存在问题
数据湖适合低成本保存原始数据,并支持数据工程和机器学习;数据仓库则擅长结构化建模、SQL 分析和 BI。为了同时获得两者的能力,很多企业采用了“数据湖 + 数据仓库”的两层架构:数据先进入湖中,再经过 ETL 复制到数据仓库。

两层架构的问题不在于数据湖或数据仓库本身,而在于表状态和数据管理被拆分到了两个系统:
- 数据需要多次复制,增加存储成本和数据延迟;
- 多级 ETL 引入更多失败点,湖与仓可能出现 Schema 和版本不一致;
- 权限、血缘、质量和成本需要跨系统重复管理。
Databricks 与学术界合作者在 CIDR 2021 论文中系统化提出了 Lakehouse 架构:在低成本、可直接访问的开放存储之上,提供 ACID 事务、数据版本、审计、索引、缓存和查询优化等数据库能力,并同时支持 BI、数据工程和机器学习(Lakehouse: A New Generation of Open Platforms)。

Lakehouse 不是简单地把数据湖和数据仓库部署在一起,而是让数仓能力直接建立在数据湖存储之上,使不同工作负载尽可能围绕同一份受管理的数据运行。实现这一目标的关键,是在开放文件之上增加事务元数据层——这正是 Hudi、Iceberg 和 Delta Lake 等事务湖表所解决的问题。
Lakehouse 的存储基础:事务湖表
Lakehouse 是一种完整的数据架构,不是某种文件格式,也不等同于某一个产品。它的核心思想,是在低成本、开放的数据湖存储之上,重新建立事务、版本、治理和查询优化等数据库能力。

其中,事务湖表层是数据湖走向 Lakehouse 的关键,因为它需要回答一个基本问题:
- 在某一个时刻,哪些文件共同构成一张合法的数据表?
假设一个分布式任务需要生成 500 个 Parquet 文件,但在写完第 420 个文件后失败。虽然对象存储中已经出现了大量新文件,但这些文件不能构成一次完整提交,Reader 也不应该看到这批半成品。当 CDC、Backfill 和 Compaction 并发运行时,系统还需要提供四类表级语义:
- 原子提交:一批文件要么全部可见,要么全部不可见;
- 版本管理:Reader 能够读取一致的当前状态或历史状态;
- 并发控制:Writer 与维护任务发生冲突时能够检测、重试或回滚;
- 数据演进:正确处理 Update、Delete、Schema 和分区变化。
事务湖表通常采用“先写不可变文件,后原子发布表状态”的提交方式:

现代对象存储即使提供强一致的单文件读写和 Listing,也不会自动提供跨多个 Data File、Delete File 和 Metadata File 的表级事务。Hudi、Iceberg 和 Delta Lake 的核心价值,就是通过事务元数据定义表状态与提交边界。
不过,事务湖表只是 Lakehouse 的存储基础。完整的 Lakehouse 还需要查询优化、缓存、统一 Catalog、权限治理以及对 BI、流处理和 AI 工作负载的支持。
Hudi 与 Iceberg:两种事务湖表内核
Databricks 以 Delta Lake 作为 Lakehouse 的事务存储基础(Delta Lake 论文),通过事务日志为 Parquet 文件增加 ACID、Schema Enforcement、Time Travel 和批流统一能力。但 Delta Lake 只是事务湖表的一种实现,Hudi 与 Iceberg 也在开放存储之上构建了各自的表状态与提交协议。
这里先简单介绍下 Hudi 与 Iceberg,是因为它们代表了两种不同的内核设计思路:

Hudi 最初关注的是如何将数据库中的持续变化低延迟写入数据湖。它通过 Record Key 和 Index 定位记录所在的 File Group,再由 Timeline 管理写入、Compaction、Clustering 和 Cleaning。它首先回答的是:
- 同一条记录发生多次变化时,如何定位、合并并增量传播这些变化?
Iceberg 最初关注的是超大规模分析表的可靠性与查询规划。它通过 Snapshot 定义表版本,使用 Manifest Tree 组织 Data File 和 Delete File,再由 Catalog 原子发布新的表状态。它首先回答的是:
- 当一张表包含海量文件时,如何快速、准确地确定某个时刻的完整文件集合?
Hudi、Iceberg 与 Delta Lake 都来自真实生产系统中的一致性、更新效率和元数据扩展问题,并不是 Lakehouse 概念出现之后才产生的。Lakehouse 在更高层次上总结了它们共同推动的架构方向,而这些事务湖表则提供了 Lakehouse 所需的存储基础。
本文接下来将深入 Hudi 的 Timeline、File Group 与 File Slice 等核心机制,分析它如何定义一致的表状态,并处理数据更新、并发写入、故障恢复与历史版本管理。
Hudi 整体架构:从写入到读取
理解 Hudi,不能只把它看成一种新的文件格式。Parquet、ORC 解决的是单个文件内部如何编码数据,而 Hudi 解决的是另一组问题:
- 一条记录更新后,Writer 如何找到旧记录所在的文件?
- 分布式任务生成多个文件后,如何原子发布?
- Reader 如何识别已经提交、尚未完成和已经失效的文件?
- 高频更新产生的 Log、小文件和历史版本如何持续维护?
- 下游如何只消费某段时间内发生变化的数据?
从整体上看,Hudi 是构建在对象存储或 HDFS 之上的增量存储引擎。它通常不原地改写已经发布的 Base File。COW 生成新的 Base File,MOR 则将变化追加或写入 Log File,最后通过 Timeline 发布新的表状态。
Hudi 的核心组件
Hudi 的整体架构可以分为访问层、表内核、元数据层、存储层和后台维护层:

其中几个核心组件分别解决不同问题:
| 组件 | 核心作用 |
|---|---|
| Timeline | 记录写入和表维护操作,决定一次变化是否可见 |
| Index | 将 Record Key 定位到已有 File Group |
| File Group / File Slice | 组织一组记录在不同时间点的物理版本 |
| File-System View | 根据 Timeline 构造 Reader 应该看到的文件集合 |
| Metadata Table | 保存文件列表、Record Index 和统计信息 |
| Record Merger | 合并同一 Record Key 的多个版本 |
| Table Services | 持续处理 Log、小文件、旧版本和数据布局 |
这些组件共同构成了 Hudi 的基本工作方式:
1 | 写入时:Index 定位记录 → 写 Base 或 Log → Timeline 发布 |
一条 Upsert 如何写入 Hudi
假设数据库中的一条商品记录发生更新:
1 | product_id = 42 |
对应到 Hudi 中,上述一次数据的 Upsert 的主路径如下:

这条路径可以概括为四个阶段。
- Writer 从输入数据中生成 Record Key 和 Partition Path,并按照配置处理同一 Key 的重复版本。
- Index 查询这条记录是否已经存在。如果找到已有位置,记录被路由到对应 File Group;如果没有找到,则作为 Insert 分配到已有小文件或新的 File Group。
Writer 根据表类型执行物理写入:

- COW 在写入时完成记录合并并生成新的 Base File;MOR 则可以先把更新写入 Log,将部分物化成本延后到读取或 Compaction 阶段。
- 数据文件写完后,Writer 更新相关元数据,并将 Timeline 上的 Instant 发布为
COMPLETED。只有完成这一步,新文件才会对 Reader 可见。
因此,Hudi 的原子性并不依赖同时修改多个数据文件,而是先生成文件,最后通过 Timeline 原子发布新的表状态。没有进入 COMPLETED 的写入,即使已经在对象存储中生成了文件,也不会成为合法表版本(Hudi Write Operations)。
Reader 如何重建表状态
Hudi Reader 不会通过文件修改时间寻找“最新文件”,也不能直接读取目录中的全部文件。Reader 首先根据 Timeline 确定目标时间点的有效 Action,再通过 Metadata Table 或存储目录 Listing 获取文件信息,构造 File-System View,最后选择对应的读取方式:

不同查询类型对应不同的物理读取路径:
| 查询类型 | 读取方式 |
|---|---|
| Snapshot | 读取最新已提交状态;MOR 需要合并 Base 与 Logs |
| Read Optimized | 只读取最新完成物化的 Base File |
| Time Travel | 读取指定历史时间点的 File Slice |
| Incremental Query | 读取一段 Timeline 范围内发生变化的记录 |
| CDC Query | 读取具体的 Insert、Update 和 Delete 事件 |
由此可以建立对 Hudi 的整体理解:

Hudi 的核心并不是简单地让 Parquet 支持 Update,而是围绕持续变化的数据,建立一套从记录定位、物理写入、原子发布到一致读取的完整路径。
Hudi 内核:表状态、版本与并发如何实现
上一节从整体上介绍了 Hudi 的读写路径,本节继续下钻 Hudi 内核,重点分析表状态、文件版本、记录合并与并发控制的实现。Hudi 内核的核心概念如下图所示:

这些机制共同维护着 Hudi 表的五个核心不变量:
- 只有进入
COMPLETED状态的 Action 才能被 Reader 看见; - 在索引定义的唯一性范围内,一条逻辑记录在任意时刻只有一个有效位置;
- 对于任意查询时间点,每个有效 File Group 只能暴露一个对应的 File Slice;
- 同一 Record Key 的多个版本必须由一致的 Merge 规则得到确定结果;
- Compaction、Clustering 等后台操作可以改变物理布局,但不能意外改变表的逻辑状态。
Timeline:表状态的操作日志
对象存储即使能够保证单个对象的原子写入和强一致读取,也不会自动为一批 Base File、Log File 和 Metadata File 提供跨对象事务。Hudi 因此不原地重写已经发布的 Base File,而是先创建 Timeline Instant,再写出新的文件版本,最后通过 Timeline 原子发布新的表状态。
Timeline 与 Instant
Timeline 是一张 Hudi 表上所有操作的有序历史,Instant 则是其中的一条记录:
1 | Timeline |
三者分别表示:
- Time:操作何时开始和完成;
- Action:执行了什么操作;
- State:操作执行到什么阶段。
例如,一张表的 Timeline 可能包含:

这里,T1、T2 是 Instant 的 Begin Time;COMMIT、COMPACTION 是 Action 类型;COMPLETED、INFLIGHT 则是执行状态。
因此,Timeline 不是一次提交,Instant 也不只是一个时间戳:Timeline 是完整的操作历史,Instant 是其中一条由时间、Action 和 State 共同定义的记录。
Instant 状态流转
一次典型写入可以概括为:
- 创建
REQUESTEDInstant,并保存可选的执行计划; - 将 Instant 转为
INFLIGHT,开始生成 Base File、Log File 及相关元数据; - 完成文件校验、并发冲突检测,并生成 Commit Metadata;
- 原子发布
COMPLETEDInstant,使新状态对 Reader 可见。
一个典型 Instant 会经历三个状态:

REQUESTED:Action 已经被计划,但尚未执行;INFLIGHT:Writer 或 Table Service 正在生成文件;COMPLETED:Action 已经成功发布,可以参与表状态构建。
如果 Writer 在 INFLIGHT 阶段失败,原 Instant 不会进入名为 ROLLBACK 的状态,而是保持未完成。后续系统可以创建一个独立的 Rollback Action,或者由 Failed-Write Cleaner 清理残留文件。
Timeline 记录哪些 Action
Timeline 不只记录数据写入,也记录表的物理维护和故障恢复:
| Action | 主要语义 |
|---|---|
| COMMIT | 通常对应 COW 写入,生成新的 Base File |
| DELTA_COMMIT | 通常对应 MOR 写入,生成 Base File 或 Log File |
| REPLACE_COMMIT | 用新 File Group 原子替换旧 File Group |
| COMPACTION | 将 Base 与 Logs 合并为新的 Base File |
| LOG_COMPACTION | 合并多个 Log File,降低读取成本 |
| CLEAN | 删除超过保留策略的旧文件版本 |
| ROLLBACK | 撤销失败或未完成的写入 |
| INDEXING | 构建或更新 Metadata Index |
文件存在不等于文件可见
假设一个分布式任务计划生成 500 个文件,但 Driver 在写完第 420 个文件后失败:

Reader 判断文件是否可见,主要依据对应 Action 是否进入 COMPLETED,而不是文件是否已经存在:
1 | VisibleActions(T) |
Begin Time 表示 Action 的开始时间,Completion Time 才表示它进入可见序列的时间。并发场景下,两者顺序可能不同,因此 Reader 和增量消费必须以完成状态及 Completion Time 为准。
Hudi 将近期 Action 保存在 Active Timeline 中,较早的 Completed Action 则逐步写入基于 Parquet 和 Manifest 的 LSM Timeline History,以控制长期扫描成本(Hudi Technical Specification)。
File Group 与 File Slice:如何组织文件版本
Timeline 定义哪些操作已经完成,File Group 与 File Slice 则把这些操作产生的文件组织成可读取的版本。
File Group:稳定的文件版本分组
File Group 是 Hudi 组织数据文件的基本单元,由 Partition Path + file_id 唯一标识。
1 | FileGroupId = Partition Path + file_id |
file_id 是 Hudi 为一组记录分配的物理标识。同一个 File Group 在不同时间生成的 Base File 和 Log File 都会沿用这个 file_id。在索引定义的唯一性范围内,一条记录在任意时刻只应映射到一个有效 File Group:
1 | Record Key |
file_id 在 File Group 的生命周期内保持稳定,但 Clustering、Insert Overwrite 等 REPLACE_COMMIT 操作可能淘汰旧 File Group,并创建具有新 file_id 的 File Group。
File Slice:File Group 在某个时间点的物理版本
一个 File Group 由多个 File Slice 组成:
1 | FileGroup(partitionPath, file_id) |
File Slice 表示该 File Group 在某个 Base Instant 上的物理状态:
- Base File 保存已经完成列式物化的数据;
- Log File 保存尚未物化的更新、删除或 CDC 变化;
- Base File 是可选的,因此 MOR 也可能存在 Log-only File Slice。
COW 与 MOR 的核心差异,就是它们创建和更新 File Slice 的方式不同:

COW:写入时创建新版本
在 COW 表中,更新一个 File Group 时,Writer 读取旧 Base File、合并变更记录,再写出新的 Base File。
1 | Slice T1:Base T1 |
旧 Base File 不会被原地修改。只有 T2 对应的 Action 进入 COMPLETED 后,Reader 才会选择 T2 作为最新 File Slice;旧 Slice 在 Cleaner 回收前仍可用于历史读取。
COW 将合并成本放在写入路径,Reader 通常只需要读取最新 Base File。代价是即使只更新少量记录,也可能需要重写整个受影响的 Base File。
MOR:在当前 Slice 中累积变化
在 MOR 表中,对已有 File Group 的更新通常先写入 Log File:
1 | Base Instant T1 |
T2、T3 是 Delta Commit 的 Instant,但这些 Log 仍属于以 T1 为 Base Instant 的 File Slice,并不会各自创建新的 Slice。读取时有两条路径:
- Snapshot Query 合并 Base 与有效 Logs,返回最新状态;
- Read Optimized Query 只读取最近物化的 Base File,可能看不到尚未 Compaction 的更新。
当 T4 执行 Compaction 时,Hudi 将 Base T1 与 Logs T2/T3 合并为 Base T4,并在同一个 File Group 中创建新的 File Slice:
1 | Slice T1:Base T1 + Logs T2/T3 |
Compaction 改变的是物理布局,不应改变前后的逻辑记录状态。
COW 与 MOR 的权衡
| 维度 | COW | MOR |
|---|---|---|
| 更新方式 | 重写受影响的 Base File | 通常先写入 Log |
| 新 File Slice | 写出新 Base File 时产生 | 通常由 Compaction 或新 Base File 产生 |
| Snapshot 读取 | 最新 Base File | Base 与 Logs 合并 |
| 写放大 | 较高 | 较低 |
| 读放大 | 较低 | 取决于 Log 数量 |
| 后台维护 | 不需要 MOR Compaction | 需要控制 Compaction 周期 |
File Group 越大,文件数量越少,但 COW 重写成本和并发冲突范围可能越大;File Group 越小,更新粒度更细,却容易产生小文件和更多元数据。MOR 还需要控制 Log 的累积速度,避免把写入成本全部转移到查询和 Compaction。
Log Block 与 Record Merger:MOR 如何恢复最新状态
Hudi Log File 不是普通的 Avro 文件,而是由多个 Log Block 组成的原生容器。Block 的内容可以使用 Avro、HFile 或 Parquet 编码。一个 Log Block 的结构可以简化为:

Magic 和长度字段用于识别 Block 边界及损坏数据;Header 可以携带 Instant、Schema、Record Positions、Block Identifier 和 Partial Update 等信息。常见 Block 类型包括:
| Block | 作用 |
|---|---|
| Data Block | 保存完整或部分的 Insert、Update 记录 |
| Delete Block | 保存带有 Ordering Value 的删除 Tombstone |
| Command Block | ROLLBACK_BLOCK 可撤销前一个失败写入产生的 Block |
| CDC Block | 按 CDC 配置保存操作类型及 Before/After Image |
Corrupted Block 不会由 Writer 持久化,而是 Reader 对损坏或未完整写入 Block 的内部表示。MOR Snapshot Reader 会把 Base File、Data Block 和 Delete Tombstone 共同交给 Record Merger:

Record Merger 并不总是由 Ordering Field 决定。Hudi 主要提供三种 Merge Mode:
| Merge Mode | 冲突解决方式 |
|---|---|
| COMMIT_TIME_ORDERING | 更晚提交的记录胜出 |
| EVENT_TIME_ORDERING | Ordering Field 更大的记录胜出 |
| CUSTOM | 使用自定义 HoodieRecordMerger |
例如:
1 | (key=42, source_lsn=101, price=80) |
使用 EVENT_TIME_ORDERING 时,即使 source_lsn=102 最晚到达,最终仍保留 source_lsn=103 的版本。
对于需要处理乱序、重试和重复事件的 Event-time 或自定义 Merge,合并逻辑应保持确定性,并尽可能满足交换性、结合性和幂等性,避免写入、查询与 Compaction 得到不同结果(Hudi Technical Specification、Hudi Record Mergers)。
Index 与 Metadata Table:如何定位记录和文件
Index 介绍
Hudi 的 Upsert 需要先判断输入记录是 Insert 还是 Update。Index 负责将输入从:
1 | <Record Key, New Record> |
转换为:
1 | <Record Key, Partition Path, File Group ID> |
这个过程通常称为 tagLocation:

常见的定位方式包括:
| Index | 定位方式 | 主要权衡 |
|---|---|---|
| Bloom Index | 使用 Key Range 与 Bloom Filter 筛选候选文件 | 灵活,但随机更新可能扫描较多文件 |
| Bucket Index | 通过 Record Key 的 Hash 映射 Bucket | 定位稳定,但需要处理 Bucket 数量与数据倾斜 |
| Record-Level Index | 直接保存 Record Key 到 File Group 的映射 | 定位快,但需要维护索引状态 |
Record-Level Index 近似保存:
1 | Record Key → Partition Path + File Group ID |
它可以具有两种唯一性范围:
- Global RLI:Record Key 在全表范围唯一;
- Partitioned RLI:
Partition Path + Record Key在分区内唯一。
如果分区字段可能变化,例如 product_id=42 从 region=cn 移动到 region=us,分区内索引需要显式删除旧分区版本,否则两个分区可能同时保留该 Key。
Metadata Table:Hudi 的内部元数据存储
Metadata Table 位于 .hoodie/metadata,它本身是一张内部 MOR 表,用不同 Partition 保存文件与索引信息:

它主要解决两个扩展性问题:
- 避免 Reader 和 Writer 对对象存储执行大规模目录 Listing;
- 避免查询规划阶段读取大量 Data File Footer 才能获得统计信息。
Data Table 与 Metadata Table 的 Timeline 具有父子关联,使数据文件、文件列表和索引随写入共同演进。这里应与 Catalog 区分:Catalog 负责表发现、Schema 和权限入口,Metadata Table 则保存 Hudi 表内部的文件与索引状态(Hudi Table Metadata、Hudi Indexing)。
并发控制与 Table Services:如何安全地重写文件
Hudi 默认采用 Single Writer。启用多 Writer 后,OCC 允许多个 Writer 先在锁外生成文件,只在校验和提交阶段进行短暂协调:

OCC 的冲突判断可以概念化为:
1 | Current Write Set |
两个 Writer 修改不重叠的 File Group 时可以分别提交;写入范围重叠时,通常只有一个能够成功。实际冲突范围还取决于 Action 类型和配置的 Conflict Resolution Strategy。
Hudi 还使用 MVCC 协调 Writer 与 Table Services。NBCC 则允许受支持的 MOR 与 Bucket Index 组合并发写入同一 File Group,再由 Reader 和 Compactor 解决记录冲突,但它不适用于任意表类型、索引和后台操作(Hudi Concurrency Control)。持续 Upsert 会产生 Log、小文件、旧 File Slice 和 Timeline 历史,Table Services 负责处理这些存储债务:
| Table Service | 核心作用 |
|---|---|
| Compaction | 合并同一 File Group 的 Base 与 Logs,生成新 Base File |
| Log Compaction | 合并多个 Log File,降低 Log 扫描成本 |
| Clustering | 调整文件大小、排序和数据局部性 |
| Cleaning | 按保留策略回收旧 File Slice |
| Timeline Maintenance | 控制 Active Timeline 的长期规模 |
Compaction 与 Clustering 的区别最重要:

- Compaction 在同一个 File Group 内物化 Log,保持记录的逻辑状态;
- Clustering 通常通过
REPLACE_COMMIT用新的 File Groups 替换旧 File Groups; - Cleaning 只回收超过保留策略的历史文件,不参与 Record Merge;
- 这些操作主要改变物理布局,不应被下游解释为新的业务 CDC 事件。
生产实践:用 Hudi 构建订单实时湖表
订单表是典型的持续更新数据:订单会依次经历创建、支付、发货、完成或退款,数据库还可能产生重复、乱序和删除事件。
如果将每次变化直接追加为 Parquet 文件,数据湖中会同时存在同一订单的多个版本。Reader 必须扫描并合并全部历史记录,才能恢复最新状态;如果每天全量重写,又会带来较大的计算和写放大。
Hudi 的核心价值,就是在对象存储上持续维护一张可更新、可增量消费的订单表。
从数据库 CDC 到 Hudi
一条典型的生产链路如下:

表中可以保留以下关键字段:
1 | order_id |
Writer 收到更新后,先通过 Index 找到订单所在的 File Group。MOR 表通常将变化追加到当前 File Slice 的 Log File,随后把本次写入发布为一个 Completed Instant。
Reader 不直接根据文件是否存在判断数据是否有效,而是根据 Timeline 选择已经完成的 Action,再读取对应的 Base File 和 Log File。
Hudi 如何处理更新、乱序与失败
假设同一订单依次到达以下事件:
1 | T1:order_id=42,source_lsn=101,status=CREATED |
T3 虽然最后到达,但它的 source_lsn 小于 103。使用 Event-time Ordering 时,Record Merger 仍会保留 PAID,避免乱序事件覆盖新状态。
整个过程可以概括为:

如果 Writer 在文件写出后、Instant 完成前失败,这些文件虽然已经存在,但不会被 Reader 看见。后续 Writer 可以根据未完成 Instant 和 Heartbeat 判断写入是否失效,并通过 Rollback 清理残留文件;Cleaner 则主要负责回收超过保留策略的已提交历史版本。
随着 Log 不断累积,Compaction 会将 Base File 与 Logs 合并成新的 Base File。Compaction 改变的是物理布局,不应改变订单的最终业务状态。
Hudi 的核心价值
这个案例体现了 Hudi 的三个核心能力:
- 通过 Record Key、Index 和 File Group,将数据库的持续更新映射到对象存储;
- 通过 Timeline 隔离未完成写入,原子发布一批文件组成的新表状态;
- 同时提供 Snapshot Query 和 Incremental Query,让分析系统读取最新状态,让实时系统消费增量变化。
因此,Hudi 最具代表性的场景不是单纯保存离线文件,而是数据库 CDC、订单状态、用户画像和实时特征等持续更新的数据集。
总结
Hudi 并不是简单地在 Parquet 文件之上增加事务日志,而是围绕持续变化的数据构建了一套完整的湖表内核:Timeline 定义操作的提交与可见性,Index 将记录定位到 File Group,File Slice 组织记录的物理版本,Record Merger 与 Table Services 则负责合并变化并持续优化存储布局。它最核心的价值,是让 CDC、Upsert 和增量处理能够可靠地运行在开放存储之上。
但 Hudi 只是事务湖表的一种设计路线。面对“如何定义表状态、如何发布原子提交、如何组织更新以及如何控制读写放大”这些共同问题,不同系统给出了不同答案。
后续文章将继续分析三种代表性架构:Iceberg 如何通过 Snapshot、Manifest 与 Catalog Commit 管理超大规模分析表;Paimon 如何在 Snapshot 之下引入 Bucket 与 LSM Tree,构建面向流式主键更新的实时湖存储;以及 Lance 如何针对 AI 数据重新设计列式存储、随机访问与多模态索引。通过这些系统的对比,我们将进一步理解数据湖架构如何从通用分析,逐步演进到实时处理与 AI 数据基础设施。
