如何并行读取多目录至多个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
相关产品推荐
相关产品推荐

