【发布时间】:2016-03-22 07:30:01
【问题描述】:
我遇到了一个问题,我需要在加入之前转换从 spark 读取的两个流。
一旦我做了转换,我就不能再加入了,我猜类型不再是 DStream[(String, String)] 而是 DStream[Map[String, String]]
val windowStream1 = act1Stream.window(Seconds(5)).transform{rdd => rdd.map(_._2).map(l =>(...toMap)}
val windowStream2 = act2Stream.window(Seconds(5)).transform{rdd => rdd.map(_._2).map(l =>(...toMap)}
val joinedWindow = windowStream1.join(windowStream2) //can't join
有什么想法吗?
【问题讨论】:
标签: scala apache-spark spark-streaming scala-collections