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

