如何高效将数百万单个JSON文件导入Spark DataFrame?
最优解决方案:直接批量读取JSON文件
你的循环union方法性能差的核心原因是:每次union都会累加DataFrame的依赖链(lineage),当文件数量达到数千甚至数百万时,执行计划会变得异常复杂,调度与优化的开销呈指数级增长,同时单个小文件的逐次读取无法利用Spark的分布式并行处理能力。
Spark原生支持批量读取多文件,无需手动循环合并,下面是几种高效方案:
1. 直接读取包含所有JSON文件的目录
如果所有JSON文件都存放在同一个目录下(或统一的父目录),直接传入目录路径即可,Spark会自动遍历目录内的所有文件(默认不递归子目录,如需递归可添加recursiveFileLookup=True参数):
df = spark.read.json( path="/path/to/all_json_files", schema=my_schema, recursiveFileLookup=True # 可选:如果需要读取子目录中的文件 )
2. 使用通配符匹配目标文件
如果需要筛选特定命名规则的文件,用通配符批量匹配:
- 匹配目录下所有JSON文件:
/path/to/*.json - 递归匹配所有子目录中的JSON文件:
/path/to/**/*.json
df = spark.read.json( path="/path/to/**/*.json", schema=my_schema )
3. 传入文件路径列表(适用于需自定义筛选逻辑的场景)
如果需要通过代码筛选特定文件(比如按文件名、修改时间过滤),可以先获取符合条件的文件路径列表,再直接传给read.json:
import glob # 获取所有目标JSON文件路径(示例用glob,也可根据业务逻辑自定义筛选) file_paths = glob.glob("/path/to/target_files/*.json") # 批量读取 df = spark.read.json( path=file_paths, schema=my_schema )
额外性能优化建议
- 保持schema指定:你当前指定
schema=my_schema的做法非常关键,避免Spark自动采样推断schema(百万文件场景下推断schema会耗时极长)。 - 调整分区大小:如果单个JSON文件极小,可通过
spark.sql.files.maxPartitionBytes参数调整分区的最大字节数(默认128MB),让Spark自动合并小文件到同一个分区,减少任务调度开销:spark.conf.set("spark.sql.files.maxPartitionBytes", "256m") # 设置为256MB - 避免小文件泛滥:后续如果需要写出结果,可使用
coalesce或repartition合并分区,减少输出文件数量。
内容的提问来源于stack exchange,提问作者python152
相关产品推荐
相关产品推荐

