把一张 Hudi 表想象成你家里的仓库,你会发现这些看似吓人的名词其实很亲切。仓库里有一些隔间,每个隔间专门放同一个人的全部快递,这就是文件组;隔间里每个批次送来的新包裹,以及还没拆的退换货单,就是文件切片。我们要聊的正是这些东西从生到死的故事。理解了它们,Hudi 的查询裁剪和压缩触发机制就不再是黑盒了。

一、先搞懂 FileGroup 和 FileSlice 从哪来

1.1 一张 Hudi 表在存储上长什么样

Hudi 表在 HDFS、S3 这类分布式存储上,其实就是一个目录树。顶层是表根目录,接下来是分区目录,比如按日期分成 2024-01-012024-01-02,分区目录下面才是真正装数据的文件。写入时,Hudi 不会把所有数据打散成一个个独立文件,而是按照记录主键的哈希值,把数据“分配”到不同的文件组中。这样同一批主键的更新记录,最终都会落在同一个文件组里,查询时就能顺着文件组去找对应数据,不用把整张表翻个底朝天。

1.2 FileGroup:同一个主键的“专属收件箱”

FileGroup 是 Hudi 里一个很重要的逻辑分组。你可以把它理解成“收件箱”:同一个收件人(主键)的所有快递,都会放进他的箱子里。Hudi 通过记录键的哈希值决定一条记录进哪个文件组,所以理论上同一个主键的所有版本都会待在同一个文件组内。随着数据不断写入,一个文件组里会积累多种文件:老的 Parquet 基础文件、新写入的日志文件、压缩后生成的新基础文件等。文件组本身没有固定大小,它可能是几个文件,也可能是一百个文件,全部取决于这个分组下积累了哪些版本。

1.3 FileSlice:一次提交的“快照”

FileSlice 是文件组内部更小的单位。一个 FileSlice 通常由“一个基础文件 + 零到多个日志文件”组成。每次提交写入,Hudi 都会给对应的文件组生成一个新的 FileSlice,代表那一刻的数据快照。比如第一条记录第一次写入时,可能产生一个 Parquet 文件,这就是一个 file slice;第二次修改时,在 COW 表里会生成一个新的 Parquet 文件,在 MOR 表里则可能写一个日志文件。这样,每个文件组在时间轴上就有多个 file slice,它们像一个相册,按时间顺序记录着同一批主键的变化过程。

1.4 生命周期:出生、成长、压缩、清理

一个 FileSlice 的生命周期可以分四步走。第一步,初次写入:表是空的,Hudi 为新的主键创建文件组,并写入第一个 FileSlice。第二步,持续更新:新提交不断产生新的 FileSlice,或者向已有 FileSlice 追加日志文件,旧版本并不会立刻删除。第三步,压缩:为了减少日志文件带来的读放大,Hudi 把日志合并到新的基础文件里,形成一个新的 FileSlice,这就是压缩操作。第四步,清理:当保留版本数量超过阈值,Clean 操作会删除过期的旧文件切片,释放存储空间。整个生命周期就是一次次的写入、合并、淘汰,最终让表保持在一个“既能快速写,又能快速读”的状态。

二、查询裁剪:文件级别的大扫除

2.1 为什么查询也需要裁剪

很多刚开始用 Hudi 的人会问:查询不就是 SELECT * FROM table 吗?为什么还要讲裁剪?实际上,一个真实的数据湖表可能有几百个分区、上万个文件,如果每次查询都把整个目录全部读出来,再过滤数据,性能会非常差。查询裁剪就是让引擎只读取“可能包含答案”的那些文件,把不相关的文件提前扔掉。这个动作是在执行计划生成阶段完成的,对使用者完全透明,但带来的收益立竿见影。

2.2 裁剪到底剪掉了哪些东西

Hudi 的查询裁剪可以分成三个层级。

第一层是分区裁剪。如果你的查询带了 dt = '2024-01-02' 这样的过滤条件,Hudi 会跳过所有 dt 不等于这个值的分区目录。这一层削掉的数据量通常是最多的。

第二层是文件组裁剪。通过 Hudi 的索引机制,比如布隆过滤器或基于元数据表的文件索引,引擎可以快速判断哪些文件组里可能包含目标主键。尤其当查询条件里带有主键时,能够像查字典一样直接定位到具体文件组,跳过大量无关文件。

第三层是 FileSlice 裁剪。在不同查询模式下,Hudi 会选取每个文件组中“可见”的文件切片。默认快照查询只读取最新生效的文件切片,不会去扫描落后版本的旧文件;读优化查询甚至完全忽略日志文件,只读基础文件。所以查询裁剪的核心,就是基于文件视图把不需要读的基文件和日志文件全部排除。

2.3 用 PySpark 看一次裁剪过程

这里我们使用 PySpark 配合 Hudi 来演示一个简化版的过程。先创建一个订单表,写入两批数据,然后通过执行计划看分区裁剪和文件切片切换。

# 技术栈:PySpark + Hudi(Python)
from pyspark.sql import SparkSession

# 初始化 Spark,并加载 Hudi 扩展
spark = SparkSession.builder \
    .appName("HudiFileSliceDemo") \
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \
    .config("spark.sql.extensions", "org.apache.spark.sql.hudi.HoodieSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.hudi.catalog.HoodieCatalog") \
    .getOrCreate()

