【问题标题】:How to split entries of RDD[(String,List[(String,String,String,String)])]如何拆分 RDD[(String,List[(String,String,String,String)])] 的条目
【发布时间】:2016-08-22 07:43:51
【问题描述】:

我需要对具有以下格式的 rdd 进行一些处理:

RDD[(String,List[(String,String,String,String)])]

以下是来自 RDD 的示例条目:

(600,List((5,111,1,1), (15,111,1,5), (38,111,2,null))
(700,List((5,111,1,1), (35,111,1,5), (39,111,2,null))

我需要根据在列表中每个元组的第一个元素中找到的时间戳值将每个条目拆分为多个条目。每个条目应包含 20 分钟间隔内的时间戳。

例如,第一个条目应该分成2个条目:

List((5,111,1,1), (15,111,1,5))
List((38,111,2,null))

最终结果应该是RDD[(String,List[(String,String,String,String)])]:

(600,List((5,111,1,1), (15,111,1,5)))
(600,List((38,111,2,null)))
(700,List((5,111,1,1))
(700,List((35,111,1,5), (39,111,2,null))

任何提示如何做到这一点以及应用哪些功能?

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    您可以创建一个 splitList 函数,根据您想要的行为从给定记录中拆分列表(不确定我是否准确地遵循它,描述有点模棱两可),然后使用 flatMap 将每个键值记录“拆分”为多个记录:

    def doStuff() = {
      val input: RDD[(String,List[(String,String,String,String)])] = sc.parallelize(Seq(
        ("600",List(("5","111","1","1"), ("15","111","1","5"), ("38","111","2",null))),
        ("700",List(("5","111","1","1"), ("35","111","1","5"), ("39","111","2",null)))
      ))
    
      def splitList(l: List[(String,String,String,String)]): Iterable[List[(String,String,String,String)]] = {
        l.groupBy(_._1.toInt / 20).values // or any other logic
      }
    
      val result = input.flatMap { case (k, l) => splitList(l).map(sublist => (k, sublist)) }
    
      result.foreach(println)
      // prints: 
      // (600,List((38,111,2,null)))
      // (600,List((5,111,1,1), (15,111,1,5)))
      // (700,List((35,111,1,5), (39,111,2,null)))
      // (700,List((5,111,1,1)))
    }
    

    【讨论】:

    • 应用您的代码时,出现错误Tas not serializable, Caused by: java.io.NotSerializableException
    • 我的代码 sn-p 是:val rdd = sc.textFile("...") val splitted = rdd.map(line => line.split(",")) val processed = splitted.map(x=>(x(1),List((x(0),x(2),x(3),x(4))))) val separated = processed.flatMap { case (k, l) => splitList(l).map(sublist => (k, sublist)) }
    • 哦,这取决于splitList 的定义位置——如果它是不可序列化的 的方法,它会失败。最好将它放在 object 中(这意味着它不必被序列化),或者另一个方法中的定义,它应该可以工作。
    • 好的,现在它可以与object 一起使用,但只是好奇这两个选项中哪个更有效?另外,你能告诉我definition inside another method是什么意思吗(任何小例子)?
    • 效率上的差异可以忽略不计,如果它甚至存在的话。另一种方法中的定义只是将这个def 放在调用flatMap 的方法的主体内 - 更新了我的答案以显示这一点。
    猜你喜欢
    • 1970-01-01
    • 2018-07-06
    • 1970-01-01
    • 1970-01-01
    • 2021-10-16
    • 2017-02-17
    • 2021-09-28
    • 2017-08-22
    • 2015-12-11
    相关资源
    最近更新 更多