Hudi中历史数据覆盖最新数据问题及解决方法咨询
Hudi迟来历史数据覆盖最新数据的问题处理
问题现象
先批量写入初始数据到Hudi表,后续每日写入增量数据,但当迟来的历史数据(dms_timestamp更早)到达时,表中已有的最新数据被历史数据覆盖,不符合预期。
首次写入数据
+---+-----+-------------+ | id| req|dms_timestamp| +---+-----+-------------+ | 1| one| 2022-12-17| | 2| two| 2022-12-17| | 3|three| 2022-12-17| +---+-----+-------------+
首次写入配置
"className"-> "org.apache.hudi", "hoodie.datasource.write.precombine.field"-> "dms_timestamp", "hoodie.datasource.write.recordkey.field"-> "id", "hoodie.table.name"-> "hudi_test", "hoodie.consistency.check.enabled"-> "false", "hoodie.datasource.write.reconcile.schema"-> "true", "path"-> basePath, "hoodie.datasource.write.keygenerator.class"-> "org.apache.hudi.keygen.ComplexKeyGenerator", "hoodie.datasource.write.partitionpath.field"-> "", "hoodie.datasource.write.hive_style_partitioning"-> "true", "hoodie.upsert.shuffle.parallelism"-> "1", "hoodie.datasource.write.operation"-> "upsert", "hoodie.cleaner.policy"-> "KEEP_LATEST_COMMITS", "hoodie.cleaner.commits.retained"-> "5",
首次写入后表数据
+-------------------+---------------------+------------------+----------------------+------------------------------------------------------------------------+---+-----+-------------+ |_hoodie_commit_time|_hoodie_commit_seqno |_hoodie_record_key|_hoodie_partition_path|_hoodie_file_name |id |req |dms_timestamp| +-------------------+---------------------+------------------+----------------------+------------------------------------------------------------------------+---+-----+-------------+ |20221214130513893 |20221214130513893_0_0|id:3 | |005674e6-a581-419a-b8c7-b2282986bc52-0_0-36-34_20221214130513893.parquet|3 |three|2022-12-17 | |20221214130513893 |20221214130513893_0_1|id:1 | |005674e6-a581-419a-b8c7-b2282986bc52-0_0-36-34_20221214130513893.parquet|1 |one |2022-12-17 | |20221214130513893 |20221214130513893_0_2|id:2 | |005674e6-a581-419a-b8c7-b2282986bc52-0_0-36-34_20221214130513893.parquet|2 |two |2022-12-17 | +-------------------+---------------------+------------------+----------------------+------------------------------------------------------------------------+---+-----+-------------+
写入迟来的历史数据
+---+----+-------------+ | id| req|dms_timestamp| +---+----+-------------+ | 1|null| 2019-01-01| +---+----+-------------+
此次写入配置
"hoodie.table.name"-> "hudi_test", "hoodie.datasource.write.recordkey.field" -> "id", "hoodie.datasource.write.precombine.field" -> "dms_timestamp", // get_common_config "className"-> "org.apache.hudi", "hoodie.datasource.hive_sync.use_jdbc"-> "false", "hoodie.consistency.check.enabled"-> "false", "hoodie.datasource.write.reconcile.schema"-> "true", "path"-> basePath, // get_partitionDataConfig -- no partitionfield "hoodie.datasource.write.keygenerator.class"-> "org.apache.hudi.keygen.ComplexKeyGenerator", "hoodie.datasource.write.partitionpath.field"-> "", "hoodie.datasource.write.hive_style_partitioning"-> "true", // get_incrementalWriteConfig "hoodie.upsert.shuffle.parallelism"-> "1", "hoodie.datasource.write.operation"-> "upsert", "hoodie.cleaner.policy"-> "KEEP_LATEST_COMMITS", "hoodie.cleaner.commits.retained"-> "5",
写入后表数据(异常结果)
+-------------------+---------------------+------------------+----------------------+------------------------------------------------------------------------+---+-----+-------------+ |_hoodie_commit_time|_hoodie_commit_seqno |_hoodie_record_key|_hoodie_partition_path|_hoodie_file_name |id |req |dms_timestamp| +-------------------+---------------------+------------------+----------------------+------------------------------------------------------------------------+---+-----+-------------+ |20221214131440563 |20221214131440563_0_0|id:3 | |37dee403-6077-4a01-bf28-7afd65ef390a-0_0-18-21_20221214131555500.parquet|3 |three|2022-12-17 | |20221214131555500 |20221214131555500_0_1|id:1 | |37dee403-6077-4a01-bf28-7afd65ef390a-0_0-18-21_20221214131555500.parquet|1 |null |2019-01-01 | |20221214131440563 |20221214131440563_0_2|id:2 | |37dee403-6077-4a01-bf28-7afd65ef390a-0_0-18-21_20221214131555500.parquet|2 |two |2022-12-17 | +-------------------+---------------------+------------------+----------------------+------------------------------------------------------------------------+---+-----+-------------+
解决方案
核心原因是:Hudi默认upsert操作的预合并逻辑仅作用于当前写入批次内的记录,不会自动与表中已存的最新记录的precombine字段对比。当写入单条历史记录时,批次内无同key记录,Hudi直接将其作为最新版本写入,覆盖原有数据。
1. 启用乐观并发控制(OCC)
添加以下配置,让Hudi在upsert时基于precombine字段判断是否更新,仅当新记录时间戳晚于已有记录时才执行更新:
"hoodie.write.concurrency.mode"-> "optimistic_concurrency_control", "hoodie.write.lock.provider"-> "org.apache.hudi.client.transaction.lock.InProcessLockProvider", // 单机测试用,生产推荐ZooKeeper/HBase锁 "hoodie.write.conflict.resolution.strategy"-> "version_strategy"
2. 自定义预合并逻辑
在写入Hudi前,手动关联表中已有最新数据,保留时间戳更大的记录:
import org.apache.hudi.DataSourceReadOptions._ import org.apache.spark.sql.functions._ // 读取Hudi表最新快照数据 val existingLatestDF = spark.read .format("org.apache.hudi") .option(QUERY_TYPE.key, QUERY_TYPE_SNAPSHOT_OPT_VAL) .load(basePath) .select("id", "req", "dms_timestamp") // 本次待写入的历史数据 val incomingDF = spark.createDataFrame(Seq( (1, null, "2019-01-01") )).toDF("id", "req", "dms_timestamp") // 合并数据,保留时间戳最新的记录 val combinedDF = incomingDF.join(existingLatestDF, Seq("id"), "full_outer") .withColumn("latest_dms", greatest( coalesce(incomingDF("dms_timestamp"), lit("1970-01-01")), coalesce(existingLatestDF("dms_timestamp"), lit("1970-01-01")) )) .withColumn("final_req", when(col("latest_dms") === incomingDF("dms_timestamp"), incomingDF("req")).otherwise(existingLatestDF("req"))) .select("id", "final_req", "latest_dms") .toDF("id", "req", "dms_timestamp") // 将合并后的数据写入Hudi combinedDF.write.format("org.apache.hudi") .options(yourUpsertConfig) .mode("append") .save(basePath)
3. 改用insert操作(不推荐)
如果业务允许同key存在多版本记录,可将写入操作改为insert,但后续查询需指定时间范围获取最新数据,不利于日常使用。
内容的提问来源于stack exchange,提问作者awadhesh14
相关产品推荐
相关产品推荐

