PySpark中如何遍历读取S3存储桶内的所有文件?
在PySpark中遍历读取S3存储桶的所有文件
如果你需要逐个处理S3桶中的每个文件,有两种常用方案,结合你熟悉的boto3或直接用Spark生态的API:
方案1:结合boto3列出文件后逐个读取
利用你已掌握的boto3列出S3对象,再循环用PySpark读取每个文件:
示例代码
import boto3 from pyspark.sql import SparkSession # 初始化SparkSession spark = SparkSession.builder.appName("ReadS3Files").getOrCreate() # 配置boto3 S3客户端 s3_client = boto3.client('s3') bucket_name = "your-bucket-name" # 分页列出桶内所有文件(处理大量文件的情况) files = [] response = s3_client.list_objects_v2(Bucket=bucket_name) while True: for obj in response.get('Contents', []): # 跳过S3虚拟文件夹(以/结尾) if not obj['Key'].endswith('/'): files.append(f"s3://{bucket_name}/{obj['Key']}") if not response['IsTruncated']: break response = s3_client.list_objects_v2(Bucket=bucket_name, ContinuationToken=response['NextContinuationToken']) # 逐个读取并处理文件 for file_path in files: # 根据你的文件格式调整format(比如csv、parquet) df = spark.read.format("csv").option("header", "true").load(file_path) # 这里添加你的文件处理逻辑,比如打印结构、写入数据库等 df.printSchema()
方案2:使用Spark原生API获取文件列表
不需要依赖boto3,直接用Spark或Hadoop的API获取文件路径:
方式A:用wholeTextFiles提取路径
适合文本类文件,可快速获取所有文件路径:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("ListS3Files").getOrCreate() bucket_path = "s3://your-bucket-name" # 获取所有文件的路径(键为路径,值为文件内容) file_rdd = spark.sparkContext.wholeTextFiles(f"{bucket_path}/*") file_paths = file_rdd.keys().collect() # 逐个读取处理 for path in file_paths: df = spark.read.format("parquet").load(path) # 处理逻辑...
方式B:Hadoop FileSystem API
适合更底层的文件操作,支持递归遍历:
from pyspark.sql import SparkSession from org.apache.hadoop.fs import Path from org.apache.hadoop.conf import Configuration spark = SparkSession.builder.appName("ListS3Files").getOrCreate() conf = Configuration() fs = Path.getFileSystem(conf) bucket_path = Path("s3://your-bucket-name") # 递归列出所有文件状态 file_statuses = fs.listStatus(bucket_path) file_paths = [status.getPath().toString() for status in file_statuses if not status.isDirectory()] # 逐个读取处理 for path in file_paths: df = spark.read.format("json").load(path) # 处理逻辑...
额外提示
- 如果所有文件格式统一且不需要单独处理,直接用
spark.read.format("xxx").load("s3://your-bucket-name/")一次性合并读取,效率远高于逐个读取。 - 确保Spark环境已配置S3访问权限(如IAM角色、access key配置),避免权限错误。
内容的提问来源于stack exchange,提问作者PythonDeveloper
相关产品推荐
相关产品推荐

