【问题标题】:Did Spark 2.3 change the way it process small files?Spark 2.3 是否改变了它处理小文件的方式?
【发布时间】:2018-09-06 22:43:13
【问题描述】:

我刚开始使用 Spark 2+(2.3 版本),在查看 Spark UI 时发现了一些奇怪的现象。 我有一个 HDFS 集群中的目录列表,其中包含总共 24000 个小文件。

当我想对它们运行 Spark 操作时,Spark 1.5 会为每个输入文件生成一个单独的任务,就像我之前使用的那样。我知道每个 HDFS 块(在我的情况下,一个小文件就是一个块)在 Spark 中生成一个分区,每个分区由一个单独的任务处理。

另外,my_dataframe.rdd.getNumPartitions() 命令输出 24000。

关于 Spark 2.3 在同一输入上,命令 my_dataframe.rdd.getNumPartitions() 输出 1089。Spark UI 还为我的 Spark 操作生成 1089 个任务。您还可以看到 spark 2.3 中生成的作业数量大于 1.5

两个 Spark 版本的代码相同(我需要稍微更改数据框、路径和列名,因为它是工作代码):

%pyspark
dataframe = sqlContext.\
                read.\
                parquet(path_to_my_files)
dataframe.rdd.getNumPartitions()
dataframe.\
    where((col("col1") == 21379051) & (col("col2") == 2281643649) & (col("col3") == 229939942)).\
    select("col1", "col2", "col3").\
    show(100, False)

这是由

生成的物理计划
dataframe.where(...).select(...).explain(True)
Spark 1.5
== Physical Plan ==
Filter (((col1 = 21379051) && (col2 = 2281643649)) && (col3 = 229939942))
 Scan ParquetRelation[hdfs://cluster1ns/path_to_file][col1#27,col2#29L,col3#30L]
Code Generation: true

Spark 2.3
== Physical Plan ==
*(1) Project [col1#0, col2#2L, col3#3L]
+- *(1) Filter (((isnotnull(col1#0) && (col1#0 = 21383478)) && (col2 = 2281643641)) && (col3 = 229979603))
   +- *(1) FileScan parquet [col1,col2,col3] Batched: false, Format: Parquet, Location: InMemoryFileIndex[hdfs://cluster1ns/path_to_file..., PartitionFilters: [], PushedFilters: [IsNotNull(col1)], ReadSchema: struct<col1:bigint,col2:bigint,col3:bigint>....

以上作业是使用 pyspark 从 zeppelin 生成的。 还有其他人用 spark 2.3 遇到过这种情况吗? 我不得不说我喜欢处理多个小文件的新方法,但我也想了解 Spark 内部可能发生的变化。

我在互联网上搜索了最新一本书“Spark 权威指南”,但没有找到任何有关 Spark 生成工作物理计划的新方法的信息。

如果您有任何链接或信息,将会很有趣。 谢谢!

【问题讨论】:

  • 能否提供职位代码?
  • 嗨@addmeaning,我在帖子上下文中添加了代码。谢谢!
  • 代码非常简单。您能否尝试为两个版本的 spark 运行 dataframe.explain(True),以测试代码是否转换为不同的操作集?
  • 好主意,我在选择的结果数据帧上运行解释(真)。我添加了物理计划(再次进行了一些代码混淆)
  • 2.x doesn't use Hadoop configuration to compute splits。所以这说明了分区的数量。

标签: apache-spark pyspark bigdata apache-spark-2.0


【解决方案1】:

Spark 2.3 configuration

|spark.files.maxPartitionBytes| 134217728 (128 MB) |读取文件时打包到单个分区的最大字节数。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-02-22
    • 1970-01-01
    • 2013-11-03
    • 2017-09-16
    • 1970-01-01
    • 1970-01-01
    • 2019-12-21
    • 1970-01-01
    相关资源
    最近更新 更多