PySpark用coalesce(1)合并大量小JSON文件过慢,求更快替代方案?
解决PySpark合并大量小JSON文件的性能问题
核心问题分析
你当前的coalesce(1)确实是性能瓶颈:它会强制将所有数据压缩到单个Executor节点处理,不仅引发不必要的Shuffle,还会让单节点承担所有的IO和数据处理工作,面对百万级小文件时效率极低。
下面是针对Databricks+GCS场景的几种高效优化方案,按性能从高到低排序:
方案1:直接合并为换行分隔JSON(NDJSON),跳过JSON解析
如果接受换行分隔格式,直接读取文本文件合并是最快的方式——无需解析JSON结构,完全绕开Spark的DataFrame序列化/反序列化开销:
def merge_to_ndjson(bucket_path: str, data_type: str, date: str, mode: str = "overwrite"): input_path = f"{bucket_path}/{data_type}/{date}" output_path = f"{input_path}/merged_ndjson" # 读取所有小JSON文件为文本行(每个文件内容作为一行,对应一个JSON对象) text_df = spark.read.text(input_path, wholetext=True) # 直接合并为单个文本文件 text_df.coalesce(1).write.mode(mode).text(output_path) # (可选)将Spark生成的临时part文件重命名为目标文件名 temp_file = [f.path for f in dbutils.fs.ls(output_path) if f.path.endswith(".txt")][0] final_file = f"{input_path}/merged.json" dbutils.fs.mv(temp_file, final_file) dbutils.fs.rm(output_path, recurse=True)
优势:
- 速度最快,避免JSON解析的CPU开销
- 无需处理Schema一致性问题(即使小文件Schema有差异也能合并)
方案2:优化Spark读取配置,减少Shuffle开销
如果需要保留DataFrame操作(比如过滤、转换数据),可以先调整Spark读取参数,让读取阶段就合并小文件,减少后续分区压力:
def optimized_merge_df(bucket_path: str, data_type: str, date: str, mode: str = "overwrite"): input_path = f"{bucket_path}/{data_type}/{date}" output_path = input_path # 调整Spark配置:读取时自动合并小文件,减少初始分区数 spark.conf.set("spark.sql.files.maxPartitionBytes", "128m") # 按128MB为单位合并分区 spark.conf.set("spark.sql.files.openCostInBytes", "1073741824") # 降低小文件的打开成本阈值 # 读取JSON,此时Spark会自动合并小文件到更少的分区 df = spark.read.format("json").load(input_path) # 先合并到少量分区(根据总数据量估算,比如总数据10GB就用10个分区),再压缩到1个 # 先repartition减少分区数,比直接coalesce(1)的Shuffle开销小得多 df = df.repartition(10) df.coalesce(1).write.mode(mode).option("lineSep", "\n").json(output_path)
关键配置说明:
spark.sql.files.maxPartitionBytes:控制每个读取分区的大小,默认128MB,可根据小文件总容量调整spark.sql.files.openCostInBytes:设置打开一个文件的虚拟成本,值越大,Spark越倾向于合并更多小文件到一个分区
方案3:利用Databricks本地磁盘+Shell命令合并
如果小文件总容量在Worker节点本地磁盘范围内,可以先将文件下载到本地,用Shell命令高效合并后再上传回GCS:
import os def merge_via_local_disk(bucket_path: str, data_type: str, date: str): input_path = f"{bucket_path}/{data_type}/{date}" local_temp_dir = f"/tmp/{data_type}_{date}" merged_local_path = f"{local_temp_dir}/merged.json" output_gcs_path = f"{input_path}/merged.json" # 1. 将GCS上的小文件下载到本地临时目录 dbutils.fs.cp(input_path, f"file:{local_temp_dir}", recurse=True) # 2. 用Shell命令合并所有JSON文件为换行分隔格式 os.system(f"cat {local_temp_dir}/*.json > {merged_local_path}") # 3. 将合并后的文件上传回GCS dbutils.fs.cp(f"file:{merged_local_path}", output_gcs_path) # 4. 清理本地临时文件 dbutils.fs.rm(f"file:{local_temp_dir}", recurse=True)
优势:
- 利用Shell命令的原生文件合并能力,比Spark单分区处理更快
- 适合总数据量不大(比如几十GB以内)的场景
方案4:生成单个JSON数组(需注意Driver内存)
如果必须输出单个JSON数组,可以在读取文本后,在Driver端拼接数组(仅适合总数据量较小的场景,避免Driver OOM):
def merge_to_json_array(bucket_path: str, data_type: str, date: str): input_path = f"{bucket_path}/{data_type}/{date}" output_path = f"{input_path}/merged_array.json" # 读取所有小文件为文本内容 text_rdd = spark.sparkContext.wholeTextFiles(input_path).values() # 收集到Driver端,拼接成标准JSON数组 json_objects = text_rdd.collect() json_array = f"[{','.join(json_objects)}]" # 将数组写入GCS dbutils.fs.put(output_path, json_array, overwrite=True)
注意事项:
- 必须确保Driver节点内存足够容纳所有JSON内容
- 总数据量超过Driver内存时会引发OOM,不适合百万级小文件的大体积场景
内容的提问来源于stack exchange,提问作者Omer P
相关产品推荐
相关产品推荐

