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

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 17:54:03