在Java开发中,异步任务链指的是将多个异步任务按照特定顺序或逻辑关系组合执行,前一个任务的结果可以作为后一个任务的输入,同时支持并行执行、异常统一处理等能力,能够有效提升复杂业务场景下的执行效率。实现异步任务链的核心工具是Java 8引入的CompletableFuture类,它提供了丰富的API来支持任务编排。

基础准备:自定义线程池
默认情况下CompletableFuture会使用ForkJoinPool.commonPool()线程池,实际生产中建议自定义线程池来控制并发资源,避免共用线程池导致的资源争抢问题。
import java.util.concurrent.*;
public class AsyncTaskChainDemo {
// 自定义线程池,核心线程数10,最大线程数20,队列容量100
private static final ThreadPoolExecutor CUSTOM_THREAD_POOL = new ThreadPoolExecutor(
10,
20,
60L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<>(100),
new ThreadFactory() {
private int count = 0;
@Override
public Thread newThread(Runnable r) {
return new Thread(r, "async-task-thread-" + count++);
}
},
new ThreadPoolExecutor.CallerRunsPolicy()
);
}
串行异步任务链实现
串行任务链是指前一个异步任务执行完成后,再执行下一个异步任务,并且可以传递前一个任务的结果。使用thenApplyAsync方法可以实现带结果传递的串行任务。
import java.util.concurrent.CompletableFuture;
public class AsyncTaskChainDemo {
// 自定义线程池省略,同上
public static void main(String[] args) throws Exception {
// 第一个异步任务:查询用户ID
CompletableFuture<Integer> queryUserIdTask = CompletableFuture.supplyAsync(() -> {
System.out.println("执行查询用户ID任务,线程名:" + Thread.currentThread().getName());
try {
Thread.sleep(1000); // 模拟耗时操作
} catch (InterruptedException e) {
e.printStackTrace();
}
return 1001; // 返回用户ID
}, CUSTOM_THREAD_POOL);
// 第二个异步任务:根据用户ID查询用户信息,接收上一个任务的结果
CompletableFuture<String> queryUserInfoTask = queryUserIdTask.thenApplyAsync(userId -> {
System.out.println("执行查询用户信息任务,用户ID:" + userId + ",线程名:" + Thread.currentThread().getName());
try {
Thread.sleep(1500); // 模拟耗时操作
} catch (InterruptedException e) {
e.printStackTrace();
}
return "用户ID:" + userId + ",用户名:张三,年龄:25";
}, CUSTOM_THREAD_POOL);
// 第三个异步任务:处理用户信息并保存,接收上一个任务的结果
CompletableFuture<String> saveUserInfoTask = queryUserInfoTask.thenApplyAsync(userInfo -> {
System.out.println("执行保存用户信息任务,用户信息:" + userInfo + ",线程名:" + Thread.currentThread().getName());
try {
Thread.sleep(500); // 模拟耗时操作
} catch (InterruptedException e) {
e.printStackTrace();
}
return "保存成功:" + userInfo;
}, CUSTOM_THREAD_POOL);
// 获取最终任务结果
String result = saveUserInfoTask.get();
System.out.println("最终执行结果:" + result);
// 关闭线程池
CUSTOM_THREAD_POOL.shutdown();
}
}
并行异步任务组合
如果多个异步任务之间没有依赖关系,可以并行执行,等待所有任务完成后再执行后续操作,使用CompletableFuture.allOf可以实现这个功能。
import java.util.concurrent.CompletableFuture;
public class AsyncTaskChainDemo {
// 自定义线程池省略,同上
public static void main(String[] args) throws Exception {
// 并行任务1:查询商品库存
CompletableFuture<Integer> queryStockTask = CompletableFuture.supplyAsync(() -> {
System.out.println("执行查询商品库存任务,线程名:" + Thread.currentThread().getName());
try {
Thread.sleep(1200);
} catch (InterruptedException e) {
e.printStackTrace();
}
return 50; // 返回库存数量
}, CUSTOM_THREAD_POOL);
// 并行任务2:查询商品价格
CompletableFuture<Double> queryPriceTask = CompletableFuture.supplyAsync(() -> {
System.out.println("执行查询商品价格任务,线程名:" + Thread.currentThread().getName());
try {
Thread.sleep(800);
} catch (InterruptedException e) {
e.printStackTrace();
}
return 299.9; // 返回商品价格
}, CUSTOM_THREAD_POOL);
// 并行任务3:查询商品评价
CompletableFuture<String> queryCommentTask = CompletableFuture.supplyAsync(() -> {
System.out.println("执行查询商品评价任务,线程名:" + Thread.currentThread().getName());
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
e.printStackTrace();
}
return "好评率98%"; // 返回评价信息
}, CUSTOM_THREAD_POOL);
// 等待所有并行任务完成
CompletableFuture<Void> allTask = CompletableFuture.allOf(queryStockTask, queryPriceTask, queryCommentTask);
// 所有任务完成后,组合结果执行后续操作
CompletableFuture<String> combineResultTask = allTask.thenApplyAsync(v -> {
try {
int stock = queryStockTask.get();
double price = queryPriceTask.get();
String comment = queryCommentTask.get();
return "商品库存:" + stock + ",价格:" + price + ",评价:" + comment;
} catch (Exception e) {
e.printStackTrace();
return "获取商品信息失败";
}
}, CUSTOM_THREAD_POOL);
System.out.println("商品组合信息:" + combineResultTask.get());
CUSTOM_THREAD_POOL.shutdown();
}
}
异步任务链的异常处理
异步任务链执行过程中如果出现异常,需要统一捕获处理,避免异常丢失。可以使用exceptionally方法处理单个任务的异常,也可以使用handle方法处理所有任务的结果和异常。
import java.util.concurrent.CompletableFuture;
public class AsyncTaskChainDemo {
// 自定义线程池省略,同上
public static void main(String[] args) throws Exception {
CompletableFuture<String> taskChain = CompletableFuture.supplyAsync(() -> {
System.out.println("执行第一个任务");
return "第一步结果";
}, CUSTOM_THREAD_POOL).thenApplyAsync(result -> {
System.out.println("执行第二个任务,接收结果:" + result);
// 模拟异常
if (true) {
throw new RuntimeException("第二个任务执行失败");
}
return "第二步结果";
}, CUSTOM_THREAD_POOL).exceptionally(ex -> {
// 捕获异常,返回默认值
System.out.println("捕获到异常:" + ex.getMessage());
return "任务执行失败,使用默认结果";
}).thenApplyAsync(result -> {
System.out.println("执行第三个任务,接收结果:" + result);
return "最终处理结果:" + result;
}, CUSTOM_THREAD_POOL);
System.out.println("任务链执行结果:" + taskChain.get());
CUSTOM_THREAD_POOL.shutdown();
}
}
注意事项
- 使用CompletableFuture时尽量指定自定义线程池,避免默认线程池的资源争抢问题。
- 调用
get()方法会阻塞当前线程,实际开发中可以使用thenAccept或者回调方式处理最终结果,避免阻塞。 - 如果任务链中有不需要返回值的异步任务,可以使用
runAsync方法替代supplyAsync。 - 多个任务组合时,如果需要任意一个任务完成就执行后续操作,可以使用
CompletableFuture.anyOf方法。
Java异步任务链CompletableFuture线程池异步编程修改时间:2026-07-24 08:48:20