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