Scala aggregate function vs. Spark RDD aggregate function

Viewed 176

Here are the definitions of the function:

Scala:

aggregate[B](z: => B)(seqop: (B, A) => B, combop: (B, B) => B): B

Spark RDD:

aggregate[B](z: B)(seqop: (B, A) => B, combop: (B, B) => B): B

I know that the Scala aggregate function is designed to work on parallel collections and the Spark RDD aggregate function is designed to work on distributed collections.

But, Why the z parameter in Scala is in lazy format, while in Spark RDD is in eager format?

1 Answers

Well first of all this is a call-by-name parameter in Scala. This means that they are evaluated every time they are used, which is not the same thing as lazy, where it's evaluated only once the first time it's used and all subsequent calls use that result. (https://docs.scala-lang.org/tour/by-name-parameters.html)

So spark relies on distributed data sets which means computation can be done on multiple nodes. And I think they chose to have the zero element a call-by-value parameter (which what you referred to as 'eager'), to avoid having to recompute it on each and every node where this computation is done.

Related