经验首页 前端设计 程序设计 Java相关 移动开发 数据库/运维 软件/图像 大数据/云计算 其他经验
当前位置:技术经验 » Java相关 » Java » 查看文章
关于Java8 parallelStream并发安全的深入讲解
来源:jb51  时间:2018/11/1 9:43:53  对本文有异议

背景

Java8的stream接口极大地减少了for循环写法的复杂性,stream提供了map/reduce/collect等一系列聚合接口,还支持并发操作:parallelStream。

在爬虫开发过程中,经常会遇到遍历一个很大的集合做重复的操作,这时候如果使用串行执行会相当耗时,因此一般会采用多线程来提速。Java8的paralleStream用fork/join框架提供了并发执行能力。但是如果使用不当,很容易陷入误区。

Java8的paralleStream是线程安全的吗

一个简单的例子,在下面的代码中采用stream的forEach接口对1-10000进行遍历,分别插入到3个ArrayList中。其中对第一个list的插入采用串行遍历,第二个使用paralleStream,第三个使用paralleStream的同时用ReentryLock对插入列表操作进行同步:

  1. private static List<Integer> list1 = new ArrayList<>();
  2. private static List<Integer> list2 = new ArrayList<>();
  3. private static List<Integer> list3 = new ArrayList<>();
  4. private static Lock lock = new ReentrantLock();
  5.  
  6. public static void main(String[] args) {
  7. IntStream.range(0, 10000).forEach(list1::add);
  8.  
  9. IntStream.range(0, 10000).parallel().forEach(list2::add);
  10.  
  11. IntStream.range(0, 10000).forEach(i -> {
  12. lock.lock();
  13. try {
  14. list3.add(i);
  15. }finally {
  16. lock.unlock();
  17. }
  18. });
  19.  
  20. System.out.println("串行执行的大小:" + list1.size());
  21. System.out.println("并行执行的大小:" + list2.size());
  22. System.out.println("加锁并行执行的大小:" + list3.size());
  23. }

执行结果:

串行执行的大小:10000
并行执行的大小:9595
加锁并行执行的大小:10000

并且每次的结果中并行执行的大小不一致,而串行和加锁后的结果一直都是正确结果。显而易见,stream.parallel.forEach()中执行的操作并非线程安全。

那么既然paralleStream不是线程安全的,是不是在其中的进行的非原子操作都要加锁呢?我在stackOverflow上找到了答案:

  • https://codereview.stackexchange.com/questions/60401/using-java-8-parallel-streams
  • https://stackoverflow.com/questions/22350288/parallel-streams-collectors-and-thread-safety

在上面两个问题的解答中,证实paralleStream的forEach接口确实不能保证同步,同时也提出了解决方案:使用collect和reduce接口。

  • http://docs.oracle.com/javase/tutorial/collections/streams/parallelism.html

在Javadoc中也对stream的并发操作进行了相关介绍:

The Collections Framework provides synchronization wrappers, which add automatic synchronization to an arbitrary collection, making it thread-safe.

Collections框架提供了同步的包装,使得其中的操作线程安全。

所以下一步,来看看collect接口如何使用。

stream的collect接口

闲话不多说直接上源码吧,Stream.java中的collect方法句柄:

  1. <R, A> R collect(Collector<? super T, A, R> collector);

在该实现方法中,参数是一个Collector对象,可以使用Collectors类的静态方法构造Collector对象,比如Collectors.toList(),toSet(),toMap(),etc,这块很容易查到API故不细说了。

除此之外,我们如果要在collect接口中做更多的事,就需要自定义实现Collector接口,需要实现以下方法:

  1. Supplier<A> supplier();
  2. BiConsumer<A, T> accumulator();
  3. BinaryOperator<A> combiner();
  4. Function<A, R> finisher();
  5. Set<Characteristics> characteristics();

要轻松理解这三个参数,要先知道fork/join是怎么运转的,一图以蔽之:

上图来自:http://www.infoq.com/cn/articles/fork-join-introduction

简单地说就是大任务拆分成小任务,分别用不同线程去完成,然后把结果合并后返回。所以第一步是拆分,第二步是分开运算,第三步是合并。这三个步骤分别对应的就是Collector的supplier,accumulator和combiner。talk is cheap show me the code,下面用一个例子来说明:

