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

基于Parquet的PySpark查询未利用分区裁剪,原因何在?

问题:PySpark Parquet分区过滤未生效导致查询耗时激增

项目数据概况

  • 时序数据,包含约6个字段,每日产生100万行、22MB数据
  • 按日期(格式YYYY-MM-DD)作为一级分区,satellite_id(取值为整数1或2)作为二级分区,每日对应2个本地Parquet文件
  • 数据写入前已按整数时间戳(epoch)排序,存储在本地磁盘,集群为开发机单主节点加单个工作节点

预期与实际现象

  • 预期逻辑:查询时指定分区列temporal_partition_key后,Spark应直接跳过无关分区目录,仅扫描目标日期的文件,因此即使数据量增长,查询耗时应基本稳定
  • 实际情况:仅存一周数据时,查询小段时间切片耗时约300ms;存储一年数据(共约700个文件)时,耗时跃升至5000ms。修正查询加入分区过滤条件后,耗时无明显变化

执行代码

@classmethod
def select(cls, fields: Iterable[str], from_dt: datetime, to_dt: datetime) -> List[pyspark.Row]:
    table_name = 'someTableName'

    spark = get_spark_session()
    table_df = spark.read.parquet(cls.parquet_path)
    table_df.createOrReplaceTempView(table_name)

    fields_str = ', '.join(fields)
    from_rcvtime_intg = cls.dt_to_rcvtime_intg(max(from_dt, timestamp_epoch))
    to_rcvtime_intg = cls.dt_to_rcvtime_intg(to_dt)

    query = f"""
        select {fields_str}, rcv_timestamp
        from {table_name}
        where
            {PARQUET_TEMPORAL_PARTITION_KEY} >= '{from_dt.date().isoformat()}'
            and {PARQUET_TEMPORAL_PARTITION_KEY} <= '{to_dt.date().isoformat()}'
            and rcvtime_intg >= {from_rcvtime_intg}
            and rcvtime_intg <= {to_rcvtime_intg}
    """

    records_df = spark.sql(query)
    return records_df.collect()

查询计划

== Parsed Logical Plan ==
'Project ['ang_accl_z, 'lin_accl_x, 'ang_accl_y, 'rcv_timestamp, 'lin_accl_z, 'satellite_id, 'lin_accl_y, 'ang_accl_x, 'rcv_timestamp]
+- 'Filter ((('temporal_partition_key >= 2023-06-06) AND ('temporal_partition_key <= 2023-06-06)) AND (('rcvtime_intg >= 739324800) AND ('rcvtime_intg <= 739324860)))
   +- 'UnresolvedRelation [GRACEFO_1A_full_resolution], [], false
== Analyzed Logical Plan ==
ang_accl_z: double, lin_accl_x: double, ang_accl_y: double, rcv_timestamp: timestamp, lin_accl_z: double, satellite_id: int, lin_accl_y: double, ang_accl_x: double, rcv_timestamp: timestamp
Project [ang_accl_z#7, lin_accl_x#2, ang_accl_y#6, rcv_timestamp#8, lin_accl_z#4, satellite_id#9, lin_accl_y#3, ang_accl_x#5, rcv_timestamp#8]
+- Filter (((temporal_partition_key#10 >= cast(2023-06-06 as date)) AND (temporal_partition_key#10 <= cast(2023-06-06 as date))) AND ((rcvtime_intg#0L >= cast(739324800 as bigint)) AND (rcvtime_intg#0L <= cast(739324860 as bigint))))
   +- SubqueryAlias gracefo_1a_full_resolution
      +- View (`GRACEFO_1A_full_resolution`, [rcvtime_intg#0L,rcvtime_frac#1L,lin_accl_x#2,lin_accl_y#3,lin_accl_z#4,ang_accl_x#5,ang_accl_y#6,ang_accl_z#7,rcv_timestamp#8,satellite_id#9,temporal_partition_key#10])
         +- Relation [rcvtime_intg#0L,rcvtime_frac#1L,lin_accl_x#2,lin_accl_y#3,lin_accl_z#4,ang_accl_x#5,ang_accl_y#6,ang_accl_z#7,rcv_timestamp#8,satellite_id#9,temporal_partition_key#10] parquet
== Optimized Logical Plan ==
Project [ang_accl_z#7, lin_accl_x#2, ang_accl_y#6, rcv_timestamp#8, lin_accl_z#4, satellite_id#9, lin_accl_y#3, ang_accl_x#5, rcv_timestamp#8]
+- Filter ((isnotnull(temporal_partition_key#10) AND isnotnull(rcvtime_intg#0L)) AND (((temporal_partition_key#10 >= 2023-06-06) AND (temporal_partition_key#10 <= 2023-06-06)) AND ((rcvtime_intg#0L >= 739324800) AND (rcvtime_intg#0L <= 739324860))))
   +- Relation [rcvtime_intg#0L,rcvtime_frac#1L,lin_accl_x#2,lin_accl_y#3,lin_accl_z#4,ang_accl_x#5,ang_accl_y#6,ang_accl_z#7,rcv_timestamp#8,satellite_id#9,temporal_partition_key#10] parquet
== Physical Plan ==
*(1) Project [ang_accl_z#7, lin_accl_x#2, ang_accl_y#6, rcv_timestamp#8, lin_accl_z#4, satellite_id#9, lin_accl_y#3, ang_accl_x#5, rcv_timestamp#8]
+- *(1) Filter ((isnotnull(rcvtime_intg#0L) AND (rcvtime_intg#0L >= 739324800)) AND (rcvtime_intg#0L <= 739324860))
   +- *(1) ColumnarToRow
      +- FileScan parquet [rcvtime_intg#0L,lin_accl_x#2,lin_accl_y#3,lin_accl_z#4,ang_accl_x#5,ang_accl_y#6,ang_accl_z#7,rcv_timestamp#8,satellite_id#9,temporal_partition_key#10] Batched: true, DataFilters: [isnotnull(rcvtime_intg#0L), (rcvtime_intg#0L >= 739324800), (rcvtime_intg#0L <= 739324860)], Format: Parquet, Location: InMemoryFileIndex(1 paths)[file:/nomount/masschange/data_volume_mount/full-resolution], PartitionFilters: [isnotnull(temporal_partition_key#10), (temporal_partition_key#10 >= 2023-06-06), (temporal_parti..., PushedFilters: [IsNotNull(rcvtime_intg), GreaterThanOrEqual(rcvtime_intg,739324800), LessThanOrEqual(rcvtime_int..., ReadSchema: struct&lt;rcvtime_intg:bigint,lin_accl_x:double,lin_accl_y:double,lin_accl_z:double,ang_accl_x:doubl...

疑问

是否从根本上误解了PySpark Parquet的分区机制?还是除了在SQL查询中包含分区字段外,还需要额外配置才能让Spark利用分区过滤?

内容的提问来源于stack exchange,提问作者A.D.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 16:15:53