如何自定义Hudi的_hoodie_commit_time元数据列,实现自定义时间更新管控
解决Hudi自定义时间线与延迟数据过滤问题
需求说明
需要修改Hudi默认以当前时间作为摄入时间线的行为,改用业务字段last_update_time控制数据版本。要求表仅保留最新状态记录:当延迟到达的数据其last_update_time早于已有记录的最新更新时间时,不覆盖原有记录。所有记录均包含last_update_time datetime字段,用于标识更新发生的时间。
数据示例
第1天初始数据
+---+-----+---------------------------+ |id |value|last_update_time | +---+-----+---------------------------+ |100|a |2022-11-14T13:51:39.340396Z| |101|b |2022-11-14T12:14:58.597216Z| |103|c |2022-11-14T12:14:58.597216Z| +---+-----+---------------------------+
第2天接收数据
+---+-----+---------------------------+ |id |value|last_update_time | +---+-----+---------------------------+ |100|a1 |2022-11-25T13:51:39.340396Z| <- 更新(应覆盖原有记录) |101|b1 |2022-11-12T12:14:58.597216Z| <- 延迟更新(不应覆盖原有记录) |104|d1 |2022-11-25T12:14:58.597216Z| <- 新记录(应添加) +---+-----+---------------------------+
预期输出
+---+-----+---------------------------+ |id |value|last_update_time | +---+-----+---------------------------+ |100|a1 |2022-11-25T13:51:39.340396Z| |101|b |2022-11-14T12:14:58.597216Z| |103|c |2022-11-14T12:14:58.597216Z| |104|d1 |2022-11-25T12:14:58.597216Z| +---+-----+---------------------------+
当前实际输出
+---+-----+---------------------------+ |id |value|last_update_time | +---+-----+---------------------------+ |100|a1 |2022-11-25T13:51:39.340396Z| |101|b1 |2022-11-12T12:14:58.597216Z| |103|c |2022-11-14T12:14:58.597216Z| |104|d1 |2022-11-25T12:14:58.597216Z| +---+-----+---------------------------+
当前Writer配置
{ "compression": "snappy", "hoodie.cleaner.policy": "KEEP_LATEST_COMMITS", "hoodie.datasource.write.table.type": "COPY_ON_WRITE", "hoodie.datasource.write.partitionpath.field": "", "hoodie.datasource.write.recordkey.field": "id", "hoodie.datasource.write.precombine.field": "last_update_time", "hoodie.datasource.write.row.writer.enable": "true", "hoodie.cleaner.policy.failed.writes": "LAZY", "hoodie.write.lock.provider": "org.apache.hudi.client.transaction.lock.ZookeeperBasedLockProvider", "hoodie.write.concurrency.mode": "optimistic_concurrency_control", "hoodie.fail.on.timeline.archiving": "false", "hoodie.write.lock.zookeeper.url": "zookeeper", "hoodie.write.lock.zookeeper.port": "2181", "hoodie.write.lock.num_retries": "2", "hoodie.write.lock.zookeeper.lock_key": "test", "hoodie.write.lock.zookeeper.base_path": "/lock/hudi/test" }
解决方案
关键配置修改
在现有配置中添加以下参数,让Hudi用自定义的last_update_time替代默认的提交时间作为时间线排序依据:
"hoodie.datasource.write.timestamp.ordering.field": "last_update_time"
配置说明
hoodie.datasource.write.precombine.field: 你已设置为last_update_time,该参数用于预合并阶段筛选同一recordkey下的最新记录,保留last_update_time更大的条目。hoodie.datasource.write.timestamp.ordering.field: 新增该参数后,Hudi会以last_update_time判断记录的时间先后,而非默认的提交时间。结合预合并配置,可确保延迟的旧数据(last_update_time更小)不会覆盖已有新记录。
验证
修改配置后重新运行数据写入,即可得到预期结果:id=101的旧延迟数据不会覆盖原有记录,仅更新last_update_time更大的id=100记录,同时新增id=104的记录。
内容的提问来源于stack exchange,提问作者Aleksejs R
相关产品推荐
相关产品推荐

