如何仅用Python和Boto3将Spark DataFrame导出为S3指定路径CSV
解决方案
首先修正你现有代码里的两处小错误:
spark.confi.get→spark.conf.get- SQL语句里的
selec→select,表名需要用反引号或引号包裹,比如SELECT column_1, column_2 FROMmy_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
相关产品推荐
相关产品推荐

