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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 03:05:30