【问题标题】:Scala RDD groupby count along with all columnsScala RDD groupby 计数以及所有列
【发布时间】:2017-05-28 12:26:54
【问题描述】:

我需要获取所有列以及计数。在 Scala RDD 中。

Col1 col2  col3 col4
us    A     Q1   10
us    A      Q3   10
us    A      Q2   20
us    B      Q4   10
us    B      Q5   20
uk    A      Q1   10
uk    A      Q3   10
uk    A      Q2   20
uk    B      Q4   10
uk    B      Q5   20

我想要这样的结果:

Col1    col2       col3     col4     count
us         A           Q1       10          3
us         A           Q3      10          3
us         A           Q3      10          3
us         B           Q4      10          2
us         B           Q5      20          2
uk         A           Q1       10          3
uk         A           Q3      10          3
uk         A           Q3      10          3
uk         B           Q4      10          2
uk         B           Q5      20          2

这类似于 col1、col2 的 group by 并获取计数。现在我需要 col13,col4。

我正在尝试 SCALA RDD,例如:

val Top_RDD_1 = RDD.groupBy(f=> ( f._1,f._2 )).mapValues(_.toList)

这会产生

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

只有 (col1,col2),列表 (col1,col2,col3,col14) 结果如 (us,A) List((us,a,Q1,10),(us,a,Q3,10),(us,a,Q2,20)).,,,

如何获取列表计数并访问列表值。

请帮我激发 SCALA RDD 代码。

谢谢 巴拉吉。

【问题讨论】:

    标签: scala apache-spark rdd scala-collections


    【解决方案1】:

    我看不到在 RDD 的一次“扫描”中执行此操作的方法 - 您必须使用 reduceByKey 计算计数,然后将 join 计算到原始 RDD。为了有效地做到这一点(不会导致重新计算输入),您最好在加入之前cache/persist 输入:

    val keyed: RDD[((String, String), (String, String, String, Int))] = input
      .keyBy { case (c1, c2, _, _) => (c1, c2) }
      .cache()
    
    val counts: RDD[((String, String), Int)] = keyed.mapValues(_ => 1).reduceByKey(_ + _)
    
    val result = keyed.join(counts).values.map {
      case ((c1, c2, c3, c4), count) => (c1, c2, c3, c4, count)
    } 
    

    【讨论】:

      【解决方案2】:

      这是python代码:

      sales = [["US","A","Q1", 10], ["US","A","Q2", 20], ["US","B","Q3", 10], ["UK","A","Q1", 10], ["UK","A","Q2", 20], ["UK","B","Q3", 10]]  -- Sample RDD Data
      
      def func(data):
          ldata = list(data)                  # converting iterator class to list
          size = len(ldata)                   # count(*) of the list
          return [i + [size] for i in ldata]  # adding count(*) to the list
      
      sales_count = sales.groupBy( lambda w: (w[0], w[1])).mapValues(func)
      # Result: [(('US', 'A'), [['US', 'A', 'Q1', 10, 2], ['US', 'A', 'Q2', 20, 2]]), (('US', 'B'), [['US', 'B', 'Q3', 10, 1]]), (('UK', 'A'), [['UK', 'A', 'Q1', 10, 2], ['UK', 'A', 'Q2', 20, 2]]), (('UK', 'B'), [['UK', 'B', 'Q3', 10, 1]])]
      
      finalResult = sales_count.flatMap(lambda res: res[1])
      # Result:  [['US', 'A', 'Q1', 10, 2], ['US', 'A', 'Q2', 20, 2], ['US', 'B', 'Q3', 10, 1], ['UK', 'A', 'Q1', 10, 2], ['UK', 'A', 'Q2', 20, 2], ['UK', 'B', 'Q3', 10, 1]]
      
      # Both the above operations can be combined to one statement
      finalResult = sales.groupBy( lambda w: (w[0], w[1])).mapValues(func).flatMap(lambda res: res[1])
      

      注意:自定义函数真的很有帮助,就像我做的那样。您可以轻松地将相同的代码转换为 scala 代码

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2018-06-03
        • 2021-12-03
        • 1970-01-01
        • 1970-01-01
        • 2020-01-19
        • 1970-01-01
        相关资源
        最近更新 更多