【发布时间】:2021-09-30 05:58:49
【问题描述】:
完全错误是:
org.apache.spark.SparkException:这个 RDD 缺少 SparkContext。它 可能发生在以下情况: (1) RDD 转换和动作不是由驱动程序调用的,而是在其他转换内部调用的;例如,rdd1.map(x => rdd2.values.count() * x) 无效,因为值转换 并且计数操作不能在 rdd1.map 内部执行 转型。有关详细信息,请参阅 SPARK-5063。 (2) 当 Spark Streaming 作业从检查点恢复时,如果对未定义的 RDD 的引用会触发此异常 流作业用于 DStream 操作。有关详细信息,请参阅 火花 13758。
但我认为我没有在我的代码中使用嵌套的 rdd 转换。
如何解决?
我的 scala 代码:
stream.foreachRDD { rdd => {
val nRDD = rdd.map(item => item.value())
val oldRDD = sc.textFile("hdfs://localhost:9011/recData/miniApp/mall")
val top = oldRDD.sortBy(item => {
val arr = item.split(' ')
arr(0)
}, ascending = false).take(200)
val topRDD = sc.makeRDD(top)
val unionRDD = topRDD.union(nRDD)
val validRDD = unionRDD.map(item => {
val arr = item.split(' ')
((arr(1), arr(2)), arr(3).toDouble)
})
.reduceByKey((f, s) => {
if (f > s) f else s
})
.distinct()
val ratings = validRDD.map(item => {
Rating(item._1._2.toInt, item._1._1.toInt, item._2)
})
val rank = 10
val numIterations = 5
val model = ALS.train(ratings, rank, numIterations, 0.01)
nRDD.map(item => {
val arr = item.split(' ')
arr(2)
}).toDS()
.distinct()
.foreach(item=>{
println("als recommending for user "+item)
val recommendRes = model.recommendProducts(item.toInt, 10)
for (elem <- recommendRes) {
println(elem)
}
})
nRDD.saveAsTextFile("hdfs://localhost:9011/recData/miniApp/mall")
}
}
【问题讨论】:
标签: scala apache-spark