基于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<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.
相关产品推荐
相关产品推荐

