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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 23:11:11