【发布时间】:2017-09-18 11:52:50
【问题描述】:
在模块 3 的那门课程 - 动手实验中...有一个示例(Spark 基础知识 1),我用来学习 Scala 和 Spark。
我尝试修改 Streaming 部分,以便在流媒体进入时计算移动平均线。我还没有弄清楚如何去做,但现在我面临着我不知道如何去做的问题更改数据类型。
import org.apache.log4j.Logger
import org.apache.log4j.Level
Logger.getLogger("org").setLevel(Level.OFF)
Logger.getLogger("akka").setLevel(Level.OFF)
import org.apache.spark._
import org.apache.spark.streaming._
import org.apache.spark.streaming.StreamingContext._
val ssc = new StreamingContext(sc,Seconds(1))
val lines = ssc.socketTextStream("localhost",7777)
import scala.collection.mutable.Queue
var ints = Queue[Double]()
def movingAverage(values: Queue[Double], period: Int): List[Double] = {
val first = (values take period).sum / period
val subtract = values map (_ / period)
val add = subtract drop period
val addAndSubtract = add zip subtract map Function.tupled(_ - _)
val res = (addAndSubtract.foldLeft(first :: List.fill(period - 1)(0.0)) {
(acc, add) => (add + acc.head) :: acc
}).reverse
res
}
val pass = lines.map(_.split(",")).
map(pass=>(pass(7).toDouble))
pass.getClass
类 org.apache.spark.streaming.dstream.MappedDStream
ints ++= List(pass).to[Queue]
名称:编译错误
消息:控制台:41:错误:类型不匹配;
找到:scala.collection.mutable.Queue[org.apache.spark.streaming.dstream.DStream[Double]]
必需:scala.collection.TraversableOnce[Double]
ints ++= List(pass).to[Queue] ^堆栈跟踪:
al pass2 = movingAverage(ints,2)
pass2.print()
ints.dequeue
ssc.start()
ssc.awaitTermination()
如何将流数据从传递到整数作为双精度队列获取?
【问题讨论】:
标签: scala apache-spark spark-streaming