call of distinct and map together throws NPE in spark library

apache-spark, nullpointerexception, scala

Solution

Spark does not support nested RDDs or user-defined functions that refer to other RDDs, hence the NullPointerException; see this thread on the `spark-users` mailing list.

It looks like your current code is trying to group the elements of `d` by value; you can do this efficiently with the `groupBy()` RDD method:

scala> val d = sc.parallelize(Seq("Hello", "World", "Hello"))
d: spark.RDD[java.lang.String] = spark.ParallelCollection@55c0c66a

scala> d.groupBy(x => x).collect()
res6: Array[(java.lang.String, Seq[java.lang.String])] = Array((World,ArrayBuffer(World)), (Hello,ArrayBuffer(Hello, Hello)))

Problem

I am unsure if this is a bug, so if you do something like this ``` // d:spark.RDD[String] d.distinct().map(x => d.filter(_.equals(x))) ``` you will get a Java NPE. However if you do a `collect` immediately after `distinct`, all will be fine. I am using spark 0.6.1.

Original source