预见猿份
主题
首页面试题在线工具关于我们老苗一对一私教学员评价
实战项目
项目前置基础创新WMS项目Java微服务框架与实战云岚到家项目闪聚支付项目学成在线项目青橙电商项目JVM原理与实战调优分布式事务专题Java高频面试题MySQL从入门到精通Java数据结构与算法Java 并发编程(JUC)实战课程老苗一对一私教学员评价blog
blog
  • Java 并发编程(JUC)实战课程

    • 课程介绍
    • 01 多线程入门
    • 02 线程的常用方法
    • 03 线程池
    • 04 synchronized教程
    • 05 ReentrantLock教程
    • 06 读写锁
    • 07 JMM教程
    • 08 CAS教程
    • 09 wait&notify
    • 10 常用的并发集合
    • 11 异步任务控制
    • 12 常用并发工具







----- 到底线了 -----

第11章 异步任务控制 ​

3.1 Future ​

3.1.1 定义任务的三种方式 ​

最初学习多线程时就学习了什么是任务,什么是线程,线程用来执行任务,任务包含业务逻辑,创建一个任务并执行任务有哪些方式呢?下边进行汇总。

1、继承Thread类。

Thread类本身实现了Runnable接口,继承Thread类重写run方法可以创建一个任务,调用thread对象的start即可启动一个线程并执行任务,由于Java是单继承模式,继承Thread类将无法再继承其它类,所以不推荐使用此方式。

2、实现java.lang.Runnable接口

实现java.lang.Runnable接口定义独立的任务对象,有如下方法:

1)使用Thread类的Thread(Runnable target) 构造方法

2)使用ExecutorService的void execute(Runnable command)

3、实现java.util.concurrent.Callable接口

与Runnable接口不同的是Callable接口有返回值,Callable接口如下:

java
public interface Callable<V> {
    V call() throws Exception;
}
1
2
3
4

通过ExecutorService的Future submit(Callable task)方法提交Callable任务得到一个Future返回值,Future表示未来可能返回的结果,通过Future接口的方法可以获取任务的返回值,判断任务是否完成等操作,如下:

配图21

get()和get(long,TimeUnit)两个方法可以获取任务执行结果,两个方法都是阻塞方法,其中get()直到任务完成为止,get(long,TimeUnit)方法可以设置超时时间,如果达到超时时间任务还没有结束则抛出异常。

isDone():非阻塞方法,获取任务是否完成。

isCancelled(): 非阻塞方法,如果任务完成前被终止则返回true。

cancel(boolean)取消任务。

查询API文档如下:

配图22

3.1.2 定义具有返回值的任务 ​

下边测试如何定义一个具有返回值的任务。

定义一个带有返回值的任务需要实现Callable接口,如下:

java
@FunctionalInterface
public interface Callable<V> {
    /**
     * Computes a result, or throws an exception if unable to do so.
     *
     * @return computed result
     * @throws Exception if unable to compute a result
     */
    V call() throws Exception;
}
1
2
3
4
5
6
7
8
9
10

下边实现求1到100这几个数的和,测试代码如下:

java
package com.yjoffer.javase.thread.utils;

import com.yjoffer.javase.config.Logger;

import java.util.concurrent.*;

/**
 * 测试任务控制
 * @author 预见猿份(www.yjoffer.com)
 * @version 1.0
 **/
public class FutureTest {

    private static Logger logger = Logger.getLogger(FutureTest.class);

