Spark RDD操作之ReduceByKey

一、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的原理如下图

     

 

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

swap空间的增减方法

(1)增大swap空间去激活swap交换区:swapoff v /dev/vg00/lvswap扩展交换lv:lvextend L 10G /dev/vg00/lvswap重新生成swap交换区:mkswap /dev/vg00/lvswap激活新生成的交换区:swapon v /dev/vg00/lvswap