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

如何在Azure Databricks中生成1GB大小的Parquet分区文件

解决ADLS Gen2中Parquet分区小文件问题(每个分区生成1GB文件)

你之前尝试的spark.sql.files.maxPartitionBytes是读取阶段的参数,用于控制Spark读取文件时的分区大小,对输出文件不生效;parquet.block.size是Parquet存储的块大小(文件内部的存储单元),一个文件可包含多个块,因此也无法直接控制单文件的整体大小。以下是几种可行的解决方案:

方法1:用spark.sql.files.maxRecordsPerFile控制单文件记录数

这个参数直接限制每个输出文件的最大记录数,结合你数据的平均单条记录大小,估算出1GB对应的记录数后配置即可。比如假设单条记录平均大小为1KB,1GB约等于1048576条记录:

# 根据你的数据调整记录数,确保文件接近1GB
spark.conf.set("spark.sql.files.maxRecordsPerFile", 1048576)

df.write.partitionBy("date").mode("overwrite").parquet(target_file_path)

优势:无需提前统计数据,配置简单,写入时直接生效。注意需根据实际数据的平均记录大小调整数值。

方法2:动态计算分区数,精准控制文件大小

先统计每个date分区的数据量,再动态计算每个分区需要的输出文件数,通过临时分区键实现精准拆分:

from pyspark.sql.functions import sum, length, lit, floor, rand, to_json, struct
from pyspark.sql.types import LongType

# 估算每个date分区的总字节数(用JSON序列化长度近似,可根据实际字段调整)
size_stats = df.select(
    "date",
    sum(length(to_json(struct(df.columns)))).cast(LongType()).alias("partition_size")
).groupBy("date").agg(sum("partition_size").alias("total_bytes")).collect()

# 构建日期-大小映射
date_size_map = {row["date"]: row["total_bytes"] for row in size_stats}

# 目标文件大小(1GB)
TARGET_FILE_SIZE = 1024 * 1024 * 1024

# 计算每个日期需要的分区数
date_part_count = {
    date: max(1, (size // TARGET_FILE_SIZE) + 1) 
    for date, size in date_size_map.items()
}

# 添加临时分区键,拆分每个date分区为指定数量的子分区
df_with_temp_part = df.withColumn(
    "temp_part",
    floor(rand() * lit(date_part_count[df["date"]]))
)

# 按date和temp_part重分区后写入,temp_part会被自动忽略(仅按date分区)
df_with_temp_part.repartition("date", "temp_part")\
    .write.partitionBy("date")\
    .mode("overwrite")\
    .parquet(target_file_path)

优势:精准控制每个分区的文件数量,避免全局repartition的高额开销;缺点是需要提前统计数据,增加少量前置计算时间。

方法3:写入后合并小文件(针对已生成的小文件)

如果已经生成大量小文件,可读取每个分区数据,按目标大小合并后重新写入:

from pyspark.sql.functions import input_file_name, split

# 读取现有Parquet数据,提取date分区信息
existing_df = spark.read.parquet(target_file_path)
existing_df = existing_df.withColumn("date_part", split(input_file_name(), "/").getItem(-2))

# 目标文件大小(1GB)
TARGET_FILE_SIZE = 1024 * 1024 * 1024
# 替换为你的数据平均单条记录大小(单位:字节)
average_record_size = 1024

# 逐个处理每个date分区
for date_part in existing_df.select("date_part").distinct().rdd.flatMap(lambda x: x).collect():
    partition_df = existing_df.filter(existing_df.date_part == date_part)
    # 计算需要合并的分区数
    approx_total_size = partition_df.count() * average_record_size
    num_coalesce = max(1, int(approx_total_size // TARGET_FILE_SIZE) + 1)
    # 合并后覆盖原分区
    partition_df.coalesce(num_coalesce)\
        .write.mode("overwrite")\
        .parquet(f"{target_file_path}/{date_part}")

优势:无需修改原写入逻辑,修复已有小文件;缺点是需要二次读写数据,适合离线场景。

方法4:启用Spark自适应查询执行(AQE)

开启AQE后,Spark会自动合并小的shuffle分区,减少输出文件数量:

# 启用自适应查询执行
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
# 设置合并后的最小分区大小(接近1GB)
spark.conf.set("spark.sql.adaptive.coalescePartitions.minPartitionSize", "1g")

df.write.partitionBy("date").mode("overwrite").parquet(target_file_path)

优势:无需手动调整分区,Spark自动优化;注意AQE主要针对shuffle后的分区,若写入前无shuffle操作,效果可能有限。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 01:37:11