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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 21:45:29