如何获取Hudi表的最新版本记录?附AWS Glue作业代码
获取Hudi表中每个user_id的最新版本记录方案
问题根源分析
你当前的Hudi配置中,hoodie.datasource.write.recordkey.field设置为user_id,timestamp,这会将user_id+timestamp组合作为唯一记录键。因此更新user_id=2的记录时,由于timestamp不同,Hudi会将其视为新记录插入而非更新原有记录,最终表中保留三条历史记录(符合你的预期)。但要获取每个user_id的最新版本,需要调整配置并使用正确的查询方式。
第一步:修正Hudi写入配置
修改hudiWriteConfig中的记录键配置,将recordkey.field改为仅user_id,确保同一user_id的不同版本会被关联:
hudiWriteConfig = { 'className': 'org.apache.hudi', 'hoodie.table.name': hudi_table_name, 'hoodie.datasource.write.operation': 'upsert', 'hoodie.datasource.write.table.type': 'MERGE_ON_READ', 'hoodie.datasource.write.precombine.field': 'timestamp', 'hoodie.datasource.write.recordkey.field': 'user_id', # 仅保留user_id作为记录键 # 其他原有配置保持不变 }
recordkey.field:定义唯一标识记录的字段,用user_id确保同一用户的所有版本归为同一记录组precombine.field:用timestamp作为合并时的版本判断依据,Hudi会自动保留timestamp最大的版本作为最新记录
第二步:获取最新版本的查询方式
方式1:使用Hudi内置快照视图(Spark/Glue中查询)
MERGE_ON_READ类型的Hudi表同步到Glue Catalog后,会自动生成快照视图(命名规则:[表名]_snapshot),该视图默认返回每个记录的最新版本:
SELECT user_id, name, timestamp FROM your_hudi_table_snapshot;
方式2:Spark SQL窗口函数手动筛选
如果需要自定义筛选逻辑,可以利用Hudi的元数据字段或原始timestamp字段,通过窗口函数实现:
WITH ranked_records AS ( SELECT user_id, name, timestamp, ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY timestamp DESC) AS rn FROM your_hudi_table ) SELECT user_id, name, timestamp FROM ranked_records WHERE rn = 1;
方式3:Athena中查询最新版本
在Athena中,直接使用窗口函数基于timestamp或Hudi元数据字段_hoodie_commit_time(记录数据写入Hudi的时间,也可作为版本判断依据)筛选:
WITH ranked_records AS ( SELECT user_id, name, timestamp, ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY timestamp DESC) AS rn FROM your_database.your_hudi_table ) SELECT user_id, name, timestamp FROM ranked_records WHERE rn = 1;
注意事项
- 修改配置后重新运行Glue作业,后续的upsert操作会正确关联同一user_id的记录,同时Hudi依然会保留所有历史版本
- 若需查询全部历史记录,直接查询原始Hudi表即可获取所有三条记录
内容的提问来源于stack exchange,提问作者Mee
相关产品推荐
相关产品推荐

