Dataproc上Hudi Upsert生成多Parquet文件的配置问题排查
Hudi Upsert后生成多Parquet文件的解决方案
用户问题
我是Dataproc新手,核心目标是将本地数据源的表导入后以Parquet格式存储在Cloud Storage桶中,并同步更新BigQuery中的对应表。此前已通过Dataproc/PySpark/Hudi完成本地数据的导入与Cloud Storage存储工作。当前遇到的问题:配置Hudi的upsert参数后,期望将新数据追加更新至同一个Parquet文件(避免删除历史数据,每张表仅保留一个Parquet文件),但执行Upsert代码后,Cloud Storage中生成了新的Parquet文件,想确认是否遗漏了相关配置。
Upsert代码:
table_location = "gs://bucket/{}/".format(table_name) updates = spark.read.format("jdbc") \ .option("url",url) \ .option("user", username) \ .option("password", password) \ .option("driver", "com.sap.db.jdbc.Driver") \ .option("query", query) \ .load() hudi_options = { 'hoodie.table.name': table_name, 'hoodie.datasource.write.storage.type': 'COPY_ON_WRITE', 'hoodie.datasource.write.recordkey.field': 'a,b,c,d,e', 'hoodie.datasource.write.table.name': table_name, 'hoodie.datasource.write.operation': 'upsert', 'hoodie.datasource.write.precombine.field': 'x', 'hoodie.upsert.shuffle.parallelism': 2, 'hoodie.insert.shuffle.parallelism': 2, 'path': table_location, 'hoodie.datasource.hive_sync.enable': 'true', 'hoodie.datasource.hive_sync.database': database_name, 'hoodie.datasource.hive_sync.table': table_name, 'hoodie.datasource.hive_sync.partition_extractor_class': 'org.apache.hudi.hive.MultiPartKeysValueExtractor', 'hoodie.datasource.hive_sync.use_jdbc': 'false', 'hoodie.datasource.hive_sync.mode': 'hms' } updates.write.format("hudi") \ .options(**hudi_options) \ .mode("append") \ .save()
问题根源
Hudi默认机制下,即使设置低并行度,也会基于文件大小阈值、数据分区规则以及COPY_ON_WRITE模式的特性生成新文件。COPY_ON_WRITE模式的upsert会重写包含更新记录的文件分片,而非直接追加到原文件,这是导致多文件生成的核心原因。
配置调整方案
要实现单Parquet文件存储,需添加以下配置并调整写入逻辑:
1. 修改Hudi配置参数
在hudi_options中加入以下参数:
# 设置单个文件最大容量(示例为1GB,可根据数据量调整) 'hoodie.datasource.write.file.max.size': '1073741824', # 禁用分区(无分区需求时设置,避免按分区生成多文件) 'hoodie.datasource.write.partitionpath.field': '', # 保留最新版本文件,自动清理旧版本 'hoodie.cleaner.policy': 'KEEP_LATEST_COMMIT', 'hoodie.cleaner.commits.retained': '1'
2. 强制重分区为1
在写入前将数据重分区为1个分区,确保输出仅生成单个文件:
updates.repartition(1).write.format("hudi") \ .options(**hudi_options) \ .mode("append") \ .save()
注意事项
- 单文件存储仅适合小数据量场景,数据量过大时会严重影响读写性能。
- COPY_ON_WRITE模式下,upsert仍会生成新文件版本,cleaner配置会自动删除旧版本,最终保留最新的单个文件。
- 同步BigQuery时,可通过BigQuery直接挂载Cloud Storage的Parquet文件路径,或借助Hudi的Hive Sync同步元数据后,通过BigQuery连接Hive metastore访问数据。
内容的提问来源于stack exchange,提问作者Cir02
相关产品推荐
相关产品推荐

