第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接口如下:
public interface Callable<V> {
V call() throws Exception;
}通过ExecutorService的Future submit(Callable task)方法提交Callable任务得到一个Future返回值,Future表示未来可能返回的结果,通过Future接口的方法可以获取任务的返回值,判断任务是否完成等操作,如下:

get()和get(long,TimeUnit)两个方法可以获取任务执行结果,两个方法都是阻塞方法,其中get()直到任务完成为止,get(long,TimeUnit)方法可以设置超时时间,如果达到超时时间任务还没有结束则抛出异常。
isDone():非阻塞方法,获取任务是否完成。
isCancelled(): 非阻塞方法,如果任务完成前被终止则返回true。
cancel(boolean)取消任务。
查询API文档如下:

3.1.2 定义具有返回值的任务
下边测试如何定义一个具有返回值的任务。
定义一个带有返回值的任务需要实现Callable接口,如下:
@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到100这几个数的和,测试代码如下:
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();
}
}输出:
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在future.get();处打断点,get()只是获取结果不影响 任务的执行。
get(long,TimeUnit)等待指定时间后,如果任务没有执行完成则抛出异常java.util.concurrent.TimeoutException。
下边的代码执行任务至少需要5秒,get()方法等待3秒,所以最终由于get()方法的超时提前到达导致抛出TimeoutException异常,测试代码如下:
//测试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();
}
}运行程序,输出:
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] - 返回结果3.1.3 异常处理
如果任务在执行时出现异常,外界可以拿到异常信息吗?
测试代码如下:
//测试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();
}
}运行程序,输出如下:
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)从输出可以看出:当任务执行过程出现异常,通过捕获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。如下图:

通过异步回调就可以解决长期等待结果造成线程长期阻塞的问题。
2、支持链式处理
支持异步回调链式编程,支持多个Future链式编程。
3、支持多Future组合处理
支持多Future并发执行,通过链式编程可设置多Future执行完成后的下一个阶段的工作内容。
4、统一异常处理
支持统一异常处理,通过链式编程设置异常处理逻辑。
CompletableFuture实现了Future接口和CompletionStage接口,查看源代码如下:
public class CompletableFuture<T> implements Future<T>, CompletionStage<T> {...}Future接口在前边介绍过不再赘述,CompletionStage接口定义了异步执行的阶段处理,比如:一个阶段执行完成后执行下一个阶段,两个阶段都执行成功才执行另一个阶段。
前边我们学习了Future,CompletableFuture具备Future的功能,下边两个方法可以实现异步执行一个任务,并且支持任务返回值。
//返回一个新的CompletableFuture,由给定执行器中运行的线程异步完成,支持任务返回值。
static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier, Executor executor)
//返回一个新的CompletableFuture,它在运行给定操作之后由在给定执行程序中运行的任务异步完成,没有返回值
static CompletableFuture<Void> runAsync(Runnable runnable, Executor executor)supplyAsync方法可以获取任务的返回值,任务对象需要实现Supplier接口,并且传入Executor。(如果不传入Executor则CompletableFuture会自己创建线程池)
runAsync方法支持获取任务的返回值,任务对象需要实现Runnable接口。
下边测试supplyAsync方法,代码如下:
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();
}
}运行程序,输出:
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.com3.2.3 异步回调
下边两个方法可以实现任务回调:
回调方法由同一个线程进行回调
public <U> CompletableFuture<U> thenApply(Function<? super T,? extends U> fn)
回调方法由新线程进行回调
public <U> CompletableFuture<U> thenApplyAsync(Function<? super T,? extends U> fn)带有Async后缀的回调方法由新的线程异步调用回调方法,不带Async后缀的回调方法由执行任务的线程执行回调方法
回调方法需要实现Function接口,此接口有一个R apply(T t);方法,接收任务的返回值作为参数,同时回调方法也有返回值。
测试代码如下:
//异步回调
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();
}
}运行程序,输出:
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.com3.2.4 异常处理
当任务执行过程出现异常也可以指定回调方法,使用下列方法即可实现:
public CompletableFuture<T> exceptionally(Function<Throwable,? extends T> fn)测试代码如下:
//异常处理
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();
}
}输出如下:
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=异常结果3.2.5 组合操作
在生产中存在批量执行多任务的场景 ,比如:批量采集不同城市 的数据,为了提高任务执行效率可以对多个任务进行并发处理。
CompletableFuture提供Future组合操作,如下是其中的两个方法:

allOf: 所有任务全部完成则开始执行下一个阶段,比如:通过并发执行多任务采集数据,最后进行汇总。
anyOf: 只要有一个任务完成则开始执行下一个阶段,比如:通过测试多路径的网络速度,最先测试完成的路径即是速度最快的。
allOf与anyOf返回值类型不同:allOf的返回值类型为Void,anyOf的返回值类型为Object。anyOf可以通过回调轻松拿到任务的返回值。
通过allOf()方法可完成从北京、上海、深圳、郑州上个城市的软件平台采集数据。
代码如下:
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);
}输出 :
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] - 收集数据结果:[北京平台数据, 上海平台数据, 深圳平台数据, 郑州平台数据]下边测试anyOf方法,人为控制“郑州平台”的速度最快,预期先返回郑州平台的数据。代码如下:
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);
}输出 :
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=郑州平台数据3.3 CompletionService
3.3.1 批量执行任务的问题
在“CompletableFuture”章节allOf方法批量执行任务,这种方式必须等到所有任务执行完毕才可以去执行回调方法,如果遇到执行时间特别长的任务将影响最终结果的处理。
运行下边的测试代码:
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);
}4个任务的处理时间最短是1秒,最长是5秒,收集任务结果的代码是在所有任务全部完成后才去执行。
输出 :
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] - 收集数据结果:[北京平台数据, 上海平台数据, 深圳平台数据, 郑州平台数据]3.3.2 CompletionService入门
针对上边的问题可以用CompletionService来处理,CompletionService的内部维护了一个阻塞队列,它会将执行完毕的Future放入队列由消费者来获取。
CompletionService 接口的实现类是 ExecutorCompletionService,构造方法如下:
ExecutorCompletionService(Executor executor)
使用提供的执行程序创建一个ExecutorCompletionService来执行基本任务,一个LinkedBlockingQueue作为完成队列。
ExecutorCompletionService(Executor executor, BlockingQueue<Future<V>> completionQueue)
使用提供的执行程序创建一个ExecutorCompletionService,用于执行基本任务,并将提供的队列作为其完成队列。如果不指定completionQueue默认使用LinkedBlockingQueue。
下边用CompletionService实现批量任务的处理。
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);
}
}运行程序发现4个任务的结果根据完成时间先后依次返回。
输出结果如下:
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] - 收集数据:郑州平台数据