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

PySpark DataFrame写入S3指定目录结构问题求助

解决PySpark写入CSV到指定S3目录结构的问题

你的代码存在几个关键问题,以下是针对性的修正方案:

问题分析

  1. DataFrame不可变性忽略:df.withColumn("datetime", current_timestamp())不会修改原DataFrame,必须将结果重新赋值给df才会生效。
  2. 时间格式不匹配:current_timestamp()生成的是Timestamp类型,partitionBy会默认生成带时区的目录名(如datetime=2023-05-15T03:12:23.123+00:00),不符合你需要的20230515031223格式。
  3. 写入路径逻辑错误:Spark写入CSV时,指定的路径是目录而非文件名,直接把input_file_actual.csv加到路径里会创建同名目录,而非目标文件。

修正步骤

1. 导入必要函数

先导入格式化时间的工具函数:

from pyspark.sql.functions import current_timestamp, date_format

2. 添加格式化后的datetime列

将当前时间转为yyyyMMddHHmmss格式的字符串,作为分区列:

# 生成符合要求的datetime字符串列
df = df.withColumn("datetime", date_format(current_timestamp(), "yyyyMMddHHmmss"))

3. 写入基础目录(自动生成分区结构)

指定基础路径为s3a://{bucket}/data/data,Spark会自动根据datetime列的值生成datetime=20230515031223目录:

bucket = "my-internal-bucket"
base_path = f"s3a://{bucket}/data/data"

# 写入CSV,开启表头,按datetime分区
df.write.option("header", "true").partitionBy("datetime").mode("overwrite").csv(base_path)

4. 重命名生成的part文件为指定文件名

Spark写入后会在分区目录下生成类似part-00000-xxx.csv的文件,需要将其重命名为input_file_actual.csv,可以用Hadoop文件系统API实现:

from pyspark.sql import SparkSession

# 获取Hadoop文件系统实例
spark = SparkSession.builder.getOrCreate()
fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration())

# 获取生成的datetime值,定位分区目录
datetime_value = df.select("datetime").first()[0]
partition_dir = f"{base_path}/datetime={datetime_value}"
target_file = f"{partition_dir}/input_file_actual.csv"

# 筛选分区目录下的CSV文件(排除_SUCCESS标记文件)
file_status = fs.listStatus(spark._jvm.org.apache.hadoop.fs.Path(partition_dir))
csv_files = [status.getPath().toString() for status in file_status if status.getPath().getName().endswith(".csv")]

# 重命名第一个CSV文件为目标文件名,删除其他冗余文件
if csv_files:
    fs.rename(spark._jvm.org.apache.hadoop.fs.Path(csv_files[0]), spark._jvm.org.apache.hadoop.fs.Path(target_file))
    for f in csv_files[1:]:
        fs.delete(spark._jvm.org.apache.hadoop.fs.Path(f), False)

可选:合并为单个文件再写入(小数据量适用)

如果数据量不大,可先合并为一个分区,避免生成多个part文件:

# 合并为1个分区写入临时目录
temp_path = f"{base_path}/temp_{datetime_value}"
df.coalesce(1).write.option("header", "true").mode("overwrite").csv(temp_path)

# 移动临时文件到目标路径并重命名
temp_csv_files = [status.getPath().toString() for status in fs.listStatus(spark._jvm.org.apache.hadoop.fs.Path(temp_path)) if status.getPath().getName().endswith(".csv")]
if temp_csv_files:
    fs.rename(spark._jvm.org.apache.hadoop.fs.Path(temp_csv_files[0]), spark._jvm.org.apache.hadoop.fs.Path(target_file))
    # 删除临时目录
    fs.delete(spark._jvm.org.apache.hadoop.fs.Path(temp_path), True)

完整代码示例

from pyspark.sql import SparkSession
from pyspark.sql.functions import current_timestamp, date_format

# 初始化SparkSession(未初始化时执行)
spark = SparkSession.builder.appName("WriteCSVToS3").getOrCreate()

# 读取原DataFrame(你的原有代码)
bucket = "my-internal-bucket"
directory_remote = "data"
toName_csv = "your_input_file.csv"
schema = your_schema_definition  # 替换为你的实际Schema
df = spark.read.format('csv').options(header='true').load(f"s3a://{bucket}/{directory_remote}/{toName_csv}", schema=schema)

# 添加格式化后的datetime列
df = df.withColumn("datetime", date_format(current_timestamp(), "yyyyMMddHHmmss"))

# 定义路径变量
base_path = f"s3a://{bucket}/data/data"
datetime_value = df.select("datetime").first()[0]
target_file_path = f"{base_path}/datetime={datetime_value}/input_file_actual.csv"

# 合并分区写入临时目录
temp_path = f"{base_path}/temp_{datetime_value}"
df.coalesce(1).write.option("header", "true").mode("overwrite").csv(temp_path)

# 处理文件重命名与清理
fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration())
temp_csv_files = [status.getPath().toString() for status in fs.listStatus(spark._jvm.org.apache.hadoop.fs.Path(temp_path)) if status.getPath().getName().endswith(".csv")]

if temp_csv_files:
    fs.rename(spark._jvm.org.apache.hadoop.fs.Path(temp_csv_files[0]), spark._jvm.org.apache.hadoop.fs.Path(target_file_path))
    fs.delete(spark._jvm.org.apache.hadoop.fs.Path(temp_path), True)

print(f"文件已成功写入:{target_file_path}")

注意:coalesce(1)会将所有数据集中到单个Executor,不适合超大数据量场景。如果数据量较大,建议保留多part文件,或使用合理的分区策略拆分数据,避免单文件过大。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 15:12:41