【问题标题】:Apache Spark mapPartition strange behavior (lazy evaluation?)Apache Spark mapPartition 奇怪的行为(懒惰评估?)
【发布时间】:2017-08-02 22:04:19
【问题描述】:

我正在尝试使用这样的代码(在 Scala 中)记录 RDD 上每个 mapPartition 操作的执行时间:

rdd.mapPartitions{partition =>
   val startTime = Calendar.getInstance().getTimeInMillis
   result = partition.map{element =>
      [...]
   }
   val endTime = Calendar.getInstance().getTimeInMillis
   logger.info("Partition time "+(startTime-endTime)+ "ms")
   result
}

问题是它在开始执行映射操作之前立即记录“分区时间”,所以我总是获得2毫秒的时间。

我通过查看 Spark Web UI 注意到了这一点,在日志文件中,关于执行时间的行在任务开始后立即出现,而不是像预期的那样在结束时出现。

有人能解释一下为什么吗? mapPartitions里面的代码应该是线性执行的,还是我错了?

谢谢

问候 卢卡

【问题讨论】:

  • 转换被延迟评估。
  • 好的,谢谢!我解决了在 endTime 之前放置“result.size”的问题。我认为默认情况下,mapPartitions 中的地图作为 Scala 操作并不懒惰。
  • @philantrovert 不,这不是原因,mapPartitions 中的地图不是 Spark 转换,这是纯 scala 相关的
  • @RaphaelRoth 我明白了,谢谢!
  • 如果对您有帮助,请您接受我的回答

标签: scala apache-spark


【解决方案1】:

mapPartitions 内部的partitions 是一个Iterator[Row],而Iterator 在Scala 中是惰性求值的(即当使用迭代器时)。这与 Spark 的懒惰评估无关!

调用partitions.size 将触发对映射的评估,但会消耗迭代器(因为它只能迭代一次)。一个例子

val it = Iterator(1,2,3)
it.size // 3
it.isEmpty // true

你可以做的是将 Iterator 转换为非惰性集合类型:

rdd.mapPartitions{partition =>
   val startTime = Calendar.getInstance().getTimeInMillis
   result = partition.map{element =>
      [...]
   }.toVector // now the statements are evaluated
   val endTime = Calendar.getInstance().getTimeInMillis
   logger.info("Partition time "+(startTime-endTime)+ "ms")
   result.toIterator
}

编辑:请注意,您可以使用System.currentTimeMillis()(甚至System.nanoTime())而不是Calendar。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-09-06
    • 2019-01-07
    • 2013-04-14
    • 1970-01-01
    • 1970-01-01
    • 2016-12-14
    • 2021-12-25
    相关资源
    最近更新 更多