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

如何并行读取多目录至多个Spark DataFrame?

并行读取多目录为不同Schema的Spark DataFrame

方法一:利用Spark RDD分布式并行读取

借助Spark的分布式特性,将目录列表转为RDD后,通过map操作把每个目录的读取任务分配到不同executor节点执行,避免driver端串行循环的低效问题。

代码示例:

from pyspark.sql import SparkSession

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

# 目标目录列表
dir_list = ['dir1', 'dir2', 'dir3']

# 定义单个目录读取函数
def read_single_dir(dir_path):
    # 根据实际数据情况调整参数,比如header=True、sep='\t'等
    return spark.read.csv(dir_path, inferSchema=True)

# 并行化目录列表,执行读取操作
dfs_rdd = spark.sparkContext.parallelize(dir_list).map(read_single_dir)

# 收集结果到driver端,得到DataFrame列表(实际数据仍存储在executor)
dfs_list = dfs_rdd.collect()

# 按需拆分到单独变量
df1, df2, df3 = dfs_list

方法二:Python线程池并行读取(本地/小规模场景)

如果是本地运行Spark或数据规模较小,可通过Python线程池在driver端并行发起读取请求,适合快速实现轻量并行。

代码示例:

from pyspark.sql import SparkSession
from concurrent.futures import ThreadPoolExecutor

spark = SparkSession.builder.appName("ThreadedDirRead").getOrCreate()

dir_list = ['dir1', 'dir2', 'dir3']

def read_single_dir(dir_path):
    return spark.read.csv(dir_path, inferSchema=True, header=True)

# 自定义线程数,建议根据目录数量和系统资源调整
with ThreadPoolExecutor(max_workers=4) as executor:
    dfs_list = list(executor.map(read_single_dir, dir_list))

df1, df2, df3 = dfs_list

关键注意事项

  • 每个目录Schema不同,必须单独读取,不能合并目录后统一读取再拆分
  • 若已知Schema,提前定义并传入read.csv(schema=预定义Schema),比自动推断Schema更高效准确
  • 大规模分布式场景优先用RDD方法,避免driver端线程池成为性能瓶颈
  • collect()仅将DataFrame的元数据拉回driver,实际数据仍存储在Spark集群的executor节点,无需担心driver内存溢出

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 11:50:22