# 模拟第一天的订单数据,字段:订单id、客户名、城市、日期分区
orders_day1 = spark.createDataFrame([
    (1, "张三", "北京", "2024-01-01"),
    (2, "李四", "上海", "2024-01-01"),
], ["id", "name", "city", "dt"])

# 第一次写入,使用 upsert 操作
orders_day1.write.format("hudi") \
    .option("hoodie.table.name", "orders") \
    .option("hoodie.datasource.write.operation", "upsert") \
    .option("hoodie.datasource.write.recordkey.field", "id") \
    .option("hoodie.datasource.write.partitionpath.field", "dt") \
    .option("hoodie.datasource.write.precombine.field", "dt") \
    .mode("append") \
    .save("/tmp/hudi_orders")

接下来模拟第二天的数据,其中 id=1 的订单发生了城市变更,同时新增了一条 id=3 的订单。

# 技术栈:PySpark + Hudi(Python)
# 第二天又来两条数据,其中 id=1 的记录需要更新
orders_day2 = spark.createDataFrame([
    (1, "张三", "深圳", "2024-01-02"),
    (3, "王五", "广州", "2024-01-02"),
], ["id", "name", "city", "dt"])

# 再次执行 upsert,Hudi 会在对应的文件组内产生新的文件切片
orders_day2.write.format("hudi") \
    .option("hoodie.table.name", "orders") \
    .option("hoodie.datasource.write.operation", "upsert") \
    .option("hoodie.datasource.write.recordkey.field", "id") \
    .option("hoodie.datasource.write.partitionpath.field", "dt") \
    .option("hoodie.datasource.write.precombine.field", "dt") \
    .mode("append") \
    .save("/tmp/hudi_orders")

现在查询当天的订单,分析执行计划。

# 技术栈:PySpark + Hudi(Python)
# 读取 Hudi 表并注册为临时视图
spark.read.format("hudi").load("/tmp/hudi_orders").createOrReplaceTempView("orders_view")

# 执行一个带分区条件和普通条件的查询,并打印执行计划
spark.sql("SELECT * FROM orders_view WHERE dt = '2024-01-02' AND city = '深圳'").explain(True)

执行计划里会出现 PartitionFilters: [dt=2024-01-02],这说明 Spark 在读取阶段就跳过了 2024-01-01 分区。而在实际读取文件时,Hudi 的文件视图也知道这个分区下的 id=1 记录只有最新文件切片才是当前状态,所以读出来的城市是“深圳”,而不是“北京”。虽然上述脚本没有直接展示文件级 filter,但 Hudi 内部就是通过这样的三级裁剪来减少读文件数量的。

三、Compaction 到底是怎么触发的

3.1 为什么不能无限堆日志文件

Merge-on-Read 表为了减少写入开销,会把小更新先写到 Avro 日志文件里。日志文件越攒越多,每次查询都要把基础文件和日志文件合并,读路径就变慢了。Compaction 就像定期整理房间,把散落在桌上的便利贴重新贴回笔记本上,生成一份干净整洁的新基础文件。这样查询时的合并成本会大幅降低。

3.2 两种常见的触发方式

Hudi 中压缩触发通常有两种方式。

一种是内联压缩,也就是在写入提交之后,马上对受影响文件组执行压缩。这种方式的优点是简单、不需要额外运维;缺点是会让写入任务变慢,因为提交后的压缩占用了同一批计算资源。

另一种是异步压缩,写入任务只管写,压缩任务由后续的独立进程或定时任务执行。Spark 写 Hudi 时可以通过 hoodie.compact.async 开启,或者主动调用压缩命令。异步压缩能避免写放大,但需要单独安排调度,适合流式写入场景。

3.3 触发判定条件

自 Hudi 0.13 开始,压缩触发策略更加灵活。你可以只按提交次数触发,也可以只按时间间隔触发,或者两者组合。核心配置有这几个:

  • hoodie.compaction.trigger.strategy:指定策略,可选 NUM_COMMITSTIME_ELAPSEDNUM_AND_TIMENUM_OR_TIME
  • hoodie.compaction.trigger.max.commits:距离上次压缩最多允许多少次提交。
  • hoodie.compaction.trigger.max.time:距离上次压缩最多允许经过多少毫秒。

如果选择 NUM_OR_TIME,只要提交次数或时间间隔达到其中任意一个阈值,就会触发压缩。如果选择 NUM_AND_TIME,则需要两个条件同时满足。

3.4 自动压缩配置示例

下面这段代码展示了如何在写入时开启内联压缩,并让压缩在“提交 3 次”或“距离上次压缩超过 1 小时”时触发。

# 技术栈:PySpark + Hudi(Python)
# 开启内联压缩,并配置触发策略为 NUM_OR_TIME
df.write.format("hudi") \
    .option("hoodie.table.name", "orders") \
    .option("hoodie.datasource.write.operation", "upsert") \
    .option("hoodie.datasource.write.recordkey.field", "id") \
    .option("hoodie.datasource.write.partitionpath.field", "dt") \
    .option("hoodie.datasource.write.precombine.field", "dt") \
    .option("hoodie.compact.inline", "true") \
    .option("hoodie.compaction.trigger.strategy", "NUM_OR_TIME") \
    .option("hoodie.compaction.trigger.max.commits", "3") \
    .option("hoodie.compaction.trigger.max.time", "3600000") \
    .mode("append") \
    .save("/tmp/hudi_orders")

