在Databricks中将Spark DataFrame按分区保存为单个Parquet文件
解决Spark按分区保存时每个目录仅生成单个Parquet文件的问题
当你使用partitionBy保存DataFrame时,Spark会根据集群资源和数据分布生成多个part文件。要实现每个分区目录仅保留单个文件,可采用以下两种常用方案:
方法一:按分区值循环写入(适合中小数据量)
先提取所有唯一的Filename值,然后逐个过滤出对应分区的数据,通过coalesce(1)合并为单个分区后写入目标目录。这种方式无需全局shuffle,适合数据量不大的场景。
PySpark 示例:
# 获取所有唯一的Filename值 filenames = [row.Filename for row in df.select("Filename").distinct().collect()] # 循环写入每个分区 for filename in filenames: df.filter(df.Filename == filename) \ .coalesce(1) \ .write \ .mode("overwrite") \ .parquet(f"{file_out_location}/Filename={filename}")
Scala 示例:
import org.apache.spark.sql.functions.col // 获取所有唯一的Filename值 val filenames = df.select(col("Filename")).distinct().collect().map(_.getString(0)) // 循环写入每个分区 filenames.foreach { filename => df.filter(col("Filename") === filename) .coalesce(1) .write .mode("overwrite") .parquet(s"$file_out_location/Filename=$filename") }
注意:这种方式是串行处理每个分区,数据量过大时会影响效率。
方法二:重分区后保存(适合大数据量)
通过repartition("Filename")将数据按Filename重分区,确保每个Filename对应一个Spark分区。再结合partitionBy保存,此时每个分区目录会自动生成单个part文件。这种方式利用Spark并行处理能力,适合大数据量场景,但会触发全局shuffle,需注意资源占用。
PySpark 示例:
df.repartition("Filename") \ .write \ .partitionBy("Filename") \ .mode("overwrite") \ .parquet(file_out_location)
Scala 示例:
df.repartition(col("Filename")) .write .partitionBy("Filename") .mode("overwrite") .parquet(file_out_location)
说明:repartition("Filename")会将相同Filename的数据集中到同一个分区,后续保存时每个分区对应一个目录,自然生成单个文件。
可选:自定义输出文件名
Spark默认生成以part-开头的文件名,若需将文件重命名为Filename=file1.parquet这类格式,需在保存后通过文件系统工具手动修改。以下是PySpark结合本地文件系统的示例:
import os # 先按方法二保存数据 df.repartition("Filename").write.partitionBy("Filename").mode("overwrite").parquet(file_out_location) # 遍历分区目录并重命名文件 for root, dirs, _ in os.walk(file_out_location): for dir_name in dirs: if dir_name.startswith("Filename="): filename = dir_name.split("=")[1] dir_path = os.path.join(root, dir_name) # 找到分区内的parquet part文件 part_files = [f for f in os.listdir(dir_path) if f.startswith("part-") and f.endswith(".parquet")] if part_files: old_file = os.path.join(dir_path, part_files[0]) new_file = os.path.join(dir_path, f"{filename}.parquet") os.rename(old_file, new_file)
若使用HDFS等分布式文件系统,需改用对应的API(如hdfs命令行工具或hdfs3库)执行重命名操作。
内容的提问来源于stack exchange,提问作者giri rajh
相关产品推荐
相关产品推荐

