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

如何仅用Python和Boto3将Spark DataFrame导出为S3指定路径CSV

解决方案

首先修正你现有代码里的两处小错误:

  • spark.confi.get → spark.conf.get
  • SQL语句里的 selec → select,表名需要用反引号或引号包裹,比如 SELECT column_1, column_2 FROM my_table``

接下来,基于你已经用df.collect()构建CSV内容的思路,结合Boto3完成上传,完整代码如下:

import datetime as dt
from pyspark.sql.functions import *
import boto3

# 修正配置获取的拼写错误
MY_BUCKET = spark.conf.get('spark.yarn.appMasterEnv.MY_BUCKET')

# 修正SQL语句的语法错误
df = spark.sql("SELECT column_1, column_2 FROM `my_table`")

date = dt.date.today()
file_date = date.strftime("%Y_%m_%d")

# 1. 构建完整的CSV内容(包含表头和数据行)
# 获取表头
header = ','.join(df.columns)
# 遍历collect()的行,转换为逗号分隔的字符串
rows = [','.join(str(value) for value in row) for row in df.collect()]
# 拼接表头和所有行,每行用换行符分隔
csv_content = f"{header}\n" + '\n'.join(rows)

# 2. 定义S3路径的Key部分(Bucket单独传入,Key不需要包含Bucket名)
s3_key = f"folder/filename{file_date}.csv"

# 3. 初始化Boto3 S3客户端(YARN环境下通常会自动读取IAM角色凭证)
s3_client = boto3.client('s3')

# 4. 上传CSV内容到指定S3路径
s3_client.put_object(
    Bucket=MY_BUCKET,
    Key=s3_key,
    Body=csv_content.encode('utf-8')  # 将字符串转为字节流上传
)

关键说明:

  • CSV内容构建:通过df.columns获取表头,遍历collect()返回的Row对象,将每个字段转为字符串后拼接成CSV行,最后组合表头和所有行。
  • S3路径处理:put_object方法中Bucket参数单独传入存储桶名,Key参数传入存储桶内的文件路径(不需要带存储桶前缀)。
  • 编码处理:上传时需要将字符串转为UTF-8编码的字节流,避免字符编码问题。
  • 大数据量注意:df.collect()会将整个DataFrame的数据拉取到Driver节点内存中,若数据量过大可能导致内存溢出。这种情况下可先将DataFrame重分区为1(df.repartition(1)),用df.write.csv写入临时路径,再通过Boto3将生成的part-xxxx.csv文件复制到目标路径并清理临时文件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 02:15:34