【发布时间】: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