    //测试Future
    public static void test_future(){

        ExecutorService threadPool = Executors.newCachedThreadPool();

        //定义一个future
        Future<Long> future = threadPool.submit(() -> {
            logger.debug("开始任务");
            long sum = 0;
            for (int i = 1; i <= 100; i++) {
                sum += i;
            }
            try {
                TimeUnit.SECONDS.sleep(5);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            //返回结果
            logger.debug("返回结果");
            return sum;
        });

        //获取结果
        try {
            logger.debug("获取结果");
            Long result = future.get();
            logger.debug("result="+result);
        } catch (InterruptedException e) {
            e.printStackTrace();
        } catch (ExecutionException e) {
            e.printStackTrace();
        }

    }

    public static void main(String[] args) {
        test_future();
    }
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55

输出:

java
15:42:23.120 FINE com.yjoffer.javase.thread2.basic.FutureTest[main] - 获取结果
15:42:23.120 FINE com.yjoffer.javase.thread2.basic.FutureTest[pool-1-thread-1] - 开始任务
15:42:28.142 FINE com.yjoffer.javase.thread2.basic.FutureTest[pool-1-thread-1] - 返回结果
15:42:36.38 FINE com.yjoffer.javase.thread2.basic.FutureTest[main] - result=5050
1
2
3
4

在future.get();处打断点,get()只是获取结果不影响 任务的执行。

get(long,TimeUnit)等待指定时间后,如果任务没有执行完成则抛出异常java.util.concurrent.TimeoutException。

下边的代码执行任务至少需要5秒,get()方法等待3秒,所以最终由于get()方法的超时提前到达导致抛出TimeoutException异常,测试代码如下:

java
//测试Future,测试get方法超时
    public static void test_future2(){

        ExecutorService threadPool = Executors.newCachedThreadPool();

        //定义一个future
        Future<Long> future = threadPool.submit(() -> {
            logger.debug("开始任务");
            long sum = 0;
            for (int i = 1; i <= 100; i++) {
                sum += i;
            }
            try {
                TimeUnit.SECONDS.sleep(5);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            //返回结果
            logger.debug("返回结果");
            return sum;
        });

        //获取结果
        try {
            logger.debug("获取结果");
            Long result = future.get(3,TimeUnit.SECONDS);
            logger.debug("result="+result);
        } catch (InterruptedException e) {
            e.printStackTrace();
        } catch (ExecutionException e) {
            e.printStackTrace();
        } catch (TimeoutException e) {
            e.printStackTrace();
        }

    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36

运行程序,输出:

java
15:45:43.38 FINE com.yjoffer.javase.thread2.basic.FutureTest[main] - 获取结果
15:45:43.38 FINE com.yjoffer.javase.thread2.basic.FutureTest[pool-1-thread-1] - 开始任务
java.util.concurrent.TimeoutException
	at java.util.concurrent.FutureTask.get(FutureTask.java:205)
	at com.yjoffer.javase.thread2.basic.FutureTest.test_future2(FutureTest.java:75)
	at com.yjoffer.javase.thread2.basic.FutureTest.main(FutureTest.java:88)
15:45:48.64 FINE com.yjoffer.javase.thread2.basic.FutureTest[pool-1-thread-1] - 返回结果
1
2
3
4
5
6
7

3.1.3 异常处理 ​

如果任务在执行时出现异常,外界可以拿到异常信息吗?

测试代码如下:

java
//测试Future,测试异常处理
    public static void test_future3(){

        ExecutorService threadPool = Executors.newCachedThreadPool();

        //定义一个future
        Future<Long> future = threadPool.submit(() -> {
            logger.debug("开始任务");
            long sum = 0;
            for (int i = 1; i <= 100; i++) {
                sum += i;
            }
            int i=1/0;//制造异常
            try {
                TimeUnit.SECONDS.sleep(5);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            //返回结果
            logger.debug("返回结果");
            return sum;
        });

        //获取结果
        try {
            logger.debug("获取结果");
            Long result = future.get();
            logger.debug("result="+result);
        } catch (InterruptedException e) {
            e.printStackTrace();
        } catch (Exception e) {
            logger.debug("异常处理...");
            e.printStackTrace();
        }

    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36

运行程序,输出如下:

java
15:48:17.900 FINE com.yjoffer.javase.thread2.basic.FutureTest[main] - 获取结果
15:48:17.900 FINE com.yjoffer.javase.thread2.basic.FutureTest[pool-1-thread-1] - 开始任务
15:48:17.929 FINE com.yjoffer.javase.thread2.basic.FutureTest[main] - 异常处理...
java.util.concurrent.ExecutionException: java.lang.ArithmeticException: / by zero
	at java.util.concurrent.FutureTask.report(FutureTask.java:122)
	at java.util.concurrent.FutureTask.get(FutureTask.java:192)
	at com.yjoffer.javase.thread2.basic.FutureTest.test_future3(FutureTest.java:112)
	at com.yjoffer.javase.thread2.basic.FutureTest.main(FutureTest.java:124)
Caused by: java.lang.ArithmeticException: / by zero
	at com.yjoffer.javase.thread2.basic.FutureTest.lambda$test_future3$2(FutureTest.java:98)
	at java.util.concurrent.FutureTask.run$$$capture(FutureTask.java:266)
	at java.util.concurrent.FutureTask.run(FutureTask.java)
	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
	at java.lang.Thread.run(Thread.java:748)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15

从输出可以看出:当任务执行过程出现异常,通过捕获get()方法异常即可获取。

3.2 CompletableFuture ​

3.2.1 Future存在问题 ​

在“Future”章节学习了定义带返回值的任务,但是Future存在以下问题:

1、get()是阻塞方法,如果任务执行慢且任务量大则会造成大量线程阻塞。

2、Future与Future不支持链式编程,比如通过链式式编程将一个Future将结果传给下一个Future。

3、不支持多个Future组合处理,比如:多个Future并行执行,全部执行完成后继续下一个工作。

4、没有一个统一的异常处理方法。

3.2.2 CompletableFuture入门 ​

CompletableFuture是Java8引入,它不仅具备Future的功能,还对Future进行改进。

1、支持异步回调。

什么是异步回调?

A程序调用B程序,A程序不用以阻塞方式等待B程序返回结果,A程序调用B程序后可以去做其它工作,B程序计算出结果后通过A程序的回调接口通知A。如下图:

配图23

通过异步回调就可以解决长期等待结果造成线程长期阻塞的问题。

2、支持链式处理

支持异步回调链式编程,支持多个Future链式编程。

3、支持多Future组合处理

支持多Future并发执行,通过链式编程可设置多Future执行完成后的下一个阶段的工作内容。

4、统一异常处理

支持统一异常处理,通过链式编程设置异常处理逻辑。

CompletableFuture实现了Future接口和CompletionStage接口,查看源代码如下:

java
public class CompletableFuture<T> implements Future<T>, CompletionStage<T> {...}
1

Future接口在前边介绍过不再赘述,CompletionStage接口定义了异步执行的阶段处理,比如:一个阶段执行完成后执行下一个阶段,两个阶段都执行成功才执行另一个阶段。

前边我们学习了Future,CompletableFuture具备Future的功能,下边两个方法可以实现异步执行一个任务,并且支持任务返回值。

java
//返回一个新的CompletableFuture,由给定执行器中运行的线程异步完成,支持任务返回值。 
static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier, Executor executor) 
//返回一个新的CompletableFuture,它在运行给定操作之后由在给定执行程序中运行的任务异步完成,没有返回值
static CompletableFuture<Void> runAsync(Runnable runnable, Executor executor)
1
2
3
4
5

supplyAsync方法可以获取任务的返回值,任务对象需要实现Supplier接口,并且传入Executor。(如果不传入Executor则CompletableFuture会自己创建线程池)

runAsync方法支持获取任务的返回值,任务对象需要实现Runnable接口。

下边测试supplyAsync方法,代码如下:

java
package com.yjoffer.javase.thread.basic;

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

/**
 * 测试CompletableFuture
 * @author 预见猿份(www.yjoffer.com)
 * @version 1.0
 **/
public class CompletableFutureTest {
	private static Logger logger = Logger.getLogger(CompletableFutureTest.class);

    //入门程序,异步执行任务
    public static void test_completableFuture(){

        ExecutorService threadPool = Executors.newCachedThreadPool();
        CompletableFuture<String> completableFuture = CompletableFuture.supplyAsync(() -> {
            logger.debug("开始执行任务");
            try {
                TimeUnit.SECONDS.sleep(2);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            logger.debug("任务返回值");
            return "www.yjoffer.com";

        },threadPool);


      	try {
               //获取结果
               String result = completableFuture.get();
               logger.debug("result="+result);
           } catch (InterruptedException e) {
               e.printStackTrace();
           } catch (ExecutionException e) {
               e.printStackTrace();
           }
    }

    public static void main(String[] args) throws ExecutionException, InterruptedException {
        test_completableFuture();
    }

}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49

运行程序,输出:

java
15:34:17.995 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-1] - 开始执行任务
15:34:17.996 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[main] - end...
15:34:20.21 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-1] - 任务返回值
15:34:20.21 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[main] - result=www.yjoffer.com
1
2
3
4

3.2.3 异步回调 ​

下边两个方法可以实现任务回调:

java
回调方法由同一个线程进行回调
public <U> CompletableFuture<U> thenApply(Function<? super T,? extends U> fn)
回调方法由新线程进行回调
public <U> CompletableFuture<U> thenApplyAsync(Function<? super T,? extends U> fn)
1
2
3
4

带有Async后缀的回调方法由新的线程异步调用回调方法,不带Async后缀的回调方法由执行任务的线程执行回调方法

回调方法需要实现Function接口,此接口有一个R apply(T t);方法,接收任务的返回值作为参数,同时回调方法也有返回值。

测试代码如下:

java
    //异步回调
    public static void test_completableFuture2(){

        ExecutorService threadPool = Executors.newCachedThreadPool();
        CompletableFuture<String> completableFuture = CompletableFuture.supplyAsync(() -> {
            logger.debug("开始执行任务");
            try {
                TimeUnit.SECONDS.sleep(2);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            logger.debug("任务返回值");
            return "www.yjoffer.com";

        },threadPool).thenApplyAsync((result)->{
            logger.debug("回调执行");
            return result;
        },threadPool);


        try {
               //获取结果
               String result = completableFuture.get();
               logger.debug("result="+result);
           } catch (InterruptedException e) {
               e.printStackTrace();
           } catch (ExecutionException e) {
               e.printStackTrace();
           }

    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33

运行程序,输出:

java
15:34:48.876 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-1] - 开始执行任务

15:34:50.901 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-1] - 任务返回值
15:34:50.901 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-2] - 回调执行
15:34:50.902 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[main] - result=www.yjoffer.com
1
2
3
4
5

3.2.4 异常处理 ​

当任务执行过程出现异常也可以指定回调方法,使用下列方法即可实现:

java
public CompletableFuture<T> exceptionally(Function<Throwable,? extends T> fn)
1

测试代码如下:

java
//异常处理
    public static void test_completableFuture3(){

        ExecutorService threadPool = Executors.newCachedThreadPool();
        CompletableFuture<String> completableFuture = CompletableFuture.supplyAsync(() -> {
            logger.debug("开始执行任务");
            try {
                TimeUnit.SECONDS.sleep(2);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            int i=1/0;
            logger.debug("任务返回值");
            return "www.yjoffer.com";

        },threadPool).thenApplyAsync((result)->{
            logger.debug("回调执行");
            return result;
        },threadPool).exceptionally((ex)->{
            logger.debug("统一异常处理");
            ex.printStackTrace();
            return "异常结果";

        });


		try {
               //获取结果
               String result = completableFuture.get();
               logger.debug("result="+result);
           } catch (InterruptedException e) {
               e.printStackTrace();
           } catch (ExecutionException e) {
               e.printStackTrace();
           }
    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36

输出如下:

java
15:35:19.658 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-1] - 开始执行任务

15:35:21.687 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-1] - 统一异常处理
java.util.concurrent.CompletionException: java.lang.ArithmeticException: / by zero
	at java.util.concurrent.CompletableFuture.encodeThrowable(CompletableFuture.java:273)
	at java.util.concurrent.CompletableFuture.completeThrowable(CompletableFuture.java:280)
	at java.util.concurrent.CompletableFuture$AsyncSupply.run(CompletableFuture.java:1592)
	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
	at java.lang.Thread.run(Thread.java:748)
Caused by: java.lang.ArithmeticException: / by zero
	at com.yjoffer.javase.thread2.basic.CompletableFutureTest.lambda$test_completableFuture3$5(CompletableFutureTest.java:94)
	at java.util.concurrent.CompletableFuture$AsyncSupply.run(CompletableFuture.java:1590)
	... 3 more
15:35:21.688 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[main] - result=异常结果
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15

3.2.5 组合操作 ​

在生产中存在批量执行多任务的场景 ,比如:批量采集不同城市 的数据,为了提高任务执行效率可以对多个任务进行并发处理。

CompletableFuture提供Future组合操作,如下是其中的两个方法:

配图24

allOf: 所有任务全部完成则开始执行下一个阶段,比如:通过并发执行多任务采集数据,最后进行汇总。

anyOf: 只要有一个任务完成则开始执行下一个阶段,比如:通过测试多路径的网络速度,最先测试完成的路径即是速度最快的。

allOf与anyOf返回值类型不同:allOf的返回值类型为Void,anyOf的返回值类型为Object。anyOf可以通过回调轻松拿到任务的返回值。

通过allOf()方法可完成从北京、上海、深圳、郑州上个城市的软件平台采集数据。

代码如下:

java
public static void test_allOf() throws ExecutionException, InterruptedException {

        ExecutorService threadPool = Executors.newCachedThreadPool();
        //配置各城市平台的地址
        String[] platforms = new String[]{
                "北京平台",
                "上海平台",
                "深圳平台",
                "郑州平台"
        };
        //批量定义CompletableFuture
        List<CompletableFuture<String>> collect = Arrays.asList(platforms).stream().map(v -> {

            return CompletableFuture.supplyAsync(() -> {
                //请求平台获取数据
                logger.debug("请求" + v + "获取数据");
                try {
                    TimeUnit.SECONDS.sleep(3);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                return v + "数据";
            },threadPool);

        }).collect(Collectors.toList());
        //将List转成数组
        CompletableFuture[] completableFutures = collect.toArray(new CompletableFuture[collect.size()]);
        //并发执行多任务
        CompletableFuture<Void> voidCompletableFuture = CompletableFuture.allOf(completableFutures);

        //使用回调收集采集的数据,allOf中的所有任务执行完成才会执行thenApplyAsync中的内容
        CompletableFuture<List<String>> listCompletableFuture = voidCompletableFuture.thenApplyAsync(result -> {
            return collect.stream().map(future -> {
                String s = null;
                try {
                    s = future.get();
                    logger.debug("收集"+s);

                } catch (InterruptedException e) {
                    e.printStackTrace();
                } catch (ExecutionException e) {
                    e.printStackTrace();
                }
                return s;
            }).collect(Collectors.toList());
        },threadPool);
        //获取采集数据,触发任务执行
        List<String> strings = listCompletableFuture.get();
       logger.debug("收集数据结果:"+strings);

    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51

输出 :

java
09:46:37.441 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-3] - 请求深圳平台获取数据
09:46:37.441 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-4] - 请求郑州平台获取数据
09:46:37.441 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-1] - 请求北京平台获取数据
09:46:37.441 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-2] - 请求上海平台获取数据
09:46:40.466 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-5] - 收集北京平台数据
09:46:40.466 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-5] - 收集上海平台数据
09:46:40.467 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-5] - 收集深圳平台数据
09:46:40.467 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-5] - 收集郑州平台数据
09:46:40.468 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[main] - 收集数据结果:[北京平台数据, 上海平台数据, 深圳平台数据, 郑州平台数据]
1
2
3
4
5
6
7
8
9

