org.apache.spark.SparkException: RDD缺失SparkContext无嵌套转换如何解决
问题定位
你代码中存在未察觉的嵌套RDD操作,符合报错的第一种场景:
- 你使用的是
org.apache.spark.mllib.recommendation.ALS(MLlib旧版API),其recommendProducts方法内部会触发RDD的转换和行动操作。 - 你将该方法放在了
toDS().distinct().foreach()分布式算子内部执行,这段逻辑会分发到Executor节点运行,而Executor上没有可用的SparkContext,直接触发RDD缺少SparkContext的异常。
另外还有两个隐藏的风险点:
- 同个
foreachRDD周期内,你既读取HDFS路径hdfs://localhost:9011/recData/miniApp/mall的内容,又直接把新数据写入同一路径,会存在读写冲突,可能出现读空、数据丢失的问题。 - 如果你的Spark Streaming作业开启了Checkpoint恢复,要避免引用Driver初始化阶段创建的全局静态RDD,所有流处理周期内用到的RDD都要在当前
foreachRDD内部创建。
解决方案
- 把推荐逻辑从分布式算子移到Driver端执行,先将待推荐的用户ID收集到Driver,再循环调用推荐方法,替换原有对应代码段即可:
// 先把所有待推荐的用户ID收集到Driver端 val userIds = nRDD.map(item => { val arr = item.split(' ') arr(2).toInt }).distinct().collect() // Driver端循环执行推荐,不会触发Executor端的RDD操作 userIds.foreach(item => { println("als recommending for user " + item) val recommendRes = model.recommendProducts(item, 10) for (elem <- recommendRes) { println(elem) } })
- 修复HDFS读写冲突:先将nRDD写入临时HDFS路径,等当前周期所有逻辑执行完成后,再用临时路径覆盖原目标路径。
- 如果开启了Checkpoint,确保
foreachRDD内部用到的所有RDD都通过当前周期可用的SparkContext创建,不要引用全局静态RDD。
内容的提问来源于stack exchange,提问作者chucklai
相关产品推荐
相关产品推荐

