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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 14:40:47