【问题标题】:Does Spark support true column scans over parquet files in S3?Spark 是否支持对 S3 中的 parquet 文件进行真正的列扫描?
【发布时间】:2017-02-03 19:25:53
【问题描述】:

Parquet 数据存储格式的一大好处是it's columnar。如果我有一个包含数百列的“宽”数据集,但我的查询只涉及其中的几列,则可以只读取存储这几列的数据,而跳过其余列。

据推测,此功能通过读取 parquet 文件头部的一些元数据来工作,该文件指示每列在文件系统上的位置。然后阅读器可以在磁盘上查找以仅读取必要的列。

有谁知道 spark 的默认 parquet reader 是否在 S3 上正确实现了这种选择性搜索?我认为it's supported by S3,但理论支持与正确利用该支持的实现之间存在很大差异。

【问题讨论】:

  • 我问这个是因为我注意到 spark/parquet 宣传的一些功能还没有正确实现,例如谓词下推,它只允许读取某些分区。我发现这很令人惊讶,并开始想知道实木复合地板/火花到底有多少像宣传的那样有效。

标签: apache-spark amazon-s3 apache-spark-sql parquet


【解决方案1】:

这需要分解

  1. Parquet 代码是否从 spark 获取谓词(是)
  2. parquet 是否会尝试使用 Hadoop FileSystem seek() + read()readFully(position, buffer, length) 调用选择性地仅读取这些列?是的
  3. S3 连接器是否将这些文件操作转换为高效的 HTTP GET 请求?在 Amazon EMR 中:是的。在 Apache Hadoop 中,您需要在类路径上安装 hadoop 2.8 并正确设置 spark.hadoop.fs.s3a.experimental.fadvise=random 以触发随机访问。

Hadoop 2.7 及更早版本对文件的激进 seek() 处理不好,因为它们总是启动 GET offset-end-of-file,对下一次搜索感到惊讶,不得不中止该连接,重新打开一个新的 TCP/ HTTPS 1.1 连接(慢,CPU 重),再做一次,重复。随机 IO 操作会影响批量加载 .csv.gz 等内容,但对于获得 ORC/Parquet 性能至关重要。

您无法在 Hadoop 2.7 的 hadoop-aws JAR 上获得加速。如果需要,您需要更新 hadoop*.jar 和依赖项,或者针对 Hadoop 2.8 从头开始​​构建 Spark

请注意,Hadoop 2.8+ 还有一个不错的小功能:如果您在 S3A 文件系统客户端上的日志语句中调用 toString(),它会打印出所有文件系统 IO 统计信息,包括在查找中丢弃了多少数据,中止TCP 连接 &c。帮助您弄清楚发生了什么。

2018-04-13 警告::不要尝试将 Hadoop 2.8+ hadoop-aws JAR 与 hadoop-2.7 JAR 集的其余部分一起放在类路径中,并期望看到任何加速。您将看到的只是堆栈跟踪。您需要更新所有 hadoop JAR 及其传递依赖项。

【讨论】:

  • 感谢分解!我认为分解是其他答案所缺乏的。
【解决方案2】:

免责声明:我没有明确的答案,也不想充当权威消息来源,但在 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]

我确实希望,如果它不是一个确定的答案,它会有所帮助,并且有人会在我离开的地方拿起它,以便尽快做出它。 希望最后死去 :)

【讨论】:

    【解决方案3】:

    不,不完全支持谓词下推。当然,这取决于:

    • 具体用例
    • Spark 版本
    • S3 连接器类型和版本

    为了检查您的具体用例,您可以在 Spark 中启用 DEBUG 日志级别,然后运行您的查询。然后,您可以查看 S3 (HTTP) 请求期间是否存在“搜索”以及实际发送了多少请求。像这样的:

    17/06/13 05:46:50 DEBUG wire: http-outgoing-1 >> "GET /test/part-00000-b8a8a1b7-0581-401f-b520-27fa9600f35e.snappy.parquet HTTP/1.1[\r][\n]" .... 17/06/13 05:46:50 DEBUG wire: http-outgoing-1 << "Content-Range: bytes 0-7472093/7472094[\r][\n]" .... 17/06/13 05:46:50 DEBUG wire: http-outgoing-1 << "Content-Length: 7472094[\r][\n]"

    这是由于 Spark 2.1 无法根据存储在 Parquet 文件中的元数据计算数据集中所有行的 COUNT(*) 而最近打开的问题报告示例:https://issues.apache.org/jira/browse/SPARK-21074

    【讨论】:

    • Michael,与其捆绑的 Hadoop JAR 版本不如说是火花; HDP 和 CDH 中的那些执行“惰性”搜索,并且,如果您启用随机 IO,则可以高效读取列式数据。关于SPARK-21074,JIRA等你升级后体验;如果你没有得到答案,它可能会被关闭为“固定/无法重现”
    【解决方案4】:

    spark 的 parquet reader 和其他 InputFormat 一样,

    1. 所有 inputFormat 都没有对 S3 有任何特别之处。输入格式可以从 LocalFileSystem 、 Hdfs 和 S3 读取,无需为此进行特殊优化。

    2. Parquet InpuTFormat 根据您询问的列将选择性地为您读取列。

    3. 如果你想确定无疑(尽管下推谓词在最新的 spark 版本中有效)手动选择列并编写转换和操作,而不是依赖于 SQL

    【讨论】:

    • 感谢您的回答,但即使在阅读之后,仍然不清楚最近的火花分布是否真正支持谓词下推。我正在寻找一个答案,该答案要么深入研究从 s3 读取镶木地板时调用的输入阅读器的特定实现,要么执行经验测试。见stackoverflow.com/a/41609999/189336——有一个令人惊讶的结果表明过滤器下推在 s3 上被破坏了。
    • 注意火花版本。在早期版本中存在谓词下推问题,但从 2 开始(肯定是 2.2),这已修复
    猜你喜欢
    • 2016-09-07
    • 1970-01-01
    • 2017-01-10
    • 2019-09-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-01-26
    相关资源
    最近更新 更多