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

PySpark写入AWS S3文件异常:生成文件夹而非指定CSV文件

问题原因

PySpark的save()方法默认会将DataFrame以分布式文件集合的形式写入指定路径:它会创建一个目录(即你看到的price.csv文件夹),里面生成多个以part-开头的分区文件,而非直接生成单个price.csv文件。这是Spark分布式计算的特性——数据分散在多个节点的分区中,每个分区会单独写入一个文件。

解决方案

要生成单个指定名称的CSV文件,需分两步处理:

  1. 将DataFrame合并为单个分区,确保所有数据写入一个文件
  2. 重命名生成的part-文件为目标文件名

修改后的代码

from pyspark.sql import SparkSession
from pyspark import SparkConf
import os
import sys  # 原代码遗漏了sys模块导入,需补充
from dotenv import load_dotenv
from pyspark.sql.functions import *

# Load environment variables from the .env file
load_dotenv()

os.environ['PYSPARK_PYTHON'] = sys.executable
os.environ['PYSPARK_DRIVER_PYTHON'] = sys.executable


AWS_ACCESS_KEY_ID = os.getenv("AWS_ACCESS_KEY_ID")
AWS_SECRET_ACCESS_KEY = os.getenv("AWS_SECRET_ACCESS_KEY")

# My spark configuration
conf = SparkConf()
conf.set('spark.jars.packages', 'org.apache.hadoop:hadoop-aws:3.3.2')
conf.set('spark.hadoop.fs.s3a.access.key', AWS_ACCESS_KEY_ID)
conf.set('spark.hadoop.fs.s3a.secret.key', AWS_SECRET_ACCESS_KEY)

spark = SparkSession.builder.config(conf=conf).getOrCreate()

# Create a PySpark DataFrame
df = spark.createDataFrame([(1, "John Doe", 30), (2, "Jane Doe", 35), (3, "Jim Brown", 40)], ["id", "name", "age"])

# 1. 合并为单个分区,写入临时目录
temp_path = "s3a://bucket/test/store/temp_price"
df.coalesce(1).write.format("csv").option("header","true").mode("overwrite").save(temp_path)

# 2. 找到临时目录下的part文件,重命名为price.csv
from pyspark.dbutils import DBUtils
dbutils = DBUtils(spark)

# 筛选临时目录下的目标part文件
part_files = [f.path for f in dbutils.fs.ls(temp_path) if f.path.endswith(".csv") and "part-" in f.path]
if part_files:
    target_path = "s3a://bucket/test/store/price.csv"
    dbutils.fs.mv(part_files[0], target_path)
    # 删除临时目录
    dbutils.fs.rm(temp_path, recurse=True)

# Stop the Spark context and Spark session
spark.stop()

关键说明

  • coalesce(1):将DataFrame分区数缩减为1,确保仅生成一个part-文件。注意:若数据量极大,此操作会将所有数据集中到单个节点,影响性能,需谨慎使用。
  • DBUtils:用于S3文件的列表、移动、删除操作,比原生Python文件操作更适配分布式存储环境;本地环境可替换为os模块实现相同逻辑。
  • 临时目录:先写入临时目录再重命名,避免直接指定目标文件名导致生成目录的问题。

替代方案(无需重命名)

若不需要严格的price.csv文件名,可直接读取生成目录下的part-*.csv文件——Spark及多数数据工具都支持将目录下的所有分区文件视为一个完整数据集。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 14:17:32