aggregate 函数则把我们从返回值类型必须与所操作的RDD类型相同的限制中解放出来。使用 aggregate() 时
需要提供我们期待返回的类型的初始值
然后 通过一个函数把 RDD 中的元素合并起来放入累加器
考虑到每个节点是在本地进行累加 的,最终,还需要提供第二个函数来将累加器两两合并。
//用aggregate()来计算RDD的平均值
public class Operation {
public static void main(String[] args) throws InterruptedException {
// TODO 自动生成的方法存根
SparkConf conf = new SparkConf().setMaster("local").setAppName("My App");
JavaSparkContext jsc = new JavaSparkContext(conf);
JavaRDD<Integer> lines = jsc.parallelize(Arrays.asList(1,2,3,4));
JavaRDD<Integer> line = jsc.parallelize(Arrays.asList(4,5,6));
AvgCount a =new AvgCount(0,0);
Function2<AvgCount,Integer,AvgCount> addAndCount = new Function2<AvgCount,Integer,AvgCount>(){
private static final long serialVersionUID = 1L; @Override
public AvgCount call(AvgCount arg0, Integer arg1) throws Exception {
// TODO 自动生成的方法存根
arg0.total += arg1;
arg0.num += 1;
return arg0;
}
};
Function2<AvgCount,AvgCount,AvgCount> conbine = new Function2<AvgCount,AvgCount,AvgCount>(){
private static final long serialVersionUID = 1L;
@Override
public AvgCount call(AvgCount arg0, AvgCount arg1) throws Exception {
// TODO 自动生成的方法存根
arg0.total += arg1.total;
arg0.num += arg1.num;
return arg0;
}
};
line.aggregate(a,(x,y)->{
x.total += y;
x.num += 1;
return x;
},
(x,y)->{
x.total +=y.total;
x.num +=y.num;
return x;}
);
AvgCount sum = line.aggregate(a, addAndCount, conbine);
System.out.println( sum.total+":"+sum.num+"--------avg:"+(sum.total/sum.num));
jsc.close();
}
}
class AvgCount implements Serializable{
public int total;
public int num;
private static final long serialVersionUID = 3325529460700487293L;
public AvgCount(int total,int num){
this.total = total;
this.num = num;
}
}
本文来自电脑杂谈,转载请注明本文网址:
http://www.pc-fly.com/a/sanxing/article-57960-4.html