PySpark按ID关联多个Parquet文件合并为单表的实现方案
PySpark多Parquet按公共ID合并最优方案
现有方案问题说明
- 循环Join方案:每执行一次Join就会触发一次Shuffle,n个文件需要执行n-1次Shuffle,文件量较大时会出现血缘过长、Shuffle冗余、数据倾斜等问题,性能很差。
- 直接批量读取方案:
spark.read.parquet(*file_path_list)默认是对多文件做纵向行合并(Union),只有所有文件列名完全一致时才能正确拼接。由于每个文件的业务字段名互不相同,Spark默认不会做Schema合并,最终只会返回单个文件的字段结构,完全达不到横向按ID关联的效果。
最优实现代码
核心思路:将多次Join的N次Shuffle优化为单次分组聚合的1次Shuffle,性能随文件数量增多提升越明显,代码逻辑和多表全外连接效果完全等价。
from pyspark.sql import functions as F from functools import reduce file_path_list = ["file1.parquet", "file2.parquet", "file3.parquet"] join_keys = ["id_foo", "id_bar"] # 逐个读取所有Parquet文件 df_collection = [spark.read.parquet(path) for path in file_path_list] # 纵向合并所有表,缺失字段自动填充null union_all_df = reduce( lambda df1, df2: df1.unionByName(df2, allowMissingColumns=True), df_collection ) # 按公共ID分组,聚合取出每个业务字段的非空值,等价于按ID做全量关联 business_columns = [col for col in union_all_df.columns if col not in join_keys] agg_rules = [F.first(col, ignorenulls=True).alias(col) for col in business_columns] final_merged_df = union_all_df.groupBy(join_keys).agg(*agg_rules) # 写出结果,数据量小于10G可加coalesce(1)输出单文件,大数据量去掉该参数分区输出 final_merged_df.coalesce(1).write.mode("overwrite").parquet("merged_result.parquet")
方案优势
- 性能更高:全程仅触发1次分组Shuffle,相比循环Join的N次Shuffle,文件量越大性能差距越显著,通常有3~10倍的性能提升。
- 数据更全:不需要指定某张表作为左连接基表,不会因为基表缺失ID组合导致数据遗漏,等价于所有表全外连接的效果。
- 维护简单:新增待合并文件时,只需要把路径加入
file_path_list即可,不需要修改关联逻辑。
注意事项
- 该方案默认每个
(id_foo,id_bar)组合在单个Parquet文件中最多出现1条记录,如果单文件内存在重复ID,可根据业务需求将first聚合替换为max/min/collect_list等逻辑。 - 如果所有待合并Parquet存放在同一个目录下,不需要手动枚举文件路径,直接读取目录即可:
spark.read.parquet("your/parquet/dir/path/")。 - 若合并后总数据量超过10G,建议去掉
coalesce(1)参数,让Spark自动分区输出Parquet,避免单文件过大导致读写性能下降。
内容的提问来源于stack exchange,提问作者PeCaDe
相关产品推荐
相关产品推荐

