【发布时间】:2020-01-21 23:52:24
【问题描述】:
我正在分析我的 spark 应用程序在小数据集的情况下的性能。我有一个血统图,如下所示:
someList.toDS()
.repartition(x)
.mapPartitions(func1)
.mapPartitions(func2)
.mapPartitions(func3)
.filter(cond1)
.count()
我有一个由 2 个节点组成的集群,每个节点有 8 个核心。执行器配置为使用 4 个核心。因此,当应用程序运行时,会出现四个执行程序,每个执行程序使用 4 个内核。
我观察到每个线程上至少(通常只有)1 个任务(即总共 16 个任务)比其他任务花费的时间要长得多。例如,在一次运行中,这些任务大约需要 15-20 秒,而其他任务则需要一秒或更短的时间。
在分析代码时,我发现瓶颈在上面的 func3 中:
def func3 = (partition: Iterator[DomainObject]) => {
val l = partition.toList // This takes almost all of the time
val t = doSomething(l)
}
从 Iterator 到 List 的转换几乎占用了所有时间。
分区大小非常小(在某些情况下甚至小于 50)。即使这样,不同分区的分区大小几乎是一致的,但每个线程只有一个任务占用时间。
我会假设当func3 在执行器上运行任务时,该分区中的数据已经存在于执行器上。不是这样吗? (在func3的执行过程中,它是否会遍历整个数据集以某种方式过滤掉这个分区的数据?!)
否则,为什么从少于 50 个对象的 Iterator 转换到 List 会占用这么多时间?
我注意到的另一件事(不确定这是否相关)是这些任务的 GC 时间(根据 spark ui)对于所有这 16 个任务来说也是异常一致的 2 秒,与其他任务相比(即便如此,2 秒)
【问题讨论】:
标签: scala apache-spark