如何使用PySpark合并存储在不同文件夹中的文件
PySpark合并多文件夹下的Part文件解决方案
嘿,刚好处理过类似的场景,给你几个清晰的方案来搞定这个需求:
核心思路:批量读取多文件夹文件
PySpark的textFile方法支持逗号分隔的多路径或者通配符匹配,这两种方式都能轻松读取不同文件夹下的所有part文件。
方案1:明确指定所有文件夹路径
如果你的目标文件夹没有统一命名规律,直接把每个文件夹的路径(加上*匹配下的所有文件)用逗号分隔传入:
from pyspark import SparkContext, SparkConf # 初始化Spark上下文 conf = SparkConf().setAppName("MergeMultiFolderFiles") sc = SparkContext(conf=conf) # 用逗号分隔所有目标文件夹的文件路径,*表示读取该文件夹下所有文件 input_paths = "/user/home/m_f012345/*,/user/home/m_f00120/*,/user/home/m_f123120/*" # 读取所有文件到RDD raw_rdd = sc.textFile(input_paths) # 合并为单个文件并保存(coalesce(1)不会触发shuffle,比repartition(1)更高效) raw_rdd.coalesce(1).saveAsTextFile("/user/home/merged_output") # 关闭Spark上下文 sc.stop()
方案2:用通配符批量匹配文件夹
如果你的文件夹命名有规律(比如都是m_f开头),直接用通配符*匹配所有符合条件的文件夹,再匹配里面的part文件,代码更简洁:
# 替换input_paths为通配符匹配 input_paths = "/user/home/m_f*/*" # 后续读取、合并逻辑和方案1一致
如果只想读取每个文件夹下的part开头文件,可以进一步精确匹配:
input_paths = "/user/home/m_f*/part*"
方案3:用Spark DataFrame简化操作(推荐)
如果你用的是较新的Spark版本,用SparkSession的DataFrame API会更简洁直观:
from pyspark.sql import SparkSession # 初始化SparkSession spark = SparkSession.builder.appName("MergeFilesWithDF").getOrCreate() # 读取所有目标文件到DataFrame df = spark.read.text("/user/home/m_f*/*") # 合并为单个文件并保存,mode("overwrite")表示覆盖已有输出路径 df.coalesce(1).write.mode("overwrite").text("/user/home/merged_output_df") # 关闭SparkSession spark.stop()
注意事项
- 路径中的逗号不能加空格,否则会被当成路径的一部分,导致Spark找不到文件
coalesce(1)适合数据量不是特别大的场景,它不会触发数据shuffle,效率更高;如果数据量极大,合并成单个文件可能引发磁盘IO瓶颈,此时可以调整分区数(比如coalesce(5)生成5个文件)- Spark默认会跳过文件夹下的隐藏文件,但如果有
_SUCCESS这类自动生成的文件,textFile也会读取它,如果你不需要这类文件,可以用filter过滤:raw_rdd = sc.textFile(input_paths).filter(lambda line: not line.startswith("_SUCCESS"))
内容的提问来源于stack exchange,提问作者LUZO
相关产品推荐
相关产品推荐

