如何在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
相关产品推荐
相关产品推荐

