You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 00:30:43