【问题标题】:Speed up InMemoryFileIndex for Spark SQL job with large number of input files为具有大量输入文件的 Spark SQL 作业加速 InMemoryFileIndex
【发布时间】:2018-11-02 14:10:01
【问题描述】:

我有一个用 Java 编码的 apache spark sql 作业(使用数据集),它的输入范围在 70,000 到 150,000 个文件之间。

构建 InMemoryFileIndex 似乎需要 45 分钟到 1.5 小时。

这段时间没有日志,网络使用率很低,几乎没有CPU使用率。

这是我在标准输出中看到的示例:

24698 [main] INFO org.spark_project.jetty.server.handler.ContextHandler  - Started o.s.j.s.ServletContextHandler@32ec9c90{/static/sql,null,AVAILABLE,@Spark}
25467 [main] INFO org.apache.spark.sql.execution.streaming.state.StateStoreCoordinatorRef  - Registered StateStoreCoordinator endpoint
2922000 [main] INFO org.apache.spark.sql.execution.datasources.InMemoryFileIndex  - Listing leaf files and directories in parallel under: <a LOT of file url's...>
2922435 [main] INFO org.apache.spark.SparkContext  - Starting job: textFile at SomeClass.java:103

在这种情况下,有 45 分钟基本上什么也没发生(据我所知)。

我使用以下方式加载文件:

sparkSession.read().textFile(pathsArray)

有人可以解释一下 InMemoryFileIndex 中发生了什么,我怎样才能加快这一步?

【问题讨论】:

  • 想通了,我向 sparkSession.read().textFile(70,000 个路径...) 传递了太多文件。 Spark 显然会检查为孩子传递的每条路径,以防它是文件夹或 glob 模式。这需要很长时间。
  • SimpleSam5 您能否详细说明一下您是如何发现 Spark 检查每条路径的?谢谢!
  • 我相信我找到了代码,它按顺序单独地对每个路径进行全局化。可能有办法将其并行化,但我发现最好的解决方案是尽可能少地发送路径,尽可能使用复杂的 glob 表达式。例如,我最终使用了一个充满“{foo,bar,foobar,barfoo}”表达式的路径。

标签: apache-spark apache-spark-sql


【解决方案1】:

InMemoryFileIndex 负责分区发现(并因此进行分区修剪),它正在执行文件列表,并且它可能会运行并行作业,如果您有很多文件,这可能需要一些时间,因为它必须为每个文件编制索引。执行此操作时,Spark 会收集有关文件的一些基本信息(例如它们的大小)来计算一些在查询计划期间使用的基本统计信息。如果您想在每次读取数据时避免这种情况,您可以使用 metastore 和 saveAsTable() 命令将数据保存为数据源表(Spark 2.1 支持),并且此分区发现将只执行一次并且信息将保存在元存储中。然后就可以使用 Metastore 读取数据了

sparkSession.read.table(table_name)

它应该很快,因为这个分区发现阶段将被跳过。我推荐看this Spark 峰会上讨论这个问题的演讲。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-12-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-10-12
    • 2016-01-11
    • 2019-05-09
    相关资源
    最近更新 更多