【问题标题】:org.apache.spark.SparkException: This RDD lacks a SparkContext errororg.apache.spark.SparkException:此 RDD 缺少 SparkContext 错误
【发布时间】: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


    【解决方案1】:

    错误告诉您缺少SparkContext。我猜程序在这一行失败了:

    val oldRDD = sc.textFile("hdfs://localhost:9011/recData/miniApp/mall")
    

    documentation 提供了一个创建 SparkContext 以在这种情况下使用的示例。

    来自文档:

    val stream: DStream[String] = ...
    
    stream.foreachRDD { rdd =>
    
      // Get the singleton instance of SparkSession
      val spark = SparkSession.builder.config(rdd.sparkContext.getConf).getOrCreate()
      import spark.implicits._
    
      // Do things...
    }
    

    虽然您使用的是RDDs 而不是DataFrames,但应该适用相同的原则。

    【讨论】:

      猜你喜欢
      • 2018-01-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2022-01-25
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-02-23
      相关资源
      最近更新 更多