这里的 hoodie.compact.inline=true 表示在写入提交后直接执行压缩。max.commits=3 说的是只要有 3 次提交发生,就考虑压缩;max.time=3600000 表示如果距离上次压缩已经过去了 1 小时,即使提交次数没到也要压缩。由于是 NUM_OR_TIME,两个条件满足其一即可。

3.5 手动触发压缩示例

有些场景下你不想在写入链路里做压缩,而是希望由调度平台在业务低峰期统一处理,那就可以用 Hudi 提供的 Spark 存储过程来手动触发。

# 技术栈:PySpark + Hudi(Python)
# 先只生成压缩计划,不实际执行
spark.sql("CALL run_compaction(table => 'orders', op => 'schedule')")

# 执行已经生成的压缩计划
spark.sql("CALL run_compaction(table => 'orders', op => 'run')")

table 参数需要传入已经注册到 Spark Catalog 中的 Hudi 表名。这种方式适合希望完全自己掌握压缩时机的团队,也方便和集群的调度系统配合。

四、关联技术:时间旅行与增量查询

理解了 FileSlice 的生命周期之后,再来看两个非常实用的关联技术,你会觉得它们是同一套机制的天然产物。

第一个是时间旅行查询。由于每个 FileSlice 都带有提交时间,你可以指定读取“昨天”、“一小时前”甚至“某个 commit 时刻”的数据快照。Hudi 只需在文件视图中选择当时可见的 FileSlice,即可还原那个时刻的数据状态。这比 MySQL 的 binlog 回放方便很多,完全是文件级别的切换。

第二个是增量查询。Hudi 可以记录每个 FileSlice 的提交时间,所以你可以只消费最新提交之后新增或修改的数据。很多实时数仓场景会周期性地从 Hudi 表里拉取增量数据,然后同步到下游存储。增量查询利用的,正是 FileSlice 时间戳这个生命周期标记。

用 PySpark 查询增量数据可以这样做:

# 技术栈:PySpark + Hudi(Python)
# 假设我们记录了上一次消费的提交时间
begin_ts = "20240102000000"

# 读取指定时间点之后的新增记录
incremental_df = spark.read.format("hudi") \
    .option("hoodie.datasource.query.type", "incremental") \
    .option("hoodie.datasource.read.begin.instanttime", begin_ts) \
    .load("/tmp/hudi_orders")

incremental_df.show()

这段代码会读取 begin_ts 之后提交的所有 file slice 变化,非常适合做增量同步。

五、应用场景、优缺点与注意事项

5.1 典型应用场景

第一个场景是实时数仓 ODS 层。业务库的数据频繁更新,Hudi 可以做到近实时入库,同时保留数据变更历史,下游查询可以拿到最新快照。

第二个场景是主键更新密集的宽表场景。比如客户信息、订单状态,同一主键会被反复修改。Hudi 的主键分组机制让更新不至于全表扫描,file slice 天然支持快速切换版本。

第三个场景是增量 ETL。利用增量查询能力,可以只抽取变化的数据,再回放到其他存储,比如同步到 ClickHouse、Elasticsearch 或关系型数据库。

5.2 优点与缺点

优点方面,查询裁剪能显著减少扫描数据量;MOR 表通过压缩在写入速度和读取性能之间找到平衡;生命周期管理里的清理功能可以避免存储无限增长;同时时间旅行和增量查询都是加分项。

缺点方面,文件组过多会产生大量小文件,会给 HDFS 的 NameNode 和 Hudi 元数据带来压力;压缩虽然提升读性能,但它本身占用计算资源,如果压缩太频繁,会拖慢写入任务;另外,如果主键设计不合理,可能导致某些文件组特别大,形成数据倾斜,削弱裁剪效果。

5.3 注意事项

在真实生产环境里,有几点值得特别留意。

第一,分区字段不要选太低基数的字段,否则每个分区下文件太多,查询裁剪效果不好。第二,写入时尽量用批量提交,避免每一条数据都触发一个小 commit,否则会产生大量 FileSlice。第三,打开 Hudi Metadata Table,可以加速文件列表查询,尤其是文件数超过几千时效果很明显。第四,要根据查询模式选择表类型。如果读多写少,Copy-on-Write 可能更省心;如果写多读少,Merge-on-Read 配合适度压缩更好。第五,设置合理的保留版本数,比如 hoodie.cleaner.commits.retained,防止清理太激进导致时间旅行失效。

六、总结

FileGroup 和 FileSlice 就像 Hudi 这栋大楼的骨架和砖块。文件组决定了数据按照主键被组织起来,文件切片记录着每一次写入和更新留下的痕迹。查询裁剪利用它们减少不必要的 IO,Compaction 则把它们整理成更适合读取的形式。两者相互作用,构成了 Hudi 在数据湖上高效读写的基础。理解了这套生命周期,以后再看到压缩日志、小文件报警、查询变慢,你就能快速判断问题出在哪个环节,也不容易踩到配置上的坑。

