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
相关产品推荐
相关产品推荐

