【发布时间】:2018-12-29 02:37:49
【问题描述】:
我有一个变量 inputVector :DStream[(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