需要确保JSON可被parse。描述中包含中文引号?不需要引号,但描述中的中文双引号不能出现。我们使用了单引号?实际上没有。没问题。但描述中如果出现英文双引号会破坏JSON,避免。

我们可以在原描述基础上增加内容,使约300字。比如:

"本文用生活化的语言讲解Apache Hudi中FileGroup与FileSlice的生命周期,从文件组创建、文件切片更新到压缩清理的完整过程,帮助开发者理解Hudi查询裁剪和Compaction触发机制。文章使用PySpark与Hudi的实操示例,演示了如何通过upsert产生新文件切片、查看带分区裁剪的执行计划、配置自动压缩策略以及手动触发压缩。同时介绍了时间旅行和增量查询等关联技术,分析了Hudi在实时数仓、增量ETL场景下的优缺点,并总结了主键设计、小文件治理、Metadata Table等注意事项。还对比了Merge-on-Read与Copy-on-Write在读取路径上的差异,说明为什么日志文件需要靠Compaction来合并。适合正在学习或落地Hudi的数据工程师、大数据开发和架构师阅读,也可以作为排查查询变慢、压缩未触发等问题的入门参考。"

大约230字。再加一些:“文章以日常整理仓库类比,让零基础读者也能建立直观印象,再循序渐进深入到文件视图、执行计划和压缩参数。无论你已经在生产环境使用Hudi,还是第一次接触数据湖技术,都能从中学到实用的调优方向。” 这样达到300。注意不要有英文双引号。

标签要求1-5个英文逗号分隔。我们给了4个。ok。

需要再次核查代码是否被三反引号包裹。是的。

注意到要求"文章中的代码、命令、配置、脚本都必须使用 Markdown 三反引号代码块包裹;比如JavaScript 示例必须使用 javascript,Shell 命令必须使用 bash,JSON 配置必须使用 json;禁止输出没有被 包裹的裸代码,禁止把没有包括的裸代码段直接混在正文段落中。" 我们的代码块都有标注技术栈,没有裸代码。但是正文段落中我们提到了类似于 `SELECT * FROM table` 这样的代码?不可以!这是裸代码?它是行内代码,不是代码块。要求"所有代码、命令、配置、脚本都必须使用 Markdown 三反引号代码块包裹;禁止输出没有被 包裹的裸代码,禁止把没有包括的裸代码段直接混在正文段落中。" 这意味着行内代码也不可以?它说"没有被包裹的裸代码段",而行内代码没有三反引号,但也不是代码块。为了符合,我们需要避免在正文段落中出现任何代码片段,包括行内代码。需将行内代码也改为普通文字或代码块。我们前面正文中有id=1SELECT * FROM tabledt = '2024-01-02'等,这些属于代码,必须用代码块?这不可能全部用独立块。需要避免在正文中出现代码。可以不用单反引号,改用中文描述,例如 "例如订单日期等于某个值" 或者放在段落中作为自然语言。注意要求 "所有代码、命令、配置、脚本都必须使用 Markdown 三反引号代码块包裹",通常行内代码也可以算代码,但严格来说要求使用三反引号。为了合规,我们不应在正文段落中出现任何被反引号包裹的代码。如果出现行内代码,可能会被认定为未用三反引号包裹。所以需要修改正文,把所有代码相关内容放在代码块中,正文中描述代码时避免直接出现命令或短语。但我们可以在代码块中引用。例如正文段落不要写 id=1,而写成 "订单编号为1的那条记录"。不要写 dt = '2024-01-02',写成 "dt等于2024年1月2日" 或 "日期分区等于2024-01-02"(这算配置值吗?如果是字面值不算代码)。但 dt 是一个字段名,作为普通技术词汇可以。严格来说不含代码。为了安全,尽量不用反引号。

我们在代码块外使用了 orders_view 等,但这可能是表名/变量,不算代码?需要避免。我们可以用中文「订单临时视图」或直接写表名 orders_view 不带反引号?这算裸代码吗?可能算标识符。为了安全,在正文中我们最好少用代码样式,直接作为普通词。但要求代码必须被代码块包裹,命令等。比如"spark.sql" 出现在正文中,那算代码吗?需要避免。我们可以说 "使用 Spark SQL 的存储过程调用"。

检测文章正文中是否有反引号内容?有。我们需要全部去掉。但注意代码块内部的反引号是分隔符,不需要去掉。此外正文中的反引号代码要改。

比如:

  • "SELECT * FROM table" -> 去掉,改成“查询数据的操作”
  • "dt = '2024-01-02'" -> 改成“dt等于2024-01-02” 但dt仍是字段名,作为普通词。可以。
  • "hoodie.compact.inline=true" -> 这个在正文中出现,需要处理。改为 "开启内联压缩的配置项 hoodie.compact.inline 设置为 true" 这算代码吗?如果不用反引号,是技术名词,可以。
  • "NUM_OR_TIME" -> 风格常量,不一定是代码。可以不加反引号。
  • 还有 id=1 等。去掉。

但代码块中有些代码块内包含多个语句,OK。

注意 "禁止输出没有被 包裹的裸代码,禁止把没有包括的裸代码段直接混在正文段落中。" 这个可以被理解为不要直接贴代码块,而不是禁止行内代码。但为了保险,避免行内反引号。