下边测试anyOf方法,人为控制“郑州平台”的速度最快,预期先返回郑州平台的数据。代码如下:

java
    public static void test_anyOf() throws ExecutionException, InterruptedException {
        ExecutorService threadPool = Executors.newCachedThreadPool();
        //配置各城市平台的地址
        String[] platforms = new String[]{
                "北京平台",
                "上海平台",
                "深圳平台",
                "郑州平台"
        };
        //批量定义CompletableFuture
        List<CompletableFuture<String>> collect = Arrays.asList(platforms).stream().map(v -> {
            return CompletableFuture.supplyAsync(() -> {
                //请求平台获取数据
                logger.debug("请求" + v + "获取数据");
                try {
                    if(v.equals("郑州平台")){//郑州平台速度最快
                        TimeUnit.SECONDS.sleep(1);
                    }else{
                        TimeUnit.SECONDS.sleep(3);
                    }

                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                return v + "数据";
            },threadPool);

        }).collect(Collectors.toList());
        //将List转成数组
        CompletableFuture[] completableFutures = collect.toArray(new CompletableFuture[collect.size()]);

        CompletableFuture<Object> objectCompletableFuture = CompletableFuture.anyOf(completableFutures).thenApplyAsync((result) -> {
            logger.debug("采集到数据" + result);
            return result;
        },threadPool);

        //获取结果
        Object result = objectCompletableFuture.get();
        logger.debug("result="+result);
    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40

输出 :

java
16:34:26.317 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-1] - 请求北京平台获取数据
16:34:26.317 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-2] - 请求上海平台获取数据
16:34:26.317 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-3] - 请求深证平台获取数据
16:34:26.317 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-4] - 请求郑州平台获取数据
16:34:27.338 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-5] - 采集到数据郑州平台数据
16:34:27.339 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[main] - result=郑州平台数据
1
2
3
4
5
6

