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

如何不使用groupByKey/groupBy实现Scala RDD的分组操作

不使用groupBy/groupByKey实现RDD分组聚合

你可以通过Spark的reduceByKey、aggregateByKey或combineByKey算子实现需求,这些聚合类操作无需依赖groupBy系列函数,就能将相同key的value合并为Iterable集合。以下是具体实现示例:

Scala 版本实现

1. 初始化输入RDD

import org.apache.spark.SparkContext
import org.apache.spark.SparkConf

val conf = new SparkConf().setAppName("GroupWithoutGroupBy").setMaster("local")
val sc = new SparkContext(conf)

val inputRDD = sc.parallelize(Seq(
  ("Group 1", "John"),
  ("Group 2", "Sam"),
  ("Group 1", "Mary"),
  ("Group 3", "Pam")
))

方法一:map + reduceByKey

先将每个元素转为(key, List[value])格式,再通过reduceByKey合并相同key的列表:

val resultRDD = inputRDD
  .map { case (key, value) => (key, List(value)) }
  .reduceByKey((list1, list2) => list1 ++ list2)

方法二:aggregateByKey

指定初始空列表,分别定义分区内元素累加逻辑和跨分区列表合并逻辑:

val resultRDD = inputRDD.aggregateByKey(List[String]())(
  (acc, value) => acc :+ value,  // 分区内将新value加入累加列表
  (acc1, acc2) => acc1 ++ acc2   // 合并不同分区的结果列表
)

方法三:combineByKey

更底层的聚合算子,分别定义combiner创建、单元素合并、combiner合并三个逻辑:

val resultRDD = inputRDD.combineByKey(
  (value: String) => List(value),  // 将单个value转为初始列表
  (acc: List[String], value: String) => acc :+ value,  // 向列表添加新value
  (acc1: List[String], acc2: List[String]) => acc1 ++ acc2  // 合并两个列表
)

验证结果

执行以下代码输出结果:

resultRDD.collect().foreach(println)

输出符合预期:

(Group 1,List(John, Mary))
(Group 2,List(Sam))
(Group 3,List(Pam))

Python 版本实现

1. 初始化输入RDD

from pyspark import SparkContext, SparkConf

conf = SparkConf().setAppName("GroupWithoutGroupBy").setMaster("local")
sc = SparkContext(conf=conf)

input_rdd = sc.parallelize([
    ("Group 1", "John"),
    ("Group 2", "Sam"),
    ("Group 1", "Mary"),
    ("Group 3", "Pam")
])

方法一:map + reduceByKey

result_rdd = input_rdd.map(lambda x: (x[0], [x[1]])).reduceByKey(lambda a, b: a + b)

方法二:aggregateByKey

result_rdd = input_rdd.aggregateByKey([], lambda acc, val: acc + [val], lambda acc1, acc2: acc1 + acc2)

方法三:combineByKey

def create_combiner(val):
    return [val]

def merge_value(acc, val):
    acc.append(val)
    return acc

def merge_combiners(acc1, acc2):
    return acc1 + acc2

result_rdd = input_rdd.combineByKey(create_combiner, merge_value, merge_combiners)

验证结果

for item in result_rdd.collect():
    print(item)

这些方法利用Spark的Shuffle聚合机制,先在Map端完成局部聚合,再在Reduce端合并全局结果,性能甚至可能优于groupByKey(减少了Shuffle数据量),同时满足不使用groupBy系列函数的要求。

内容的提问来源于stack exchange,提问作者yxx

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 07:25:20