Java 异步编程经历了从回调地狱到声明式链式调用,再到协程模型的演进。理解每一层抽象的设计目的与应用场景,才能在高并发系统中做出正确的选择。
1. 异步编程模型演进
同步阻塞(Tomcat 线程池)
→ Future / Callback(异步回调,代码碎片化)
→ CompletableFuture(链式组合,声明式)
→ RxJava / Reactor(响应式流,背压)
→ Virtual Threads(虚拟线程,同步写法异步执行)
2. CompletableFuture:组合式异步
2.1 创建与完成
// 异步执行任务
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
// 运行 ForkJoinPool.commonPool()
return fetchData();
});
// 指定线程池
ExecutorService executor = Executors.newFixedThreadPool(4);
CompletableFuture<String> f2 = CompletableFuture.supplyAsync(() -> fetchData(), executor);
// 立即完成 / 异常完成
future.complete("value");
future.completeExceptionally(new RuntimeException("error"));
2.2 链式组合操作
CompletableFuture<Order> orderFuture = fetchUserAsync(userId)
.thenCompose(user -> fetchCartAsync(user.getId())) // 扁平化:thenApply 返回 Future 会嵌套
.thenCompose(cart -> createOrderAsync(cart))
.thenApply(order -> enrichWithTracking(order))
.exceptionally(ex -> {
log.error("Order creation failed", ex);
return Order.failed(ex.getMessage());
})
.whenComplete((order, ex) -> metrics.recordOrder(order));
核心方法对照:
| 方法 | 输入 | 输出 | 用途 |
|---|---|---|---|
thenApply | T | U | 转换结果 |
thenAccept | T | void | 消费结果 |
thenRun | — | void | 无输入执行 |
thenCompose | T | CompletionStage<U> | 扁平化链式 |
whenComplete | T/error | T | 完成回调(不改变结果) |
exceptionally | error | T | 异常恢复 |
handle | T/error | U | 统一处理 |
2.3 组合多个 Future
// 全部完成
CompletableFuture<Void> all = CompletableFuture.allOf(f1, f2, f3);
// 任一完成
CompletableFuture<Object> any = CompletableFuture.anyOf(f1, f2, f3);
// 编排:并行获取用户信息、商品信息,再计算总价
CompletableFuture<User> userF = fetchUser(userId);
CompletableFuture<List<Product>> productsF = fetchProducts(productIds);
CompletableFuture<BigDecimal> totalF = userF
.thenCombine(productsF, (user, products) ->
calculateTotal(user.getDiscountRate(), products));
2.4 超时控制
// JDK 9+
fetchDataAsync(id)
.orTimeout(5, TimeUnit.SECONDS) // 超时抛 TimeoutException
.completeOnTimeout(defaultValue, 5, TimeUnit.SECONDS); // 超时时返回默认值
3. Project Reactor:响应式流
3.1 Mono vs Flux
// Mono: 0 或 1 个元素
Mono<User> userMono = webClient.get()
.uri("/users/{id}", id)
.retrieve()
.bodyToMono(User.class);
// Flux: 0 到 N 个元素
Flux<Order> orderFlux = reactiveRepository.findByStatus(Status.PENDING);
3.2 核心操作符
Flux.just(1, 2, 3, 4, 5)
.filter(i -> i > 2) // 3, 4, 5
.map(i -> i * 2) // 6, 8, 10
.flatMap(i -> fetchAsync(i)) // 扁平化异步结果
.collectList() // Mono<List<T>>
.subscribe(
result -> log.info("Result: {}", result),
error -> log.error("Error", error),
() -> log.info("Completed")
);
冷流 vs 热流:
| 类型 | 订阅行为 | 代表 |
|---|---|---|
| 冷流 | 每个订阅者独立从头发射 | Flux.just() / Mono.fromCallable() |
| 热流 | 订阅者共享同一数据源 | Flux.share() / Sinks.Many |
3.3 背压(Backpressure)
Flux.range(1, 1000)
.onBackpressureBuffer(100) // 缓冲 100 个,超出丢弃或报错
.delayElements(Duration.ofMillis(10))
.subscribe();
// 背压策略
.onBackpressureDrop() // 丢弃多余
.onBackpressureLatest() // 保留最新
.onBackpressureError() // 溢出报错
3.4 Spring WebFlux 实战
@RestController
public class ProductController {
@GetMapping("/products")
public Flux<Product> listProducts(
@RequestParam(required = false) String category) {
return productService.findByCategory(category)
.delayElements(Duration.ofMillis(10)) // 模拟流式响应
.onErrorResume(e -> Flux.empty());
}
@PostMapping("/orders")
public Mono<ResponseEntity<Order>> createOrder(
@RequestBody @Valid OrderRequest request) {
return orderService.create(request)
.map(order -> ResponseEntity.status(HttpStatus.CREATED).body(order))
.onErrorMap(ValidationException.class,
ex -> new ResponseStatusException(HttpStatus.BAD_REQUEST, ex.getMessage()));
}
}
WebFlux 核心优势:少量线程(CPU 核数级别)处理大量并发连接,适合 IO 密集型 + 长连接场景(SSE、WebSocket)。
3.5 R2DBC 响应式数据库
public interface UserRepository extends ReactiveCrudRepository<User, Long> {
@Query("SELECT * FROM users WHERE status = :status")
Flux<User> findByStatus(String status);
}
// 使用
userRepo.findByStatus("active")
.flatMap(user -> enrichUserData(user))
.collectList()
.subscribe();
R2DBC vs JDBC:完全非阻塞的数据库访问,线程数不再受连接池大小限制。
4. Virtual Threads(虚拟线程)
4.1 线程模型的痛点
平台线程(OS 线程)
├─ 创建成本高(~1MB 栈空间)
├─ 上下文切换开销大
└─ 线程数受限于 OS 资源
高并发 IO → 大量线程阻塞等待 → 内存爆炸 / 上下文切换 CPU 飙升
4.2 虚拟线程原理
// JDK 21+ 正式版
Thread virtualThread = Thread.startVirtualThread(() -> {
// 同步写法,但底层非阻塞
String result = httpClient.send(request, BodyHandlers.ofString());
process(result);
});
// 或使用 Executor
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
Future<String> f1 = executor.submit(() -> fetchService1());
Future<String> f2 = executor.submit(() -> fetchService2());
return combine(f1.get(), f2.get());
}
实现原理:
Virtual Thread
│ JVM 管理,轻量(~几百字节)
│ 数量可达数百万
↓
Carrier Thread(平台线程,数量 = CPU 核数)
│ 虚拟线程挂载到载体线程上执行
↓
遇到阻塞 IO → JVM 自动卸载虚拟线程
│ 载体线程去执行其他虚拟线程
↓
IO 完成 → 调度器重新挂载到空闲载体线程
4.3 虚拟线程的 Pinning 问题
// ❌ 会 Pin(钉住)载体线程,阻塞不释放
synchronized (lock) { // 或 ReentrantLock 未使用 LockSupport
blockingIoCall(); // 虚拟线程无法卸载
}
// ✅ 正确做法:使用 ReentrantLock
Lock lock = new ReentrantLock();
lock.lock();
try {
blockingIoCall();
} finally {
lock.unlock();
}
4.4 适用场景对比
| 场景 | 推荐方案 | 原因 |
|---|---|---|
| CPU 密集型计算 | 平台线程 + 线程池 | 虚拟线程无优势 |
| IO 密集型 HTTP/RPC | 虚拟线程 | 同步写法,百万并发 |
| 流式数据处理 | Reactor / 虚拟线程 | 背压控制 / 简单写法 |
| 实时推送 (SSE/WebSocket) | Reactor | 原生支持背压与流语义 |
| 复杂异步编排 | Reactor | 丰富的操作符组合 |
5. 异步事务处理
// Spring 的 @Transactional 默认不支持异步
// 方案 1:在调用方同步执行事务包裹
@Transactional
public void processOrder(Long orderId) {
CompletableFuture.allOf(
CompletableFuture.runAsync(() -> updateInventory(orderId)),
CompletableFuture.runAsync(() -> updateAccount(orderId))
).join(); // 等待完成再提交(阻塞)
}
// 方案 2:Reactor + @Transactional(R2DBC)
@Transactional
public Mono<Void> processOrderReactive(Long orderId) {
return Mono.zip(
updateInventoryReactive(orderId),
updateAccountReactive(orderId)
).then();
}
6. 性能基准
| 场景 | 平台线程池 | 虚拟线程 | Reactor |
|---|---|---|---|
| 10K 并发 HTTP | 100 线程 × 100 连接 | 10K 虚拟线程 | EventLoop (核心数线程) |
| 内存占用 | ~100MB | ~20MB | ~15MB |
| 代码复杂度 | 简单 | 最简单(同步写法) | 较复杂(流式思维) |
| 延迟 | 中等(上下文切换) | 低(无 OS 切换) | 低(事件驱动) |
延伸阅读
- Java 高并发编程精要 — JUC 线程池与同步原语
- Spring Cloud 微服务架构全家桶 — WebFlux 网关实战
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。