如何用PySpark批量读取Parquet分区数据?
在PySpark中便捷读取指定分区的Parquet数据
你当前手动指定多个分区路径的方式虽然可行,但扩展性差。PySpark有两种更便捷的方式替代:
方式一:读取根目录后过滤分区列(推荐)
Spark支持分区谓词下推,读取整个分区根目录后,用filter指定要读取的分区值,Spark会自动只加载符合条件的分区文件,不会全量扫描:
from pyspark.sql.functions import col # 定义需要读取的年份集合,可直接扩展到15个年份 target_years = [_year, _year-1, _year-2] base_dir = f"{_data['PQ_path'] + _data['PQ_name_partitioned']}/" df = spark_sessions.read.option('compression', 'gzip')\ .option("basePath", base_dir)\ .parquet(base_dir)\ .filter( # 注意:根据你的实际分区列名调整,例子里同时出现ano-mes和Anio,需确认统一命名 col("Anio").isin(target_years) )
注:
basePath参数用于确保Spark将分区目录的键值识别为DataFrame的列,避免分区列名变成路径的一部分。
方式二:动态生成分区路径列表
如果不想读取根目录,可以用Python循环自动生成目标路径,避免手动编写多条路径:
base_dir = f"{_data['PQ_path'] + _data['PQ_name_partitioned']}/" # 生成目标路径,这里假设分区列为Anio,按需调整 target_paths = [f"{base_dir}Anio={year}" for year in range(_year-2, _year+1)] df = spark_sessions.read.option('compression', 'gzip')\ .option("basePath", base_dir)\ .parquet(*target_paths)
关于你尝试的filters参数失败的原因
PySpark的parquet读取接口不支持filters参数,这个参数是Pandasread_parquet的专属参数,Spark中需要用上述两种方式实现分区过滤。
内容的提问来源于stack exchange,提问作者Lucio Diprè
相关产品推荐
相关产品推荐

