如何用PySpark高效读取过滤列名不同的大量Parquet文件?
问题分析与优化方案
你的代码耗时超3小时的核心原因是循环遍历单个文件读取+逐个Union——这完全浪费了Spark的分布式处理能力,相当于把Spark当成单进程工具在用;再加上Pandas的iterrows()本身就是效率极低的遍历方式,每次读小文件都会触发Spark的作业调度,大量重复开销直接拖慢了整体速度。
具体优化步骤
1. 批量处理路径,放弃单文件循环
Spark原生支持读取路径列表或用通配符批量读取,不需要逐个文件循环触发读取操作。
2. 统一处理列名差异
针对Data_referencia的不同命名,读取时直接重命名为固定列名,避免每个文件单独处理。
3. 推迟去重操作
原代码每个文件都做distinct()会额外增加shuffle开销,应该在所有文件合并完成后再统一去重,减少重复计算。
4. 替换低效的iterrows()
改用Pandas的itertuples()遍历路径,效率比iterrows()高数倍。
优化后代码示例
import pyspark.sql.functions as F from functools import reduce from pyspark.sql import DataFrame # 收集所有目标路径、文件名及对应列名,用itertuples替代iterrows path_info_list = [] for row in df.itertuples(): full_path = adl_gen2_full_url(DATALAKE, FILESYSTEM, '/APPLICATION/' + row.Ingested_Path) path_info_list.append( (full_path, row.Nome_Arquivo, row.Data_referencia) ) # 定义单文件读取处理函数 def process_single_file(info): path, filename, ref_col = info try: df = spark.read.parquet(path) # 统一列名,提取目标字段,添加文件名标识 return df.withColumnRenamed(ref_col, 'data_referencia') \ .select('data_referencia', 'data_upload', 'data_processamento') \ .withColumn("nome_arquivo", F.lit(filename)) except requests.exceptions.RequestException as e: print(f"连接失败: {path}") return None except Exception as e: print(f"处理出错 {path}: {str(e)}") return None # 并行处理所有文件(利用Spark分布式能力) valid_dfs = [df for df in map(process_single_file, path_info_list) if df is not None] # 合并所有DataFrame,用unionByName避免列顺序问题 if valid_dfs: final_df = reduce(DataFrame.unionByName, valid_dfs) # 统一去重+过滤无效值 final_df = final_df.distinct().filter(F.col('data_referencia') != 'NaT') else: # 空数据场景兜底 final_df = spark.createDataFrame([], schema='data_referencia string, data_upload string, data_processamento string, nome_arquivo string')
额外提速建议
- 合并小文件:如果你的Parquet都是几十MB级的小文件,先合并成1GB左右的大文件,减少Spark的文件IO开销。
- 利用分区裁剪:如果文件是按日期等字段分区存储的,读取时直接指定分区路径,减少读取的数据量。
- 缓存结果:如果后续还要对
final_df做操作,添加final_df.cache()把数据缓存到内存,避免重复计算。
内容的提问来源于stack exchange,提问作者Gizelly
相关产品推荐
相关产品推荐

