Spark RDD实现:无需join与groupByKey计算各州城市列表及人口总和
嘿,刚接触Spark的话,用aggregateByKey或者reduceByKey确实是这类聚合场景的最优选择——既能满足你不用join和groupByKey的要求,还能保证高效的执行效率,毕竟它们会先在每个分区内做局部聚合,减少跨节点shuffle的数据量。
先给你直接上能实现需求的代码,再慢慢拆解逻辑:
用aggregateByKey实现(更灵活,适合需要自定义初始值的场景)
// 基于你给出的rdd3进行聚合 val resultRdd = rdd3.aggregateByKey((0, List.empty[String]))( // 分区内聚合:把当前累加的结果和新元素合并 (accumulator, current) => (accumulator._1 + current._2, accumulator._2 ++ current._1), // 分区间聚合:把两个分区的累加结果合并 (acc1, acc2) => (acc1._1 + acc2._1, acc1._2 ++ acc2._2) ) // 触发计算并查看结果 resultRdd.collect()
逻辑拆解:
- 初始值
(0, List.empty[String]):每个州一开始的累计人口为0,城市列表为空。 - 分区内聚合函数:
accumulator是当前分区内已经聚合好的结果(总人数,城市列表),current是当前处理的单个元素(List(城市名), 城市人口)。我们只需要把总人数加上当前城市的人口,同时把城市列表和当前的单个城市列表拼接起来。 - 分区间聚合函数:当多个分区的聚合结果需要合并时,直接把两个分区的总人数相加,城市列表拼接即可。
用reduceByKey实现(更简洁,适合不需要自定义初始值的场景)
如果你的数据里每个州至少有一个城市(不会出现空key的情况),用reduceByKey会更简洁:
val resultRdd = rdd3.reduceByKey( (a, b) => (a._1 ++ b._1, a._2 + b._2) ) resultRdd.collect()
这个逻辑更直接:对于同一个州的两个元素,直接把它们的城市列表拼接,人口相加,一步步合并成最终的结果。
为什么这两种方法是最优的?
和groupByKey相比,aggregateByKey和reduceByKey都会先在每个分区内完成局部聚合,只把分区内的聚合结果进行shuffle,而不是把所有相同key的原始数据都拉到一起,大大减少了网络传输的数据量,执行效率更高,这也是Spark中推荐的聚合方式。
内容的提问来源于stack exchange,提问作者Sree51
相关产品推荐
相关产品推荐

