hudi源码走读
- index 用来检索哪部分文件的?
- Flink dataStream.transform 的作用
- Flink hudi pipeline
- hoodie key生成规则:recordKey + partition Key
- recordKeyField = primary key
- FileId如何生成的?
km文章-Apache Hudi 快速入门
https://km.woa.com/group/51977/articles/show/510278
- 增量更新 如何发现变动的数据,定期启动下游任务
- 增量更新 实时任务如何触发
- 现有小时粒度更新,如何处理延迟数据
- 面向应用层的数据每天输出一次,如何确定时机
- 监管场景 重点需要了解
- Flink写hudi回放时,如何确保exactly once?
- timeline 支持的增量查询时间区间如何选择
- 数据已经cleans了,还能不能找到旧的timeline
- parquet文件大小与查询性能关系
- 索引是用来定位文件组,文件组是个啥?
- record key如何确定?
- hudi upsert如何添加合并逻辑
- spark和flink如何指定timeline时间来增量读取
- 三种视图无法满足业务对同一份数据的多种使用方式吧?数据摄取、查询性能和成本 不可调和,compaction异步也无法解决全部问题
- 增量视图 为什么只能是cow?
- mor模式读取时是由presto或者spark来做读时合并?合并过程中新文件如何处理?
- hudi 与 iceberg 混合查询能否可行?
日期: 5月7日, 08:10
hudi, 原名: hoodie, Hadoop Upsert Delete and Incremental
背景
Uber希望在存储上做到流批统一,需要让负责批量写入的存储系统也能支持实时写入,这就产生了update和delete的需求。为什么呢?有多种原因,例如实时计算常有的迟到数据,还有业务时效性要求以及一些合规需求(GDPR要求平台允许用户删除自己的数据)。而众所众知的是,无论是HDFS还是云平台的对象存储(例如aws的s3,阿里云的oss等),都不支持update而只能overwrite,因此要实现update和delete功能,就必须在底层存储之上做文章。Hudi于是应运而生
Upsert - hudi的招牌
如何在一个只能overwrite的文件系统上实现update操作?
hudi的思想:
把一个完整的文件拆分为多个“小文件”,当需要更新其中某条记录时,只要把包含这条记录的“小文件”给重写一遍即可
COW: Copy On Write
RECORDKEY_FIELD_OPT_KEY: 作为recordKey的字段名, 如txn_id
PARTITIONPATH_FIELD_OPT_KEY: 作为partitionPath的字段名, 如 fdate
Upsert的过程整体分为3步(这里省略了很多不太重要的步骤):
- 根据
partitionPath进行重新分区 - Tagging: 根据recordKey确定哪些记录需要插入,哪些记录需要更新。对于需要更新的记录,还需要找到旧的记录所在的文件.
tagging需要在已有的数据里寻找key相同的record,如果表的数据量比较大时会非常耗时 - 把记录写入实际的文件
MOR: Merge On Read
0.3.5版本开始引入
主要是写入性能
COW表每次在写入时,会把新写入的数据和老数据合并以后,再写成新的文件。 单单是写入的过程(不包含前期的repartition和tagging过程),就包含至少三个步骤:
- 读取老数据的parquet文件(涉及对parquet文件解码,不轻松
- 将老数据和新数据合并
- 将合并后的数据重新写成parquet文件(又涉及parquet文件编码,也不轻松)
upsert时把变更内容写入log文件,然后定期合并log文件和base文件。 这样的好处是避免了写入时读取老数据,也就避免了parquet文件不轻松的编解码过程,只需要把变更记录写入一个文件即可(而且是顺序写入)。显然是 轻松了不少
1 | warehouse |
MOR表在更新时只会把更新的那部分数据写入一个.log文件,因为.log文件不包含老数据,也不涉及tagging,又是顺序写入的,所以写入会非常快。而当客户端要读取数据时,会有两种选择:
读取时merge,同时定期地compact
- 读取时动态地把.log文件和原始数据文件(称为base文件)进行merge
- 数据保证最新,缺点是读取的性能较差
- 异步地把.log文件和base文件merge,如果merge还没完成,只能读到上个版本的数据
- 异步merge(称为compaction)有一定的延迟
Merge on read: Hudi的读取过程是实时地把base数据和log数据合并起来,并返回给用户
实时合并的实现方式是把所有log文件读入内存,放在一个HashMap里,然后遍历base文件,把base数据和缓存在内存里的log数据进行join,最后才得到合并后的结果。难免会影响到读取效率。
对于MOR表,Hudi支持3种query类型,分别是:
- Snapshot Query
- Incremental Query
- Read Optimized Query
其中: 1和3就是为了平衡读和写之间的取舍。这两者的区别是:
Snapshot Query和上文所说的一样,读取时进行“实时合并”;
Read Optimized Query则不同,只读取base文件,不读取log文件,因此读取效率和COW表相同,但读到的数据可能不是最新的。

Index
Tagging过程中,需要使用Index判断一条数据是否已经插入过。
Bloom Index:实现原理是bloom filter。优点是效率高,缺点是bloom filter固有的假阳性问题,所以Hudi对bloom filter里存在的key,还需要回溯原文件再查找一遍。Hudi默认使用的是Bloom Index。
Simple Index:实现原理是把新数据和老数据进行join。优点是实现最简单,无需额外的资源。缺点是性能比较差。
HBase Index:实现原理是把index存放在HBase里面。优点是性能最高,缺点是需要外部的系统,增加了运维压力。
global index里面存放了一张表里所有record的key,而non-global index是每个partition都有一个对应的index,里面只存放了本partition的key。如果用户使用non-global index,就必须保证同一个key的record不会出现在多个partition里面。
non-global index主要是出于index的维护成本和写入性能考虑。因为维护一个global index的难度更大,对写入性能的影响也更大。
事务: Transactional
Timeline
| 特性 | 功能 |
|---|---|
| 原子性 | 写入即使失败,也不会造成数据损坏 |
| 隔离性 | 读写分离,写入不影响读取,不会读到写入中途的数据 |
| 回滚 | 可以回滚变更,把数据恢复到旧版本 |
| 时间旅行 | 可以读取旧版本的数据(但太老的版本会被清理掉) |
| 存档 | 可以长期保存旧版本数据(存档的版本不会被自动清理) |
| 增量读取 | 可以读取任意两个版本之间的差分数据 |
Hudi在这张表的timeline里(实际存放在.hoodie目录下)会记录下v1和v2对应的文件列表。当client读取数据时,首先会查看timeline里最新的commit是哪个,从最新的commit里获得对应的文件列表,再去这些文件读取真正的数据。
多版本隔离的能力。当一个client正在读取v1的数据时,另一个client可以同时写入新的数据,新的数据会被写入新的文件里,不影响v1用到的数据文件。只有当数据全部写完以后,v2才会被commit到timeline里面。后续的client再读取时,读到的就是v2的数据。
顺带一提的是,尽管Hudi具备多版本数据管理的能力,但旧版本的数据不会无限制地保留下去。Hudi会在新的commit完成时开始清理旧的数据,默认的策略是“清理早于10个commit前的数据”。
Incremental Query(增量查询)
其实Hudi对每一条数据,都有一个隐藏字段_hoodie_commit_time用于记录commit时间,这个字段会和其他数据字段一起保存在parquet文件里
Hudi在读取parquet文件时,会同时用这个字段对结果进行过滤,把不属于时间范围内的记录都过滤掉。
不仅仅是 表格式
表格式 包括表的布局、表的Schema和对表更高的元数据跟踪
Hudi 使用 Avro 格式来存储、管理和演进表的schema
Hudi 强制执行 schema-on-write

timeline
索引时间线
索引机制
[[hudi-indexing-mechanisms]]
并发控制
集成Presto
https://prestodb.io/blog/2020/08/04/prestodb-and-hudi