在Databricks中用PySpark获取S3文件列表 替代boto提升百万级文件查询速度
PySpark快速枚举S3百万级文件方案
完全可以通过PySpark实现该需求,速度远快于单节点boto调用,核心是Spark利用分布式能力并行调用S3的批量列表接口,避免了单进程串行分页拉取的性能瓶颈。
前置依赖
- Spark集群已配置S3访问权限(通过IAM角色、核心站点配置AK/SK均可,无需在代码中硬编码密钥)
- 集群Hadoop版本已适配S3A驱动(主流Spark发行版均默认支持)
代码实现
完全复现原有boto函数的功能,返回结构和原函数一致,同时支持直接转为Spark DataFrame用于后续处理:
from pyspark.sql import SparkSession from datetime import datetime import os def get_s3_files_spark(source_directory, file_type="json"): # 初始化SparkSession,可根据集群情况自定义配置 spark = SparkSession.builder.appName("ListS3Files").getOrCreate() sc = spark.sparkContext # 路径解析逻辑和原有代码完全对齐 source_dir_str = str(source_directory) path_parts = source_directory.parts file_prepend_path = f"/{'/'.join(path_parts[1:4])}" bucket_name = path_parts[3] prefix = "/".join(path_parts[4:]) s3_full_prefix = f"s3a://{bucket_name}/{prefix}" # 调用Hadoop FileSystem接口访问S3 hadoop_conf = sc._jsc.hadoopConfiguration() path = sc._jvm.org.apache.hadoop.fs.Path(s3_full_prefix) fs = path.getFileSystem(hadoop_conf) # 递归枚举所有文件,第二个参数True代表遍历所有子目录 file_status_iter = fs.listFiles(path, True) s3_source_files = [] current_time = str(datetime.now()) target_suffix = f".{file_type}" while file_status_iter.hasNext(): status = file_status_iter.next() # 过滤指定后缀的文件 if not status.getPath().getName().endswith(target_suffix): continue # 构造和原函数格式一致的全路径 file_key = status.getPath().toString().replace(f"s3a://{bucket_name}/", "") full_file_path = f"{file_prepend_path}/{file_key}" s3_source_files.append( ( full_file_path, status.getLen(), source_dir_str, current_time ) ) # 可选:直接返回Spark DataFrame,无需拉回Driver端即可分布式处理 # df = spark.createDataFrame(s3_source_files, schema=["file_path", "file_size", "source_dir", "list_time"]) # return df return s3_source_files
性能优化建议
- 千万级以上文件场景可将前缀拆分为多个子前缀,通过
sc.parallelize分发到多个Executor并行枚举,性能可再提升数倍 - 避免将全量文件列表拉回Driver端,直接构造Spark DataFrame在集群内分布式处理,可避免Driver内存溢出
- 可在Hadoop配置中添加
fs.s3a.list.version=2,启用S3批量拉取接口,单次请求最多拉取1000条结果,降低请求开销
内容的提问来源于stack exchange,提问作者CodingInCircles
相关产品推荐
相关产品推荐

