如何使用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
相关产品推荐
相关产品推荐

