【问题标题】:Why can't i use UpdateStateByKey in my program?为什么我不能在我的程序中使用 UpdateStateByKey?
【发布时间】:2018-12-29 02:37:49
【问题描述】:

我有一个变量 inputVectorDStream[(Int, BDM[Double]),其中 BDM 是 Breeze Matrix。我想对它使用UpdateStateByKey,但是当我尝试使用它时,我得到Cannot resolve symbol UpdateStateByKey

我是Spark 的新手,但据我所知,您必须只有key-value 对才能使用它。

我错过了什么?

我的代码是:

val ssc = new StreamingContext(conf, Seconds(3))
val lines = ssc.socketTextStream("localhost", 9999)
ssc.checkpoint("./checkpoints/")

var inputRdd = lines.map(x => x.split(","))

var arr = inputRdd.transform(x => x.groupBy(_ (1)).mapValues(x => x
                  .foldLeft(Array.ofDim[Double](C, T)) { (a, b) => {
                   var c = a
                   c(b(2).toInt)(findNextEmpty(a,b(2).toInt, T)) += b(3).toDouble
                   c  }}))

var inputVector = arr.transform(x => x.map(y=> (y._1.toInt, BDM(y._2.map(_.toArray):_*))))

var example = inputVector.updateStateByKey(somefunc)

【问题讨论】:

    标签: scala apache-spark spark-streaming stateful


    【解决方案1】:

    方法的名称是updateStateByStream,用小写的u,而不是UpdateStateByStream。 Scala 区分大小写。

    【讨论】:

    • 那是我的错,我编辑了它。我也用小写的u试过了,还是一样的错误。
    【解决方案2】:

    问题是 spark-streaming 库没有完全添加到项目的依赖项中。所以,我只是将spark-streaming.jar 添加到来自File->Project Structure->Modules 的依赖项中。

    【讨论】:

      猜你喜欢
      • 2014-05-24
      • 1970-01-01
      • 2022-01-11
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-04-28
      • 2018-06-22
      相关资源
      最近更新 更多