一、reduceByKey作用
reduceByKey将RDD中所有K,V对中,K值相同的V进行合并,而这个合并,仅仅根据用户传入的函数来进行,下面是wordcount的例子。
1import java.util.Arrays; 2import java.util.List; 3 4import org.apache.spark.SparkConf; 5import org.apache.spark.api.java.JavaPairRDD; 6import org.apache.spark.api.java.JavaSparkContext; 7import org.apache.spark.api.java.function.Function2; 8import org.apache.spark.api.java.function.VoidFunction; 9 10import scala.Tuple2; 11 12public class WordCount { 13 14 public static void main(String[] args) { 15 SparkConf conf = new SparkConf().setAppName("spark WordCount!").setMaster("local[*]"); 16 JavaSparkContext javaSparkContext = new JavaSparkContext(conf); 17 List<Tuple2<String, Integer>> list = Arrays.asList(new Tuple2<String, Integer>("hello", 1), 18 new Tuple2<String, Integer>("word", 1), new Tuple2<String, Integer>("hello", 1), 19 new Tuple2<String, Integer>("simple", 1)); 20 JavaPairRDD<String, Integer> listRDD = javaSparkContext.parallelizePairs(list); 21 22 /** 23 * spark的shuffle是hash-based的,也就是reduceByKey算子的两个入参一个是来源于hashmap,一个来源于从map端拉取的数据,对于wordcount例子而言,进行如下运行 24 * hashMap.get(Key)+ Value,计算结果重新put回hashmap,循环往复,就迭代出了最后结果 25 */ 26 JavaPairRDD<String, Integer> wordCountPair = listRDD.reduceByKey(new Function2<Integer, Integer, Integer>() { 27 @Override 28 public Integer call(Integer v1, Integer v2) throws Exception { 29 return v1 + v2; 30 } 31 }); 32 wordCountPair.foreach(new VoidFunction<Tuple2<String, Integer>>() { 33 @Override 34 public void call(Tuple2<String, Integer> tuple) throws Exception { 35 System.out.println(tuple._1 + ":" + tuple._2); 36 } 37 }); 38 } 39 40}
计算结果:

二、reduceByKey的原理如下图