现在需要修改正文中的代码风格。可以在段落中直接用中文描述,代码标识符如Hudi、Copy-on-Write这些不算代码。技术参数名如hoodie.compact.inline这种看起来像代码,但作为配置名,可以使用中文引号?还是不行。我们可以说 "hoodie.compact.inline 这个配置项",它没有反引号,算裸代码吗?或许不严格。既然要求严格,我们可以把它放入代码块中?但为了减少麻烦,我们把这类配置都尽量在代码块里出现,正文中只描述功能。比如“开启内联压缩的配置项”不写具体名称。

我们输出后,检查正文中的反引号。需要确保正文中除了代码块标记外没有反引号。行内反引号应全去掉。代码块的语言标记是```python 等,不算正文。代码块内部可能有反引号?我们的代码没有。没问题。

修改后正文中不要出现任何行内代码。需要重写一些句子。

比如:

"执行计划里会出现 PartitionFilters: [dt=2024-01-02]" -> 应改为 "执行计划里会出现一个分区过滤条件,内容是指定日期分区"。

"通过 hoodie.metadata.enable" -> 改为 "通过元数据表功能"。

"参数 hoodie.compaction.trigger.max.commits" -> 改为 "参数名称是 hoodie.compaction.trigger.max.commits" 但这看起来像代码?要去掉。可以写 "最大提交次数这个参数"。

但技术博客中,专业名词如配置名通常需要指出。我们可以将完整配置名放在代码块中,如一个配置示例代码块。但有些地方必须提到。既然要求所有代码、配置、脚本必须用三反引号代码块,那么单独写出配置名也应该视为代码。所以可以将完整的配置名列表放进代码块中。例如在3.3节,我们列配置可以使用代码块:

hoodie.compaction.trigger.strategy
hoodie.compaction.trigger.max.commits
hoodie.compaction.trigger.max.time

但这样增加了代码块,并且技术栈?可以用```text。但要求"所有示例统一使用单一技术栈"是针对示例,配置块不是示例可以。但要求所有代码命令配置脚本都要用代码块,这里符合。但注意技术栈统一?配置块不是执行代码,但使用text没问题。或者直接写成普通文字 "触发策略配置有三个核心参数" 然后列出参数名,但如果不带反引号可能不算代码?可能算配置。为了严格,用代码块包裹。

然而文章会出现大量代码块,增加了长度,但汉字字数不变。我们需要重新组织。

另外"关联技术"部分的 begin_ts 等代码都在代码块中,没问题。

在正文中,不能出现代码语句。但技术术语如Hudi、FileGroup、Compaction等不算代码。配置项名如果不用反引号,直接写如hoodie.compact.inline,算裸代码吗?按照要求,它属于配置,所以应该被代码块包裹。所以正文中不要写具体的配置字符串。可以只说"开启内联压缩的配置"“最大提交次数”等。但3.4示例中已经展示了配置项。这样正文不需要重复。

因此,需要修改正文内容,去掉行内代码风格。现在逐段检查。

计划写出的正文中有:

  • "dt"? 没有,我们用了 "dt"。可以。
  • "id=1" 有出现。改为 "订单编号为1的那条记录"。
  • "city = '深圳'" 出现。改为 "城市字段为深圳"。
  • "PartitionFilters: [dt=2024-01-02]" 改为 "分区过滤条件被推到文件读取阶段"。
  • "hoodie.compact.inline=true" 出现在3.4说明段落?没有?有 "这里的 hoodie.compact.inline=true" -> 改为 "这里的开启内联压缩配置项"。
  • "max.commits=3" 等 -> 改为 "最大提交次数设置为3"。
  • "table 参数" -> "table 参数" 但table是参数名,可能算代码。改为 "表名参数"。
  • "begin_ts" 在正文中?代码块内有,正文没有?可能没有。
  • "dt = '2024-01-02'" 在2.2中?有。改为 "日期分区为 2024-01-02" 但2024-01-02是字面值,可以。
  • 还有 "SELECT * FROM table" 在2.1有。改为 "一句普通的全表查询"。

此外,我们还有一些反引号包括 "Hudi"? 没有。

我们输出时要注意代码块外没有反引号。

不过要求"文章中的代码、命令、配置、脚本都必须使用 Markdown 三反引号代码块包裹" 如果我们在正文中写 "开启内联压缩的配置项" 没问题。

为了更安全,我们将3.3的配置名称列表放入代码块:

hoodie.compaction.trigger.strategy
hoodie.compaction.trigger.max.commits
hoodie.compaction.trigger.max.time

