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

如何使用Pyspark读取S3存储桶指定文件夹内最新的CSV文件

实现方案

整体逻辑:先通过AWS S3接口拉取目标文件夹下所有CSV文件的元数据,按最后修改时间排序取最新文件,再调用PySpark接口读取该文件。

前置依赖

  • 环境已安装配置boto3库用于操作S3元数据
  • PySpark环境已配置S3访问权限(EMR集群默认自带配置,本地环境需提前配置AWS凭证)

完整代码示例

import boto3
from pyspark.sql import SparkSession

# 初始化SparkSession
spark = SparkSession.builder.appName("ReadLatestS3CSV").getOrCreate()

# 初始化S3客户端
s3 = boto3.client('s3')

# 配置目标存储桶和前缀参数
bucket_name = "data"
prefix = "first/"

# 拉取目标路径下所有对象
response = s3.list_objects_v2(Bucket=bucket_name, Prefix=prefix)

# 筛选CSV文件,排除文件夹本身
csv_files = [
    obj for obj in response["Contents"]
    if obj["Key"].endswith(".csv") and obj["Key"] != prefix
]

# 按最后修改时间倒序排序,取最新的第一个文件
latest_file = sorted(csv_files, key=lambda x: x["LastModified"], reverse=True)[0]
latest_file_s3_path = f"s3://{bucket_name}/{latest_file['Key']}"

# 读取为PySpark DataFrame,可根据实际CSV格式调整读取参数
df = spark.read.csv(
    latest_file_s3_path,
    header=True,  # 若CSV无表头则设置为False
    inferSchema=True
)

# 验证读取结果
df.show(5)

注意事项

  • 若目标文件夹下文件数量超过1000,需要补充list_objects_v2的分页逻辑,调用ContinuationToken参数遍历所有对象
  • 读取CSV的参数可根据实际文件调整,常用自定义参数包括分隔符sep、编码encoding、时间格式timestampFormat等
  • 出现S3访问权限报错时,可在初始化boto3客户端时传入aws_access_key_id和aws_secret_access_key参数,不建议硬编码凭证到代码中

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 21:18:03