Future与异步编程
开篇:同步调用的痛点
想象你去餐厅吃饭,点了三道菜。同步模式就是:你点了第一道菜,然后干等着,等第一道上了,再点第二道,又干等着......三道菜总共等了 30 分钟。
异步模式就是:你一次性把三道菜都点了,厨房同时做三道菜,10 分钟后全上齐。从 30 分钟变成 10 分钟,效率提升 3 倍。
在程序中也是一样。假设你需要调用三个远程服务,每个耗时 1 秒:
// 同步方式:总共 3 秒
String result1 = callServiceA(); // 1s
String result2 = callServiceB(); // 1s
String result3 = callServiceC(); // 1s如果三个服务之间没有依赖关系,完全可以并行调用,只需要 1 秒。这就是异步编程的价值。
一、Future:异步的雏形
Java 5 引入了 Future 接口,它代表一个异步计算的结果。你可以把任务交给线程池执行,拿到一个 Future,在需要结果的时候再去取。
ExecutorService executor = Executors.newFixedThreadPool(3);
Future<String> future = executor.submit(() -> {
Thread.sleep(1000);
return "任务完成";
});
// 做其他事情...
System.out.println("主线程继续做别的事情");
// 需要结果的时候再取
String result = future.get(); // 阻塞等待
System.out.println(result);
executor.shutdown();Future 提供了几个方法:
| 方法 | 作用 |
|---|---|
get() | 阻塞等待结果 |
get(timeout, unit) | 带超时的等待 |
isDone() | 判断任务是否完成 |
cancel(boolean) | 取消任务 |
isCancelled() | 判断是否被取消 |
Future 的局限性在于:
- get() 会阻塞:调了 get() 就卡在那里,和同步没区别
- 没有回调机制:任务完成后不能自动通知你,只能你主动去问
- 不能组合:多个 Future 之间无法方便地编排(先做 A,A 完成后做 B)
- 异常处理不友好:异常被包装在 ExecutionException 中,处理起来很别扭
用轮询代替阻塞稍微好一点,但本质上还是在浪费 CPU:
while (!future.isDone()) {
System.out.println("还没好,再等等...");
Thread.sleep(100);
}
System.out.println("结果:" + future.get());我们需要的是:任务完成后主动通知我,多个异步任务可以灵活编排。 这就是 CompletableFuture 要解决的问题。
二、CompletableFuture:异步编程利器
Java 8 引入的 CompletableFuture 实现了 Future 和 CompletionStage 两个接口,提供了强大的异步编程能力。它的核心理念是:把异步操作串成一条链,每个阶段完成后自动触发下一个阶段。
2.1 创建异步任务
两个核心方法:
// 无返回值
CompletableFuture<Void> f1 = CompletableFuture.runAsync(() -> {
System.out.println("异步执行,无返回值");
});
// 有返回值
CompletableFuture<Integer> f2 = CompletableFuture.supplyAsync(() -> {
System.out.println("异步执行,有返回值");
return 42;
});如果不指定线程池,默认使用 ForkJoinPool.commonPool()。生产环境建议指定自定义线程池:
ExecutorService myPool = Executors.newFixedThreadPool(4);
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
return "使用自定义线程池";
}, myPool);2.2 链式处理
CompletableFuture 最强大的地方就是链式调用,告别回调地狱:
thenApply:转换结果(类似 map)
CompletableFuture<String> future = CompletableFuture
.supplyAsync(() -> 100)
.thenApply(num -> num * 2) // 200
.thenApply(num -> "结果是 " + num); // "结果是 200"
System.out.println(future.get()); // 结果是 200thenAccept:消费结果,无返回值
CompletableFuture.supplyAsync(() -> "Hello")
.thenAccept(s -> System.out.println("收到:" + s));thenRun:不关心上一步结果,只是接着执行
CompletableFuture.supplyAsync(() -> "Hello")
.thenRun(() -> System.out.println("上一步做完了,我接着做"));thenCompose:链接两个异步操作(类似 flatMap)
CompletableFuture<String> future = CompletableFuture
.supplyAsync(() -> "用户ID-123")
.thenCompose(userId -> CompletableFuture.supplyAsync(
() -> "查询用户 " + userId + " 的信息"
));thenCombine:合并两个独立的异步结果
CompletableFuture<Integer> price = CompletableFuture.supplyAsync(() -> 100);
CompletableFuture<Integer> discount = CompletableFuture.supplyAsync(() -> 20);
CompletableFuture<Integer> finalPrice = price.thenCombine(discount,
(p, d) -> p - d);
System.out.println(finalPrice.get()); // 802.3 异常处理
异步链中出了异常怎么办?CompletableFuture 提供了优雅的处理方式:
exceptionally:出异常时提供一个兜底值
CompletableFuture<Integer> future = CompletableFuture
.supplyAsync(() -> {
if (true) throw new RuntimeException("出错了");
return 100;
})
.exceptionally(ex -> {
System.out.println("捕获异常:" + ex.getMessage());
return -1; // 兜底值
});
System.out.println(future.get()); // -1handle:同时处理正常结果和异常
CompletableFuture<String> future = CompletableFuture
.supplyAsync(() -> {
int result = 10 / 0; // 故意抛异常
return result;
})
.handle((result, ex) -> {
if (ex != null) {
return "发生异常:" + ex.getMessage();
}
return "正常结果:" + result;
});
System.out.println(future.get()); // 发生异常:...whenComplete:无论成功失败都会执行,但不改变结果
CompletableFuture<Integer> future = CompletableFuture
.supplyAsync(() -> 100)
.whenComplete((result, ex) -> {
if (ex == null) {
System.out.println("成功,结果:" + result);
} else {
System.out.println("失败,异常:" + ex.getMessage());
}
});2.4 多任务编排
allOf:等待所有任务完成
CompletableFuture<String> f1 = CompletableFuture.supplyAsync(() -> "任务1");
CompletableFuture<String> f2 = CompletableFuture.supplyAsync(() -> "任务2");
CompletableFuture<String> f3 = CompletableFuture.supplyAsync(() -> "任务3");
CompletableFuture.allOf(f1, f2, f3).thenRun(() -> {
System.out.println(f1.join()); // 任务1
System.out.println(f2.join()); // 任务2
System.out.println(f3.join()); // 任务3
});anyOf:任何一个任务完成就返回
CompletableFuture<Object> fastest = CompletableFuture.anyOf(f1, f2, f3);
System.out.println("最快完成的:" + fastest.get());三、实战:电商比价
需求:同一款产品,同时查询京东、淘宝、天猫、拼多多等多个平台的价格,返回价格清单。
先看同步方式——假设每个查询耗时 1 秒,5 个平台就要 5 秒:
static List<NetMall> mallList = Arrays.asList(
new NetMall("jd"), new NetMall("tb"), new NetMall("tm"),
new NetMall("dw"), new NetMall("pdd")
);
// 同步方式:逐个查询
public static List<String> getPriceSync(List<NetMall> list, String product) {
return list.stream()
.map(mall -> String.format("%s in %s price is %.2f",
product, mall.getName(), mall.calcPrice(product)))
.collect(Collectors.toList());
}再看异步方式——所有查询并行执行,只需 1 秒多:
// 异步方式:并行查询
public static List<String> getPriceAsync(List<NetMall> list, String product) {
return list.stream()
.map(mall -> CompletableFuture.supplyAsync(() ->
String.format("%s in %s price is %.2f",
product, mall.getName(), mall.calcPrice(product))))
.collect(Collectors.toList()) // 先收集所有 Future
.stream()
.map(CompletableFuture::join) // 再统一等待结果
.collect(Collectors.toList());
}注意这里的两段式 stream:先把所有查询提交出去(异步启动),再统一获取结果。不能写成一个 stream,否则会变成"提交一个等一个",失去了并行的意义。
实测结果:同步方式约 5 秒,异步方式约 1 秒。平台越多,差距越大。
四、常见面试题精选
Q1:CompletableFuture 的底层是如何实现的?
CompletableFuture 内部采用了**链式结构(Completion 链)**来管理异步阶段。核心机制包括:
- Completion 链:每个 CompletableFuture 关联一个 Completion 链,每个节点代表一个异步阶段
- 事件驱动:当一个阶段完成后,会触发完成事件,自动执行下一个阶段
- ForkJoinPool:默认使用 ForkJoinPool.commonPool() 来执行异步任务
- CAS 操作:内部通过 CAS 来保证状态更新的线程安全
Q2:如何避免回调地狱?
传统的回调方式,多个异步操作嵌套后代码会层层缩进,难以阅读和维护——这就是"回调地狱"。
CompletableFuture 通过链式调用解决了这个问题:
// 回调地狱写法(不推荐)
service.call(result1 -> {
service2.call(result1, result2 -> {
service3.call(result2, result3 -> {
// 越嵌越深...
});
});
});
// CompletableFuture 链式写法
CompletableFuture
.supplyAsync(() -> callService1())
.thenApply(r1 -> callService2(r1))
.thenApply(r2 -> callService3(r2))
.thenAccept(r3 -> System.out.println("最终结果:" + r3));Q3:如何对多线程进行编排?
使用 CompletableFuture 的组合方法:
- 串行:
thenApply/thenCompose - 并行 + 全等:
allOf - 并行 + 先到先得:
anyOf - 两两合并:
thenCombine
Q4:CompletableFuture 使用时要注意什么?
- 指定自定义线程池:不要依赖默认的 ForkJoinPool.commonPool(),它是全局共享的,一旦被占满会影响其他功能
- 处理异常:不要吞掉异常,至少用
exceptionally或handle来处理 - 避免阻塞:不要在异步链中调用
get(),用thenApply/thenAccept代替 - 注意线程安全:虽然 CompletableFuture 本身是线程安全的,但它操作的数据可能不是
- 主线程退出问题:如果使用默认线程池(守护线程),主线程结束后异步任务可能被中断
小结
| 工具 | 阻塞/非阻塞 | 回调 | 组合 | 异常处理 | 推荐 |
|---|---|---|---|---|---|
| Future + get() | 阻塞 | 不支持 | 不方便 | ExecutionException | 简单场景 |
| CompletableFuture | 非阻塞 | 支持 | 非常方便 | exceptionally/handle | 生产首选 |
核心要点:
- Future 是异步的基础,但
get()会阻塞,限制了它的使用场景 - CompletableFuture 通过链式回调和丰富的编排 API,真正实现了非阻塞的异步编程
- 在需要并行查询、多服务编排的场景中,CompletableFuture 可以带来数倍的性能提升
- 始终记得:指定线程池、处理异常、避免在链中阻塞
附录:异步编程实用模式
模式一:异步批量通知
在实际业务中,经常需要异步通知多个下游系统,全部通知成功后再更新状态:
// 异步通知下游系统
CompletableFuture<Void> allNotices = CompletableFuture.allOf(
noticeDetails.stream()
.map(detail -> CompletableFuture.supplyAsync(() -> {
sendNotice(detail);
return null;
}))
.toArray(CompletableFuture[]::new)
);
// 所有通知完成后更新状态
allNotices.whenComplete((v, e) -> {
if (e == null) {
noticeOrder.setState("SUCCESS");
noticeOrderService.updateById(noticeOrder);
} else {
log.error("通知失败", e);
}
});模式二:超时控制
CompletableFuture 在 Java 9 之后支持原生超时:
CompletableFuture<String> future = CompletableFuture
.supplyAsync(() -> {
// 模拟耗时操作
try { Thread.sleep(5000); } catch (InterruptedException e) {}
return "结果";
})
.orTimeout(2, TimeUnit.SECONDS) // 超时抛异常
.exceptionally(ex -> "超时了,返回默认值");模式三:Async 后缀方法
注意 CompletableFuture 的很多方法都有 Async 后缀版本,如 thenApply 和 thenApplyAsync。区别在于:
- 不带 Async:在上一个任务的线程中执行(可能是 main 线程)
- 带 Async:提交到线程池异步执行
在链路较长时,使用 Async 版本可以避免阻塞主线程,但也会增加线程切换的开销。要根据具体场景选择。