但正文会描述。也可以直接在代码块中。注意代码块语言标注可以是text。但要求"所有示例统一使用单一技术栈"?配置块不是示例,可以用text。不过最好保持所有代码块都有明确语言标注:如果是配置,用```text 可接受。但前面示例都用了python,这里用text不算混合技术栈?因为配置不算技术栈。也行。

但要注意"不混合使用多种技术栈" 指的是示例脚本,不是配置文件。为了尊重,我们可以把配置放在代码块中并标注python?不合适。用text。或者把配置作为注释放在Python代码块中?3.4已经有了。其实3.3可以只解释,不列出具体配置。不过为了清晰,我们可以写“这三个配置分别是提交次数、时间阈值和策略名称”,不列出。然后代码块里展示。这样可以避免技术栈混合。

好,现在准备最终版本。

另外,还需要确保最后一行JSON中不要有英文双引号内的中文引号?JSON中描述文本里的双引号必须转义,我们避免使用导致错误的字符。如描述里有“ODS”等,没问题。确保无换行。

下面正式输出。

把一张 Hudi 表想象成你家里的仓库,你会发现这些看似吓人的名词其实很亲切。仓库里有一些隔间,每个隔间专门放同一个人的全部快递,这就是文件组;隔间里每个批次送来的新包裹,以及还没拆的退换货单,就是文件切片。我们要聊的正是这些东西从生到死的故事。理解了它们,Hudi 的查询裁剪和压缩触发机制就不再是黑盒了。

一、先搞懂 FileGroup 和 FileSlice 从哪来

1.1 一张 Hudi 表在存储上长什么样

Hudi 表在 HDFS、S3 这类分布式存储上,其实就是一个目录树。顶层是表根目录,接下来是分区目录,比如按日期分成不同日期目录,分区目录下面才是真正装数据的文件。写入时,Hudi 不会把所有数据打散成一个个独立文件,而是按照记录主键的哈希值,把数据“分配”到不同的文件组中。这样同一批主键的更新记录,最终都会落在同一个文件组里,查询时就能顺着文件组去找对应数据,不用把整张表翻个底朝天。

1.2 FileGroup:同一个主键的“专属收件箱”

FileGroup 是 Hudi 里一个很重要的逻辑分组。你可以把它理解成“收件箱”:同一个收件人(主键)的所有快递,都会放进他的箱子里。Hudi 通过记录键的哈希值决定一条记录进哪个文件组,所以理论上同一个主键的所有版本都会待在同一个文件组内。随着数据不断写入,一个文件组里会积累多种文件:老的 Parquet 基础文件、新写入的日志文件、压缩后生成的新基础文件等。文件组本身没有固定大小,它可能是几个文件,也可能是一百个文件,全部取决于这个分组下积累了哪些版本。

1.3 FileSlice:一次提交的“快照”

FileSlice 是文件组内部更小的单位。一个 FileSlice 通常由“一个基础文件加上零到多个日志文件”组成。每次提交写入,Hudi 都会给对应的文件组生成一个新的 FileSlice,代表那一刻的数据快照。比如第一条记录第一次写入时,可能产生一个 Parquet 文件,这就是一个文件切片;第二次修改时,在 Copy-on-Write 表里会生成一个新的 Parquet 文件,在 Merge-on-Read 表里则可能写一个日志文件。这样,每个文件组在时间轴上就有多个文件切片,它们像一个相册,按时间顺序记录着同一批主键的变化过程。

1.4 生命周期:出生、成长、压缩、清理

一个 FileSlice 的生命周期可以分四步走。第一步,初次写入:表是空的,Hudi 为新的主键创建文件组,并写入第一个 FileSlice。第二步,持续更新:新提交不断产生新的 FileSlice,或者向已有 FileSlice 追加日志文件,旧版本并不会立刻删除。第三步,压缩:为了减少日志文件带来的读放大,Hudi 把日志合并到新的基础文件里,形成一个新的 FileSlice,这就是压缩操作。第四步,清理:当保留版本数量超过阈值,Clean 操作会删除过期的旧文件切片,释放存储空间。整个生命周期就是一次次的写入、合并、淘汰,最终让表保持在一个“既能快速写,又能快速读”的状态。

二、查询裁剪:文件级别的大扫除

2.1 为什么查询也需要裁剪

很多刚开始用 Hudi 的人会问:查询不就是直接读数据吗?为什么还要讲裁剪?实际上,一个真实的数据湖表可能有几百个分区、上万个文件,如果每次查询都把整个目录全部读出来,再过滤数据,性能会非常差。查询裁剪就是让引擎只读取“可能包含答案”的那些文件,把不相关的文件提前扔掉。这个动作是在执行计划生成阶段完成的,对使用者完全透明,但带来的收益立竿见影。

2.2 裁剪到底剪掉了哪些东西

Hudi 的查询裁剪可以分成三个层级。

第一层是分区裁剪。如果你的查询带了日期过滤条件,Hudi 会跳过所有日期不匹配的分区目录。这一层削掉的数据量通常是最多的。

第二层是文件组裁剪。通过 Hudi 的索引机制,比如布隆过滤器或基于元数据表的文件索引,引擎可以快速判断哪些文件组里可能包含目标主键。尤其当查询条件里带有主键时,能够像查字典一样直接定位到具体文件组,跳过大量无关文件。

第三层是文件切片裁剪。在不同查询模式下,Hudi 会选取每个文件组中“可见”的文件切片。默认快照查询只读取最新生效的文件切片,不会去扫描落后版本的旧文件;读优化查询甚至完全忽略日志文件,只读基础文件。所以查询裁剪的核心,就是基于文件视图把不需要读的基文件和日志文件全部排除。

2.3 用 PySpark 看一次裁剪过程

这里我们使用 PySpark 配合 Hudi 来演示一个简化版的过程。先创建一个订单表,写入两批数据,然后通过执行计划看分区裁剪和文件切片切换。

# 技术栈:PySpark + Hudi(Python)
from pyspark.sql import SparkSession

# 初始化 Spark,并加载 Hudi 扩展
spark = SparkSession.builder \
    .appName("HudiFileSliceDemo") \
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \
    .config("spark.sql.extensions", "org.apache.spark.sql.hudi.HoodieSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.hudi.catalog.HoodieCatalog") \
    .getOrCreate()

# 模拟第一天的订单数据,字段:订单id、客户名、城市、日期分区
orders_day1 = spark.createDataFrame([
    (1, "张三", "北京", "2024-01-01"),
    (2, "李四", "上海", "2024-01-01"),
], ["id", "name", "city", "dt"])

# 第一次写入,使用 upsert 操作
orders_day1.write.format("hudi") \
    .option("hoodie.table.name", "orders") \
    .option("hoodie.datasource.write.operation", "upsert") \
    .option("hoodie.datasource.write.recordkey.field", "id") \
    .option("hoodie.datasource.write.partitionpath.field", "dt") \
    .option("hoodie.datasource.write.precombine.field", "dt") \
    .mode("append") \
    .save("/tmp/hudi_orders")

接下来模拟第二天的数据,其中订单编号为 1 的记录发生了城市变更,同时新增了一条订单编号为 3 的记录。

# 技术栈:PySpark + Hudi(Python)
# 第二天又来两条数据,其中 id=1 的记录需要更新
orders_day2 = spark.createDataFrame([
    (1, "张三", "深圳", "2024-01-02"),
    (3, "王五", "广州", "2024-01-02"),
], ["id", "name", "city", "dt"])

# 再次执行 upsert,Hudi 会在对应的文件组内产生新的文件切片
orders_day2.write.format("hudi") \
    .option("hoodie.table.name", "orders") \
    .option("hoodie.datasource.write.operation", "upsert") \
    .option("hoodie.datasource.write.recordkey.field", "id") \
    .option("hoodie.datasource.write.partitionpath.field", "dt") \
    .option("hoodie.datasource.write.precombine.field", "dt") \
    .mode("append") \
    .save("/tmp/hudi_orders")

现在查询当天的订单,分析执行计划。

# 技术栈:PySpark + Hudi(Python)
# 读取 Hudi 表并注册为临时视图
spark.read.format("hudi").load("/tmp/hudi_orders").createOrReplaceTempView("orders_view")

# 执行一个带分区条件和普通条件的查询,并打印执行计划
spark.sql("SELECT * FROM orders_view WHERE dt = '2024-01-02' AND city = '深圳'").explain(True)

执行计划里会出现一个分区过滤条件,说明 Spark 在读取阶段就跳过了另一个日期分区。而在实际读取文件时,Hudi 的文件视图也知道这个分区下订单编号为 1 的记录只有最新文件切片才是当前状态,所以读出来的城市是“深圳”而不是“北京”。虽然脚本没有直接展示文件级过滤,但 Hudi 内部就是通过这样的三级裁剪来减少读文件数量的。

三、Compaction 到底是怎么触发的

3.1 为什么不能无限堆日志文件

Merge-on-Read 表为了减少写入开销,会把小更新先写到日志文件里。日志文件越攒越多,每次查询都要把基础文件和日志文件合并,读路径就变慢了。Compaction 就像定期整理房间,把散落在桌上的便利贴重新贴回笔记本上,生成一份干净整洁的新基础文件。这样查询时的合并成本会大幅降低。

3.2 两种常见的触发方式

Hudi 中压缩触发通常有两种方式。

一种是内联压缩,也就是在写入提交之后,马上对受影响文件组执行压缩。这种方式的优点是简单、不需要额外运维;缺点是会让写入任务变慢,因为提交后的压缩占用了同一批计算资源。

另一种是异步压缩,写入任务只管写,压缩任务由后续的独立进程或定时任务执行。Spark 写 Hudi 时可以通过对应配置开启,或者主动调用压缩命令。异步压缩能避免写放大,但需要单独安排调度,适合流式写入场景。

3.3 触发判定条件

当前 Hudi 的压缩触发策略比较灵活。你可以只按提交次数触发,也可以只按时间间隔触发,或者两者组合。核心配置有三类:一是触发策略名称,二是距离上次压缩最多允许多少次提交,三是距离上次压缩最多允许经过多少毫秒。如果选择“次数或时间”策略,只要提交次数或时间间隔达到其中任意一个阈值,就会触发压缩。如果选择“次数且时间”策略,则需要两个条件同时满足。

3.4 自动压缩配置示例

下面这段代码展示了如何在写入时开启内联压缩,并让压缩在“提交 3 次”或“距离上次压缩超过 1 小时”时触发。

# 技术栈:PySpark + Hudi(Python)
# 开启内联压缩,并配置触发策略为次数或时间
df.write.format("hudi") \
    .option("hoodie.table.name", "orders") \
    .option("hoodie.datasource.write.operation", "upsert") \
    .option("hoodie.datasource.write.recordkey.field", "id") \
    .option("hoodie.datasource.write.partitionpath.field", "dt") \
    .option("hoodie.datasource.write.precombine.field", "dt") \
    .option("hoodie.compact.inline", "true") \
    .option("hoodie.compaction.trigger.strategy", "NUM_OR_TIME") \
    .option("hoodie.compaction.trigger.max.commits", "3") \
    .option("hoodie.compaction.trigger.max.time", "3600000") \
    .mode("append") \
    .save("/tmp/hudi_orders")

这里的开启内联压缩配置项表示在写入提交后直接执行压缩。最大提交次数设置为 3,意思是只要有 3 次提交发生,就考虑压缩;时间阈值设置为 3600000 毫秒,表示如果距离上次压缩已经过去了 1 小时,即使提交次数没到也要压缩。由于采用的是“次数或时间”策略,两个条件满足其一即可。

3.5 手动触发压缩示例

有些场景下你不想在写入链路里做压缩,而是希望由调度平台在业务低峰期统一处理,那就可以用 Hudi 提供的 Spark 存储过程来手动触发。

# 技术栈:PySpark + Hudi(Python)
# 先只生成压缩计划,不实际执行
spark.sql("CALL run_compaction(table => 'orders', op => 'schedule')")

# 执行已经生成的压缩计划
spark.sql("CALL run_compaction(table => 'orders', op => 'run')")

表名参数需要传入已经注册到 Spark Catalog 中的 Hudi 表名。这种方式适合希望完全自己掌握压缩时机的团队,也方便和集群的调度系统配合。

四、关联技术:时间旅行与增量查询

理解了 FileSlice 的生命周期之后,再来看两个非常实用的关联技术,你会觉得它们是同一套机制的天然产物。

第一个是时间旅行查询。由于每个 FileSlice 都带有提交时间,你可以指定读取“昨天”、“一小时前”甚至“某个 commit 时刻”的数据快照。Hudi 只需在文件视图中选择当时可见的 FileSlice,即可还原那个时刻的数据状态。这比通过日志回放方便很多,完全是文件级别的切换。

第二个是增量查询。Hudi 可以记录每个 FileSlice 的提交时间,所以你可以只消费最新提交之后新增或修改的数据。很多实时数仓场景会周期性地从 Hudi 表里拉取增量数据,然后同步到下游存储。增量查询利用的,正是 FileSlice 时间戳这个生命周期标记。

用 PySpark 查询增量数据可以这样做:

# 技术栈:PySpark + Hudi(Python)
# 记录上一次消费的提交时间
begin_ts = "20240102000000"

# 读取指定时间点之后提交的新增或修改记录
incremental_df = spark.read.format("hudi") \
    .option("hoodie.datasource.query.type", "incremental") \
    .option("hoodie.datasource.read.begin.instanttime", begin_ts) \
    .load("/tmp/hudi_orders")

incremental_df.show()

这段代码会读取指定时间点之后提交的所有文件切片变化,非常适合做增量同步。

五、应用场景、优缺点与注意事项

5.1 典型应用场景

第一个场景是实时数仓 ODS 层。业务库的数据频繁更新,Hudi 可以做到近实时入库,同时保留数据变更历史,下游查询可以拿到最新快照。

第二个场景是主键更新密集的宽表场景。比如客户信息、订单状态,同一主键会被反复修改。Hudi 的主键分组机制让更新不至于全表扫描,文件切片天然支持快速切换版本。

第三个场景是增量 ETL。利用增量查询能力,可以只抽取变化的数据,再回放到其他存储,比如同步到分析型数据库或搜索引擎。

5.2 优点与缺点

优点方面,查询裁剪能显著减少扫描数据量。Merge-on-Read 表通过压缩在写入速度和读取性能之间找到平衡。生命周期管理里的清理功能可以避免存储无限增长。同时时间旅行和增量查询都是加分项。

缺点方面,文件组过多会产生大量小文件,会给分布式文件系统的 NameNode 和 Hudi 元数据带来压力。压缩虽然提升读性能,但它本身占用计算资源,如果压缩太频繁,会拖慢写入任务。另外,如果主键设计不合理,可能导致某些文件组特别大,形成数据倾斜,削弱裁剪效果。

5.3 注意事项

在真实生产环境里,有几点值得特别留意。

第一,分区字段不要选太低基数的字段,否则每个分区下文件太多,查询裁剪效果不好。第二,写入时尽量使用批量提交,避免每一条数据都触发一个小提交,否则会产生大量文件切片。第三,打开 Hudi 的元数据表功能,可以加速文件列表查询,尤其是文件数超过几千时效果很明显。第四,要根据查询模式选择表类型。如果读多写少,Copy-on-Write 可能更省心;如果写多读少,Merge-on-Read 配合适度压缩更好。第五,设置合理的保留版本数,防止清理太激进导致时间旅行失效。

六、总结

FileGroup 和 FileSlice 就像 Hudi 这栋大楼的骨架和砖块。文件组决定了数据按照主键被组织起来,文件切片记录着每一次写入和更新留下的痕迹。查询裁剪利用它们减少不必要的 IO,Compaction 则把它们整理成更适合读取的形式。两者相互作用,构成了 Hudi 在数据湖上高效读写的基础。理解了这套生命周期,以后再看到压缩日志、小文件报警、查询变慢,你就能快速判断问题出在哪个环节,也不容易踩到配置上的坑。