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

Glue作业中PySpark DataFrame保存至S3时自定义CSV文件名方法

PySpark保存CSV到S3时自定义文件名问题解决

问题场景

使用以下代码保存PySpark DataFrame到S3时:

df.coalesce(1).write\
        .format("csv")\
        .mode("append")\
        .save(f"s3://{bucket_output}/{dirname}/{filename}", header=True, nullValue = '\u0000', emptyValue = '\u0000')

实际生成的是名为filename的目录,目录内是part-(some_numbers).csv格式的文件,而非直接生成指定文件名的CSV。

解决方案

方法一:调整写入逻辑,直接生成目标文件名

Spark的save方法默认将传入路径视为目录,因此需要先写入临时目录,再通过S3操作重命名文件。Glue作业中可直接用boto3实现:

import boto3

# 初始化S3客户端(Glue作业默认已配置权限)
s3 = boto3.client('s3')

# 定义路径参数
bucket = bucket_output
temp_dir = f"{dirname}/temp_{filename}"
target_file = f"{dirname}/{filename}.csv"

# 1. 将数据写入临时目录
df.coalesce(1).write\
    .format("csv")\
    .mode("append")\
    .save(f"s3://{bucket}/{temp_dir}", header=True, nullValue='\u0000', emptyValue='\u0000')

# 2. 查找临时目录下的CSV文件
response = s3.list_objects_v2(Bucket=bucket, Prefix=temp_dir)
part_files = [obj['Key'] for obj in response['Contents'] if obj['Key'].endswith('.csv')]

if part_files:
    # 3. 复制part文件到目标路径
    s3.copy_object(
        Bucket=bucket,
        CopySource={'Bucket': bucket, 'Key': part_files[0]},
        Key=target_file
    )
    # 4. 删除临时目录下的所有文件(包括_SUCCESS等辅助文件)
    for obj in response['Contents']:
        s3.delete_object(Bucket=bucket, Key=obj['Key'])

注意:coalesce(1)会将数据合并为单个分区,若数据量过大可能引发性能问题,需根据实际情况调整。

方法二:通过S3移动操作修正现有文件

如果已经生成错误的目录结构,可直接移动文件修正:

import boto3

s3 = boto3.client('s3')
bucket = bucket_output
source_dir = f"{dirname}/{filename}"
target_file = f"{dirname}/{filename}.csv"

# 查找目录下的CSV文件
response = s3.list_objects_v2(Bucket=bucket, Prefix=source_dir)
part_files = [obj['Key'] for obj in response['Contents'] if obj['Key'].endswith('.csv')]

if part_files:
    # 复制文件到目标路径
    s3.copy_object(
        Bucket=bucket,
        CopySource={'Bucket': bucket, 'Key': part_files[0]},
        Key=target_file
    )
    # 删除原目录下的所有文件(S3虚拟目录会随文件删除自动消失)
    for obj in response['Contents']:
        s3.delete_object(Bucket=bucket, Key=obj['Key'])

原因说明

Spark的DataFrameWriter.save()方法接收的是目录路径,而非文件路径。传入s3://xxx/filename时,Spark会创建该目录,然后在目录下生成分区文件、_SUCCESS等辅助文件,这是分布式写入的默认行为。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 01:18:24