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

Spark读取分区Parquet时的过滤下推与S3目录遍历疑问

Spark分区Parquet文件过滤下推与S3遍历问题

我在Spark中尝试读取一组按date、hour分区的Parquet文件,为简化过滤逻辑编写了包装函数,支持根据分区列进行可选过滤。不清楚需满足哪些条件才能让Spark将过滤下推,避免遍历存储Parquet的S3存储桶所有叶子节点——是否必须在实例化events变量时就执行过滤?

def readFromParquet(session: SparkSession, start: Option[Timestamp] = Option.empty, end: Option[Timestamp] = Option.empty): Dataset[Event] = {
    import session.implicits._

    var events = session
      .read
      .parquet("s3a://precomputes/events")

    if (start.isDefined) {
      val startDate = Date.valueOf(start.get.toLocalDateTime.toLocalDate)
      val hour = start.get.toInstant.atZone(ZoneOffset.UTC).getHour
      events = events.where(
        'date > startDate ||
          ('date === startDate && 'hour >= hour)
      )
    }
    if (end.isDefined) {
      val endDate = Date.valueOf(end.get.toLocalDateTime.toLocalDate)
      val hour = end.get.toInstant.atZone(ZoneOffset.UTC).getHour
      events = events.where(
        'date < endDate ||
          ('date === endDate && 'hour <= hour)
      )
    }

    events
      .as[Event]
      .alias("events")
  }

2023-04-05 更新

回应Andrew Long的评论,物理计划如下:

== Physical Plan ==
*(1) ColumnarToRow
+- FileScan parquet [...] Batched: true, DataFilters: [], Format: Parquet, Location: InMemoryFileIndex(1 paths)[s3a://precomputes/events], PartitionFilters: [((date#31 > 2023-04-05) OR ((date#31 = 2023-04-05) AND (hour#32 >= 12))), ((date#31 < 2023-04-05..., PushedFilters: [], ReadSchema: struct<...

看起来分区过滤已生效,但Spark仍遍历所有叶子节点,使用分区过滤时不应避免这种情况吗?

内容的提问来源于stack exchange,提问作者maxwellray

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 01:55:18