【问题标题】:Number of partitions when creating a Spark dataframe创建 Spark 数据帧时的分区数
【发布时间】:2020-01-07 02:53:01
【问题描述】:

这个问题已在其他帖子中提出,但似乎我的问题不适合其中任何一个。

我在本地模式下使用 Spark 2.4.4,我将 master 设置为 local[16] 以使用 16 个内核。我还在 Web UI 中看到已分配 16 个内核。

我创建了一个导入大约 8MB 的 csv 文件的数据框,如下所示:

val df = spark.read.option("inferSchema", "true").option("header", "true").csv("Datasets/globalpowerplantdatabasev120/*.csv")

最后我打印出数据框的分区数:

df.rdd.partitions.size

res5: Int = 2

答案是 2。

为什么?据我阅读,分区的数量取决于默认设置为等于核心数(16)的执行器数量。

我尝试使用 spark.default.Parallelism = 4 和/或 spark.executor.instances = 4 设置 esecutor 的数量并启动了一个新的 spark 对象,但分区数量没有任何变化。

有什么建议吗?

【问题讨论】:

  • 8MB 太低了。

标签: apache-spark apache-spark-sql


【解决方案1】:

当您使用 Spark 读取文件时,分区数计算为 defaultMinPartitions 与基于 hadoop 输入拆分大小除以块大小计算的拆分数之间的最大值。由于您的文件很小,因此您获得的分区数为 2,这是两者中的最大值。

默认的 defaultMinPartitions 计算为

def defaultMinPartitions: Int = math.min(defaultParallelism, 2)

请查看https://github.com/apache/spark/blob/e9f983df275c138626af35fd263a7abedf69297f/core/src/main/scala/org/apache/spark/SparkContext.scala#L2329

【讨论】:

  • 默认MinPartitions,根据公式可以是1或2。但这是最小个分区数。不是分区本身的数量。问题依然存在,分区数是怎么设置的?
  • 分区数为最大值(defaultMinPartitions,根据hadoop输入分割大小除以块大小计算的分割数);
  • 感谢您的澄清。这对普通文件系统也有效吗?我在哪里可以找到这些参数?
  • Spark 在后台使用 Hadoop InputFilFormat,它将按输入块读取分区。所以你需要阅读 Hadoop 块大小到达逻辑。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-02-21
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多