【问题标题】:Spark streaming - transform two streams and joinSpark 流式传输 - 转换两个流并加入
【发布时间】: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


    【解决方案1】:

    这并不能解决您的问题,但可以使其更易于消化。您可以通过使用预期类型定义临时 val/def/var 标识符来拆分方法链并记录您在每个步骤中预期的类型。通过这种方式,您可以轻松发现类型不再符合您的期望的地方。

    例如我希望您的act1Streamact2Stream 实例的类型为DStream[(String, String)],我暂时将其称为s1s2。如果不是这样,请评论我。

    def joinedWindow(
          s1: DStream[(String, String)], 
          s2: DStream[(String, String)]
        ): DStream[...] = {
      val w1 = windowedStream(s1)
      val w2 = windowedStream(s2)
      w1.join(w2)
    }
    def windowedStream(actStream: DStream[(String, String)]): DStream[Map[...]] = {
      val windowed: DStream[(String, String)] = actStream.window(Seconds(5))
      windowed.transform( myTransform )
    }
    def myTransform(rdd: RDD[(String, String)]): RDD[Map[...]] = {
      val mapped: RDD[String] = rdd.map(_._2)
      // not enough information to conclude 
      // the result type from given code
      mapped.map(l =>(...toMap)) 
    }
    

    从那里可以通过填写... 部分来总结其余类型。逐行消除编译器错误,直到获得所需的结果。与

    的文档

    DStream[T]

    • def window(windowDuration: Duration): DStream[T]
    • def transform[U](transformFunc: (RDD[T]) ⇒ RDD[U])(implicit arg0: ClassTag[U]): DStream[U]

    PairDStreamFunctions[K,V]

    • def join[W](other: DStream[(K, W)])(implicit arg0: ClassTag[W]): DStream[(K, (V, W))]

    RDD[T]

    • def map[U](f: (T) ⇒ U)(implicit arg0: ClassTag[U]): RDD[U]

    至少通过这种方式,您可以准确地知道预期类型和生成的类型不匹配。

    【讨论】:

    • 是的,act1Streamact2Stream 属于 DStream[Map[String, String]] 类型,但让我在这里解释一下。我有两个DStream[(String, String)] 类型的不同流,但在加入之前,我需要转换每个流,并且转换方法返回一个DStream[Map[String, String]] 并且无法加入它......所以有没有办法加入两个流DStream[Map[String, string]]
    • 我不是spark专家,但似乎join操作只定义在PairDStreamFunctions[K,V]s上。签名告诉我,您需要有两个 DStream[(K,V)] 类型的实例才能应用连接操作。也许你可以以某种方式改变你的 transform 和内部 maps 以返回 DStream[(String,Map[String,String])] 实例,如果这在你的用例中有意义的话......
    猜你喜欢
    • 2018-05-27
    • 2020-07-17
    • 1970-01-01
    • 2016-08-25
    • 2016-09-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-12-12
    相关资源
    最近更新 更多