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

Hudi长期数据留存的性能与数据完整性问题求助

Hudi全量加载场景下的数据完整性与性能优化问题

问题背景

项目要求每日执行全量加载,并保留版本供后续查询,采用Hudi维护6年数据,初始配置如下:

"hoodie.cleaner.policy": "KEEP_LATEST_BY_HOURS",
"hoodie.cleaner.hours.retained": "52560", # 24小时*365天*6年

运行约30次后,出现数据完整性受损问题:读取时数据版本混乱,产生重复记录,给依赖这些表的S3数据湖带来严重影响。

为解决该问题,调整了提交记录的最大/最小保留数,配置如下:

"hoodie.keep.max.commits": "2300", # (365天*6年)+增量值
"hoodie.keep.min.commits": "2200", # (365天*6年)+增量值2

但该方案存在长期成本过高的问题:按日分区模拟运行脚本,小型表1年内平均运行时间从25秒增至2分30秒,随着6年数据留存,耗时将进一步失控。

复现步骤

  1. 创建示例DataFrame:
from pyspark.sql import Row
from pyspark.sql.functions import lit, to_date
from datetime import datetime, timedelta

data = [
    Row(SK=-6698625589789238999, DSC='A', COD=1), 
    Row(SK=8420071140774656230,  DSC='B', COD=2), 
    Row(SK=-8344648708406692296, DSC='C', COD=4), 
    Row(SK=504019808641096632,   DSC='D', COD=5), 
    Row(SK=-233500712460350175,  DSC='E', COD=6), 
    Row(SK=2786828215451145335,  DSC='F', COD=7), 
    Row(SK=-8285521376477742517, DSC='G', COD=8), 
    Row(SK=-2852032610340310743, DSC='H', COD=9), 
    Row(SK=-188596373586653926,  DSC='I', COD=10), 
    Row(SK=890099540967675307,   DSC='J', COD=11), 
    Row(SK=72738756111436295,    DSC='K', COD=12), 
    Row(SK=6122947679528380961,  DSC='L', COD=13), 
    Row(SK=-3715488255824917081, DSC='M', COD=14), 
    Row(SK=7553013721279796958,  DSC='N', COD=15)
]
dataframe = spark.createDataFrame(data)
  1. 使用如下Hudi配置:
hudi_options = {
    "hoodie.table.name": "example_hudi",
    "hoodie.datasource.write.recordkey.field": "SK",
    "hoodie.datasource.write.table.name": "example_hudi",
    "hoodie.datasource.write.operation": "insert_overwrite_table",
    "hoodie.datasource.write.partitionpath.field": "LOAD_DATE",
    "hoodie.datasource.hive_sync.database": "default",
    "hoodie.datasource.hive_sync.table": "example_hudi",
    "hoodie.datasource.hive_sync.partition_fields": "LOAD_DATE",
    "hoodie.cleaner.policy": "KEEP_LATEST_BY_HOURS",
    "hoodie.cleaner.hours.retained": "52560",
    "hoodie.keep.max.commits": "2300", 
    "hoodie.keep.min.commits":"2200",  
    "hoodie.datasource.write.precombine.field":"",
    "hoodie.datasource.hive_sync.partition_extractor_class":"org.apache.hudi.hive.MultiPartKeysValueExtractor",
    "hoodie.datasource.hive_sync.enable":"true",
    "hoodie.datasource.hive_sync.use_jdbc":"false",
    "hoodie.datasource.hive_sync.mode":"hms",
}
  1. 写入日期范围数据:
date = datetime.strptime('2023-06-02', '%Y-%m-%d')  # 初始日期(yyyy-mm-dd)
final_date = datetime.strptime('2023-11-01', '%Y-%m-%d')  # 结束日期(yyyy-mm-dd)
while date <= final_date:
    dataframe = dataframe.withColumn("LOAD_DATE", to_date(lit(date.strftime('%Y-%m-%d'))))
    dataframe.write.format("hudi"). \
        options(**hudi_options). \
        mode("append"). \
        save(basePath)
    date += timedelta(days=1)
  1. 分析每次加载耗时,可观察到耗时逐步增长,若按此趋势,大型表的耗时将完全失控。

预期行为

  • 30次提交后无重复文件产生;
  • 执行时间不会随时间显著增加;
  • 元数据遵循hoodie.cleaner.policy KEEP_LATEST_BY_HOURS配置的行为。

环境信息

  • Hudi版本:0.12.2
  • Spark版本:3.3.1
  • Hive版本:3.1.3
  • 存储:S3(EMRFS)
  • 平台:AWS EMR

现寻求可行的解决方案,以兼顾数据完整性与长期运行性能。


内容的提问来源于stack exchange,提问作者Luiz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 11:42:54