免责声明:我没有明确的答案,也不想充当权威消息来源,但在 Spark 2.2+ 中的拼花支持上花了一些时间,希望我的回答能帮助我们所有人获得更接近正确答案。
S3 上的 Parquet 是避免从 S3 中提取未使用列的数据,并且只检索它需要的文件块,还是提取整个文件?
我使用的是我今天从 master 构建的 Spark 2.3.0-SNAPSHOT。
parquet 数据源格式由ParquetFileFormat 处理,FileFormat。
如果我是正确的,阅读部分由buildReaderWithPartitionValues 方法处理(覆盖FileFormat's)。
buildReaderWithPartitionValues 仅用于当FileSourceScanExec 物理运算符被请求用于所谓的输入 RDD,实际上是单个 RDD 以在执行WholeStageCodegenExec 时生成内部行。
话虽如此,我认为回顾 buildReaderWithPartitionValues 所做的事情可能会让我们更接近最终答案。
当您查看 the line 时,您可以确信我们的方向是正确的。
// 启用过滤器下推时尝试下推过滤器。
该代码路径取决于spark.sql.parquet.filterPushdown Spark 属性,即is turned on by default。
spark.sql.parquet.filterPushdown 设置为 true 时启用 Parquet 过滤器下推优化。
这将我们引向 parquet-hadoop 的 ParquetInputFormat.setFilterPredicateiff 过滤器已定义。
if (pushed.isDefined) {
ParquetInputFormat.setFilterPredicate(hadoopAttemptContext.getConfiguration, pushed.get)
}
稍后,当代码回退到 parquet-mr 时使用过滤器(而不是使用所谓的矢量化 parquet 解码阅读器),代码会变得更有趣。那是我不太明白的部分(除了我在代码中看到的)。
请注意,矢量化 parquet 解码阅读器由默认打开的 spark.sql.parquet.enableVectorizedReader Spark 属性控制。
提示:要了解使用了if 表达式的哪一部分,请为org.apache.spark.sql.execution.datasources.parquet.ParquetFileFormat 记录器启用DEBUG 记录级别。
要查看所有下推过滤器,您可以打开 INFO 记录级别的 org.apache.spark.sql.execution.FileSourceScanExec 记录器。你应该see the following in the logs:
INFO Pushed Filters: [pushedDownFilters]
我确实希望,如果它不是一个确定的答案,它会有所帮助,并且有人会在我离开的地方拿起它,以便尽快做出它。 希望最后死去 :)