【发布时间】:2014-05-21 20:26:25
【问题描述】:
sessionIdList 的类型为:
scala> sessionIdList
res19: org.apache.spark.rdd.RDD[String] = MappedRDD[17] at distinct at <console>:30
当我尝试运行以下代码时:
val x = sc.parallelize(List(1,2,3))
val cartesianComp = x.cartesian(x).map(x => (x))
val kDistanceNeighbourhood = sessionIdList.map(s => {
cartesianComp.filter(v => v != null)
})
kDistanceNeighbourhood.take(1)
我收到异常:
14/05/21 16:20:46 ERROR Executor: Exception in task ID 80
java.lang.NullPointerException
at org.apache.spark.rdd.RDD.filter(RDD.scala:261)
at $line94.$read$$iwC$$iwC$$iwC$$iwC$$anonfun$1.apply(<console>:38)
at $line94.$read$$iwC$$iwC$$iwC$$iwC$$anonfun$1.apply(<console>:36)
at scala.collection.Iterator$$anon$11.next(Iterator.scala:328)
at scala.collection.Iterator$$anon$10.next(Iterator.scala:312)
at scala.collection.Iterator$class.foreach(Iterator.scala:727)
但是,如果我使用:
val l = sc.parallelize(List("1","2"))
val kDistanceNeighbourhood = l.map(s => {
cartesianComp.filter(v => v != null)
})
kDistanceNeighbourhood.take(1)
那么就不显示异常了
两个代码sn-ps的区别在于第一个sn-p sessionIdList的类型是:
res19: org.apache.spark.rdd.RDD[String] = MappedRDD[17] at distinct at <console>:30
在第二个 sn-p "l" 是类型
scala> l
res13: org.apache.spark.rdd.RDD[String] = ParallelCollectionRDD[32] at parallelize at <console>:12
为什么会出现这个错误?
我是否需要将 sessionIdList 转换为 ParallelCollectionRDD 才能解决此问题?
【问题讨论】:
-
你能让你的代码独立吗?
-
@IvanVergiliev 除了填充的 ParallelCollectionRDD 之外的所有代码都包括在内以重新创建异常。我不知道如何创建一个填充的 ParallelCollectionRDD
标签: scala apache-spark