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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:14:52