【问题标题】:Using reducedByKey instead of GroupBy使用 reduceByKey 而不是 GroupBy
【发布时间】:2018-04-12 21:58:00
【问题描述】:

对于存储为 RDD 的数据,如何使用 reduceByKey 而不是 GroupBy?

目的是按键分组,然后对值求和。

我有一个有效的 Scala 进程来查找优势比。

问题:

由于内存/磁盘问题,我们提取到脚本中的数据急剧增加并开始失败。这里的主要问题是由于“GROUP BY”而造成的大量洗牌。

样本数据:

(543040000711860,543040000839322,0,0,0,0)
(543040000711860,543040000938728,0,0,1,1)
(543040000711860,543040000984046,0,0,1,1)
(543040000711860,543040001071137,0,0,1,1)
(543040000711860,543040001121115,0,0,1,1)
(543040000711860,543040001281239,0,0,0,0)
(543040000711860,543040001332995,0,0,1,1)
(543040000711860,543040001333073,0,0,1,1)
(543040000839322,543040000938728,0,1,0,0)
(543040000839322,543040000984046,0,1,0,0)
(543040000839322,543040001071137,0,1,0,0)
(543040000839322,543040001121115,0,1,0,0)
(543040000839322,543040001281239,1,0,0,0)
(543040000839322,543040001332995,0,1,0,0)
(543040000839322,543040001333073,0,1,0,0)
(543040000938728,543040000984046,0,0,1,1)
(543040000938728,543040001071137,0,0,1,1)
(543040000938728,543040001121115,0,0,1,1)
(543040000938728,543040001281239,0,0,0,0)
(543040000938728,543040001332995,0,0,1,1)

这是转换我的数据的代码:

var groupby = flags.groupBy(item =>(item._1, item._2) )
var counted_group = groupby.map(item => (item._1, item._2.map(_._3).sum, item._2.map(_._4).sum, item._2.map(_._5).sum, item._2.map(_._6).sum))

结果:

((3900001339662,3900002247644),6,12,38,38)

((543040001332995,543040001352893),112,29,57,57)

((3900001572602,543040001071137),1,0,1,1)

((3900001640810,543040001281239),2,1,0,0)

((3900001295323,3900002247644),8,21,8,8)

我需要将其转换为“REDUCE BY KEY”,以便在将数据发送回之前减少每个分区中的数据。我使用的是 RDD,所以没有直接的方法来做 REDUCE BY。

【问题讨论】:

    标签: scala apache-spark mapreduce


    【解决方案1】:

    我想我通过使用 aggregateByKey 解决了这个问题。

    重新映射RDD以生成键值对

    val rddPair = flags.map(item => ((item._1, item._2), (item._3, item._4, item._5, item._6)))
    

    然后对结果应用aggregateByKey函数,现在每个分区返回聚合结果而不是分组结果。

    rddPair.aggregateByKey((0, 0, 0, 0))(
        (iTotal, oisubtotal) => (iTotal._1 + oisubtotal._1, iTotal._2 +  oisubtotal._2,  iTotal._3 +  oisubtotal._3,  iTotal._4 +  oisubtotal._4 ),
        (fTotal, iTotal) => (fTotal._1 + iTotal._1, fTotal._2 + iTotal._2, fTotal._3 + iTotal._3, fTotal._4 + iTotal._4)
      )
    

    【讨论】:

      【解决方案2】:

      reducyByKey 需要 RDD[(K, V)]键值对,因此您应该首先创建一个 rdd 对

      val rddPair = flags.map(item => ((item._1, item._2), (item._3, item._4, item._5, item._6)))
      

      然后你可以在上面的rddPair上使用reduceByKey作为

      rddPair.reduceByKey((x, y)=> (x._1+y._1, x._2+y._2, x._3+y._3, x._4+y._4))
      

      希望回答对你有帮助

      【讨论】:

      • 答案对@geek 没有帮助吗?
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-03
      • 1970-01-01
      • 2017-11-23
      相关资源
      最近更新 更多