3.3 CompletionService ​

3.3.1 批量执行任务的问题 ​

在“CompletableFuture”章节allOf方法批量执行任务,这种方式必须等到所有任务执行完毕才可以去执行回调方法,如果遇到执行时间特别长的任务将影响最终结果的处理。

运行下边的测试代码:

java
public static void test_allOf() throws ExecutionException, InterruptedException {

        ExecutorService threadPool = Executors.newCachedThreadPool();
        //配置各城市平台的地址
        String[] platforms = new String[]{
                "北京平台",
                "上海平台",
                "深圳平台",
                "郑州平台"
        };
        //批量定义CompletableFuture
        AtomicInteger atomicInteger = new AtomicInteger(5);
        List<CompletableFuture<String>> collect = Arrays.asList(platforms).stream().map(v -> {

            return CompletableFuture.supplyAsync(() -> {
                //请求平台获取数据
                logger.debug("请求" + v + "获取数据");
                try {
                    TimeUnit.SECONDS.sleep(atomicInteger.getAndDecrement());
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                return v + "数据";
            },threadPool);

        }).collect(Collectors.toList());
        //将List转成数组
        CompletableFuture[] completableFutures = collect.toArray(new CompletableFuture[collect.size()]);
        //并发执行多任务
        CompletableFuture<Void> voidCompletableFuture = CompletableFuture.allOf(completableFutures);
        //使用回调收集采集的数据
        CompletableFuture<List<String>> listCompletableFuture = voidCompletableFuture.thenApplyAsync(result -> {
            return collect.stream().map(future -> {
                String s = null;
                try {
                    s = future.get();
                    logger.debug("收集"+s);

                } catch (InterruptedException e) {
                    e.printStackTrace();
                } catch (ExecutionException e) {
                    e.printStackTrace();
                }
                return s;
            }).collect(Collectors.toList());
        },threadPool);
        //获取采集数据,触发任务
        List<String> strings = listCompletableFuture.get();

        logger.debug("收集数据结果:"+strings);
    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51

4个任务的处理时间最短是1秒,最长是5秒,收集任务结果的代码是在所有任务全部完成后才去执行。

输出 :

java
11:34:33.857 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-4] - 请求郑州平台获取数据
11:34:33.857 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-2] - 请求上海平台获取数据
11:34:33.857 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-3] - 请求深圳平台获取数据
11:34:33.857 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-1] - 请求北京平台获取数据
11:34:38.881 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-2] - 收集北京平台数据
11:34:38.882 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-2] - 收集上海平台数据
11:34:38.882 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-2] - 收集深圳平台数据
11:34:38.882 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-2] - 收集郑州平台数据
11:34:38.883 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[main] - 收集数据结果:[北京平台数据, 上海平台数据, 深圳平台数据, 郑州平台数据]
1
2
3
4
5
6
7
8
9