输入是一个10个整型数字的ArrayList,通过计算转换成double类型的Set,首先定义一个计算组件:

Compute.java:

  1. public class Compute {
  2. public Double compute(int num) {
  3. return (double) (2 * num);
  4. }
  5. }

接下来在Main.java中定义输入的类型为ArrayList的nums和类型为Set的输出结果result:

  1. private List<Integer> nums = new ArrayList<>();
  2. private Set<Double> result = new HashSet<>();

定义转换list的run方法,实现Collector接口,调用内部类Container中的方法,其中characteristics()方法返回空set即可:

  1. public void run() {
  2. // 填充原始数据,nums中填充0-9 10个数
  3. IntStream.range(0, 10).forEach(nums::add);
  4. //实现Collector接口
  5. result = nums.stream().parallel().collect(new Collector<Integer, Container, Set<Double>>() {
  6.  
  7. @Override
  8. public Supplier<Container> supplier() {
  9. return Container::new;
  10. }
  11.  
  12. @Override
  13. public BiConsumer<Container, Integer> accumulator() {
  14. return Container::accumulate;
  15. }
  16.  
  17. @Override
  18. public BinaryOperator<Container> combiner() {
  19. return Container::combine;
  20. }
  21.  
  22. @Override
  23. public Function<Container, Set<Double>> finisher() {
  24. return Container::getResult;
  25. }
  26.  
  27. @Override
  28. public Set<Characteristics> characteristics() {
  29. // 固定写法
  30. return Collections.emptySet();
  31. }
  32. });
  33. }

构造内部类Container,该类的作用是一个存放输入的容器,定义了三个方法:

  • accumulate方法对输入数据进行处理并存入本地的结果
  • combine方法将其他容器的结果合并到本地的结果中
  • getResult方法返回本地的结果

Container.java:

  1. class Container {
  2. // 定义本地的result
  3. public Set<Double> set;
  4.  
  5. public Container() {
  6. this.set = new HashSet<>();
  7. }
  8.  
  9. public Container accumulate(int num) {
  10. this.set.add(compute.compute(num));
  11. return this;
  12. }
  13.  
  14. public Container combine(Container container) {
  15. this.set.addAll(container.set);
  16. return this;
  17. }
  18.  
  19. public Set<Double> getResult() {
  20. return this.set;
  21. }
  22. }

在Main.java中编写测试方法:

  1. public static void main(String[] args) {
  2. Main main = new Main();
  3. main.run();
  4. System.out.println("原始数据:");
  5. main.nums.forEach(i -> System.out.print(i + " "));
  6. System.out.println("\n\ncollect方法加工后的数据:");
  7. main.result.forEach(i -> System.out.print(i + " "));
  8. }

输出:

原始数据:
0 1 2 3 4 5 6 7 8 9

collect方法加工后的数据:
0.0 2.0 4.0 8.0 16.0 18.0 10.0 6.0 12.0 14.0

我们将10个整型数值的list转成了10个double类型的set,至此验证成功~

本程序参考 http://blog.csdn.net/io_field/article/details/54971555。

一言蔽之

总结就是paralleStream里直接去修改变量是非线程安全的,但是采用collect和reduce操作就是满足线程安全的了。

总结

以上就是这篇文章的全部内容了,希望本文的内容对大家的学习或者工作具有一定的参考学习价值,如果有疑问大家可以留言交流,谢谢大家对w3xue的支持。

 友情链接:直通硅谷  点职佳  北美留学生论坛

本站QQ群:前端 618073944 | Java 606181507 | Python 626812652 | C/C++ 612253063 | 微信 634508462 | 苹果 692586424 | C#/.net 182808419 | PHP 305140648 | 运维 608723728

W3xue 的所有内容仅供测试,对任何法律问题及风险不承担任何责任。通过使用本站内容随之而来的风险与本站无关。
关于我们  |  意见建议  |  捐助我们  |  报错有奖  |  广告合作、友情链接(目前9元/月)请联系QQ:27243702 沸活量
皖ICP备17017327号-2 皖公网安备34020702000426号