Apache Spark RDD:对每个唯一键的二元组数值求和实现问询
解决Spark RDD按键对二元组元素分别求和的问题
嘿,这个需求其实很容易实现,咱们基于你已经写好的appNamesAndPropertiesRdd,用Spark自带的键值对操作就能搞定。核心思路就是按键分组,对每个键对应的二元组的两个Long值分别累加,下面给你两种常用的实现方式:
方法一:用reduceByKey(最简洁直接)
reduceByKey是Spark处理键值对RDD分组聚合最常用的算子之一,它会自动把同一个键下的所有值两两合并,咱们只需要定义合并规则就行:
// 对每个appName对应的(totalUsageTime, usageFrequency)分别求和 val summedAppRdd = appNamesAndPropertiesRdd.reduceByKey { case ((accTime, accFreq), (currTime, currFreq)) => (accTime + currTime, accFreq + currFreq) }
代码解释:
- 匿名函数里的
(accTime, accFreq)是之前累加的结果,(currTime, currFreq)是当前遍历到的二元组 - 每次合并时,把两个二元组的第一个元素(总使用时间)相加,第二个元素(使用频率)相加,最终得到每个appName对应的总和
用你给的示例数据测试的话,结果正好是:
(com.instagram.android,(10,3))(com.android.contacts,(9,5))
方法二:用aggregateByKey(适合需要自定义初始值的场景)
如果后续你需要给累加设置初始值(比如默认从某个非0值开始累加),可以用aggregateByKey,它支持分区内和分区间的两步合并:
// 初始值设为(0L, 0L),表示初始总时间和频率都是0 val summedAppRdd = appNamesAndPropertiesRdd.aggregateByKey((0L, 0L))( // 分区内合并:把当前元素加到累加器上 (accumulator, (time, freq)) => (accumulator._1 + time, accumulator._2 + freq), // 分区间合并:把两个分区的累加结果相加 (acc1, acc2) => (acc1._1 + acc2._1, acc1._2 + acc2._2) )
这个方法和reduceByKey的结果完全一致,只是更灵活,适合复杂的聚合场景。
额外小提示
如果你的数据后续需要转成DataFrame处理,也可以用groupBy("appName").agg(sum("totalUsageTime"), sum("usageFrequency"))来实现,不过既然你现在用的是RDD,上面两种方法就足够啦~
内容的提问来源于stack exchange,提问作者Fobi
相关产品推荐
相关产品推荐