3.3.2 CompletionService入门 ​

针对上边的问题可以用CompletionService来处理,CompletionService的内部维护了一个阻塞队列,它会将执行完毕的Future放入队列由消费者来获取。

CompletionService 接口的实现类是 ExecutorCompletionService,构造方法如下:

java
ExecutorCompletionService(Executor executor) 
使用提供的执行程序创建一个ExecutorCompletionService来执行基本任务,一个LinkedBlockingQueue作为完成队列。  
ExecutorCompletionService(Executor executor, BlockingQueue<Future<V>> completionQueue) 
使用提供的执行程序创建一个ExecutorCompletionService,用于执行基本任务,并将提供的队列作为其完成队列。
1
2
3
4
5

如果不指定completionQueue默认使用LinkedBlockingQueue。

下边用CompletionService实现批量任务的处理。

java
public static void test_CompletionService() throws ExecutionException, InterruptedException {

        ExecutorService threadPool = Executors.newCachedThreadPool();
        //配置各城市平台的地址
        String[] platforms = new String[]{
                "北京平台",
                "上海平台",
                "深圳平台",
                "郑州平台"
        };
       CompletionService<String> completionService  = new ExecutorCompletionService<>(threadPool);
        AtomicInteger atomicInteger = new AtomicInteger(5);
        //批量定义CompletableFuture
        List<Future<String>> collect = Arrays.asList(platforms).stream().map(v -> {
            //执行任务,并返回Future
            return completionService.submit(()->{
                //请求平台获取数据
                logger.debug("请求" + v + "获取数据");
                try {
                    TimeUnit.SECONDS.sleep(atomicInteger.getAndDecrement());
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                return v + "数据";
            });
        }).collect(Collectors.toList());

        for (int i = 0; i < 4; i++) {
            Future<String> take = completionService.take();
            String s = take.get();
            logger.debug("收集数据:"+s);
        }

    }
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34

运行程序发现4个任务的结果根据完成时间先后依次返回。

输出结果如下:

java
11:38:40.334 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-3] - 请求深圳平台获取数据
11:38:40.334 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-4] - 请求郑州平台获取数据
11:38:40.334 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-1] - 请求北京平台获取数据
11:38:40.334 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[pool-1-thread-2] - 请求上海平台获取数据
11:38:42.356 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[main] - 收集数据:深圳平台数据
11:38:43.356 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[main] - 收集数据:上海平台数据
11:38:44.356 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[main] - 收集数据:北京平台数据
11:38:45.356 FINE com.yjoffer.javase.thread2.basic.CompletableFutureTest[main] - 收集数据:郑州平台数据
1
2
3
4
5
6
7
8
← 10 常用的并发集合12 常用并发工具 →








如果发现文档内容有错误或排版错乱,请及时联系站长老苗修改,不胜感激。联系我们
关于我们 | 隐私政策 | 豫ICP备2026003386号-4 | 豫公网安备41010202004008号
目录

本页无章节