You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()

逻辑拆解:

  1. 初始值(0, List.empty[String]):每个州一开始的累计人口为0,城市列表为空。
  2. 分区内聚合函数:accumulator是当前分区内已经聚合好的结果(总人数,城市列表),current是当前处理的单个元素(List(城市名), 城市人口)。我们只需要把总人数加上当前城市的人口,同时把城市列表和当前的单个城市列表拼接起来。
  3. 分区间聚合函数:当多个分区的聚合结果需要合并时,直接把两个分区的总人数相加,城市列表拼接即可。

用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 07:33:33