【问题标题】:When is the data in a spark partition actually realised?spark分区中的数据是什么时候真正实现的?
【发布时间】: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


    【解决方案1】:
    1. 第一次实现是在repartition() 期间
    2. 第二个是在filter 操作之后,所有三个mapPartitions 开始执行(当调用count 操作时)。根据您在每个函数中的doSomething(),这将取决于 DAG 的创建方式以及花费时间的位置,因此您可以进行优化。

    【讨论】:

    • 为什么会在filter 中实现? filter 与其他转换有何不同?另外,您能否提供更多关于 DAG 如何依赖于传递给 mapPartitions 的函数的信息?
    • 嗨,我的意思是在filter 操作之后,我在第 2 点中详细说明了它。DAG 通常是根据您拥有的三个函数形成的,它检查哪个数据集被重用,或者如何在每个函数中重用,可能是第一个函数的中间数据集,可以在第二个和第三个函数中重用,从而形成 DAG。我试图含糊其辞地解释它,但没有看到功能很难解释如何改进执行。
    • "它检查哪个数据集被重用" -- 所以,所有的函数都是(it: Iterator[SomeDomainObject]) => Iterator[AnotherDomainObject]形式的闭包,并根据一些业务逻辑在内部转换对象(涉及一些I/O和其他东西)。在任何情况下,内部都不会产生火花。我不明白:“......哪个数据集被重用..”:只有一个数据集,不是吗?
    • 是的,但是 ds 被拆分到不同的分区,并被洗牌。这些分区上的操作由 spark 优化。
    • 嗨.. 是的,你是对的:ds 被分区和洗牌。但是,在上述沿袭中,重新分区发生后,mapPartitions 不需要改组(因此所有这些都是一项工作的一部分)。一旦分区在执行程序上实现,操作将由 spark 优化,但不应涉及任何 shuffle 写入(写入计数除外,这是本例中的操作)。我发现了减慢任务的瓶颈,它与火花无关。我将在答案中更新详细信息。 (感谢您的宝贵时间!)
    【解决方案2】:

    似乎只要任务开始执行,分区中的数据就可用(或者,至少在遍历该数据时没有任何重大成本,正如问题所显示的那样。)

    上述代码中的瓶颈实际上在 func2 中(我没有正确调查!),并且是因为 scala 中迭代器的惰性特性。这个问题根本与火花无关。

    首先,上面的mapPartitions 调用中的函数出现 被链接起来并像这样调用:func3( func2( func1(Iterator[A]) ) ) : Iterator[B]。所以作为func2 的输出生成的Iterator 直接馈送到func3

    其次,对于上述问题func1(和func2)定义为:

    func1(x: Iterator[A]) => Iterator[B] = x.map(...).filter...
    

    由于它们采用迭代器并将它们映射到不同的迭代器,因此它们不会立即执行。但是当func3 被执行时,partition.toList 导致func2 中的map 闭包被执行。在分析时,func3 似乎一直在花费时间,而 func2 的代码会减慢应用程序的速度。

    (针对上述问题,func2 包含一些案例对象到 json 字符串的序列化。它似乎执行了一些耗时的隐式代码,仅针对每个线程上的第一个对象。因为每个线程发生一次,每个线程只有一个任务,需要很长时间,并解释了上面的事件时间线。)

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-07-30
      • 1970-01-01
      • 1970-01-01
      • 2011-02-14
      • 1970-01-01
      • 2022-11-11
      • 1970-01-01
      相关资源
      最近更新 更多