如何不使用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
相关产品推荐
相关产品推荐

