Lambda表达式是Java8最重要的新特性,基础的内容这里就不说了,让我们从收集器开始。
什么是收集器
就是用于收集流运算后结果的角色。例如:
List<String> collect = list.stream().map(TestBean::getName).collect(Collectors.toList());
以上代码最后的Collectors.toList(),就是Java8原生的收集器,用于把流结果放到一个List中。
原生的收集器还有很多,大多定义在Collectors下的静态方法。
通过收集器的组合,能产生很复杂的收集效果,这里就不展开说了,有兴趣的可以去翻翻。
这里要说的是在原生收集器用的不趁手的时候,如何自定义一个收集器。
其实蛮简单的。
实现Collector接口就可以了。
拿一个例子来说吧,现在要定义把一个整数流中所有数求和的收集器。
比如对于流:1,2,3,4,5....
计算出来就是:1+2+3+....
好,现在来定义一个类,实现Collector接口:
Collector接口
1public class LongSumCollector implements 2 Collector<Long, LongSumCollector, Long> { 3 ... 4}
Collector接口有三个范型定义:
第一个,定义流的数据类;
第二个,定义保存中间计算结果的类,这里我们写的是收集器自身,也可以定义其他计算类;
第三个,定义流收集器的最后结果类型,就是计算出来的最后结果是个什么类型;
Collector接口有多个方法要实现,结构如下:
1public class LongSumCollector implements 2 Collector<Long, LongSumCollector, Long>{ 3 4 @Override 5 public Supplier<LongSumCollector > supplier() { 6 return null; 7 } 8 9 @Override 10 public BiConsumer<LongSumCollector, Long> accumulator() { 11 return null; 12 } 13 14 @Override 15 public BinaryOperator<LongSumCollector> combiner() { 16 return null; 17 } 18 19 @Override 20 public Function<LongSumCollector, Long> finisher() { 21 return null; 22 } 23 24 @Override 25 public Set<java.util.stream.Collector.Characteristics> characteristics() { 26 return null; 27 } 28}
下面来挨个说明一下:
supplier()方法
对计算器的初始化,就是接口范型中间那个定义的计算器类型,方法返回一个创建计算器的lambda:
1 @Override 2 public Supplier<LongSumCollector> supplier() { 3 return () -> new LongSumCollector(); 4 }
accumulator()方法
计算器对流中数据的处理,方法返回的也是lambda,第一个参数为计算器,第二个参数为流程中的数据:
1 @Override 2 public BiConsumer<LongSumCollector, Long> accumulator() { 3 return (collector, value) -> { 4 collector.sum += value; 5 }; 6 }
代码中的sum用来记录求合结果
combiner()方法
用于在并发计算的情况下,对各路并发计算的结果进行合并,方法返回的lambda,两个参数就是进行合并的两路计算器,lambda要求最后返回合并的结果。
1 @Override 2 public BinaryOperator<LongSumCollector> combiner() { 3 return (c1, c2) -> { 4 c1.sum += c2.sum; 5 return c1; 6 }; 7 }
finisher()方法
输出最后的结果:
1 @Override 2 public Function<LongSumCollector, Long> finisher() { 3 return (c1) -> c1.sum; 4 }
characteristics() 方法
用于标记收集器的特性,这里先返回一个空Set。
1 @Override 2 public Set<java.util.stream.Collector.Characteristics> characteristics() { 3 return new HashSet<>(); 4 }
好,现在收集器已经自定义完了,怎么用呢?
这里写一个测试代码
1// 数据流列表 2List<Long> list = new ArrayList<Long>(100 * 10000); 3 4// 生成100W个数 5LongStream.range(1, 100 * 10000).forEach((value) -> list.add(value)); 6 7// 用流进行数据计算 8Long sum = list.stream().collect(new LongSumCollector());
最后collect方法返回的值就是求和的结果,这种方式是不是很优雅?
并行
=====
在流程计算下,并行的实现简直是有够简单粗暴。
只需要用parallelStream()方法返回的流进行的计算就是并行的:
Long sum = list.parallelStream().collect(new LongSumCollector());
为了展式并行和串行的差别,加以时间计算和结果输出进行对比:
1// 数据流列表 2List<Long> dataList = new ArrayList<Long>(1000 * 10000); 3 4// 生成1000W个数 5LongStream.range(1, 1000 * 10000).forEach( 6 (value) -> dataList.add(value)); 7System.out.println("start"); 8// 用串行流进行数据计算 9long time1 = System.currentTimeMillis(); 10Long sum1 = dataList.stream().sequential() 11 .collect(new LongSumCollector()); 12time1 = System.currentTimeMillis() - time1; 13System.out.println("stream :" + time1 + ":" + sum1); 14 15// 用并行流进行数据计算 16long time2 = System.currentTimeMillis(); 17Long sum2 = dataList.parallelStream().collect(new LongSumCollector()); 18time2 = System.currentTimeMillis() - time2; 19System.out.println("parallelStream:" + time2 + ":" + sum2); 20System.out.println(time2 / (time1 * 1.0));
打印的结果如下:
1stream :68:49999995000000 2parallelStream:121:49999995000000 31.7794117647058822
什么情况,并行居然是串行的近2倍的时间!
好吧,测试用机只有两核,所以并行的效率发挥不出来。