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年数据留存,耗时将进一步失控。
复现步骤
- 创建示例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)
- 使用如下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", }
- 写入日期范围数据:
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)
- 分析每次加载耗时,可观察到耗时逐步增长,若按此趋势,大型表的耗时将完全失控。
预期行为
- 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
相关产品推荐
相关产品推荐

