【问题标题】:DStream all identical keys should be processed sequentiallyDStream 所有相同的键都应该按顺序处理
【发布时间】:2018-11-04 23:31:10
【问题描述】:

我有 (Key,Value) 类型的 dstream。

mapped2.foreachRDD(rdd => {
  rdd.foreachPartition(p => {
    p.foreach(x => {
    }
  )})
})

我需要确保所有具有相同键的项目都在一个分区中由一个核心处理..所以实际上是按顺序处理的..

如何做到这一点?我可以使用低效的 GroupBykey 吗?

【问题讨论】:

    标签: scala apache-spark spark-streaming


    【解决方案1】:

    你可以使用PairDStreamFunctions.combineByKey:

    import org.apache.spark.HashPartitioner
    import org.apache.spark.streaming.dstream.DStream
    /**
      * Created by Yuval.Itzchakov on 29/11/2016.
      */
    object GroupingDStream {
      def main(args: Array[String]): Unit = {
        val pairs: DStream[(String, String)] = ???
        val numberOfPartitions: Int = ???
    
        val groupedByIds: DStream[(String, List[String])] = pairs.combineByKey[List[String]](
          _ => List[String](), 
          (strings: List[String], s: String) => s +: strings, 
          (first: List[String], second: List[String]) => first ++ second, new HashPartitioner(numberOfPartitions))
    
        groupedByIds.foreachRDD(rdd => {
          rdd.foreach((kvp: (String, List[String])) => {
    
          })
        })
      }
    }
    

    combineByKey 的结果将是一个元组,其中第一个元素是键,第二个元素是值的集合。请注意,为了示例的简单性,我使用了(String, String),因为您没有提供任何类型。

    然后,使用foreach 迭代值列表并在需要时按顺序处理它们。请注意,如果您需要应用额外的转换,可以使用DStream.map 并对第二个元素(值列表)进行操作,而不是使用foreachRDD

    【讨论】:

    • 嗨,感谢您的回答..我可以使用另一个函数,例如 partionbykey (以避免键值对的低效分组)..那么我上面的代码是否可以保证值是串行执行的(即,每个执行者一个核心)?或者应该是分组。即,通过键值对中的值访问?
    • @mahdi62 为什么你认为 combineByKey 效率低下?它将在本地组合执行器内的所有相似键,并且仅通过网络对组合结果进行洗牌。
    • 该代码实际上为键、值对与一项提供了空列表...我猜组合器应该是 (x:String) =>List[String](x),
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-07-19
    • 2020-05-02
    • 1970-01-01
    • 1970-01-01
    • 2020-03-31
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多