08. 异步与响应式编程

从 Future 到 CompletableFuture,从 Project Reactor 到虚拟线程,掌握 Java 异步编程的演进脉络与实战技巧

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));

核心方法对照

方法输入输出用途
thenApplyTU转换结果
thenAcceptTvoid消费结果
thenRunvoid无输入执行
thenComposeTCompletionStage<U>扁平化链式
whenCompleteT/errorT完成回调(不改变结果)
exceptionallyerrorT异常恢复
handleT/errorU统一处理

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 并发 HTTP100 线程 × 100 连接10K 虚拟线程EventLoop (核心数线程)
内存占用~100MB~20MB~15MB
代码复杂度简单最简单(同步写法)较复杂(流式思维)
延迟中等(上下文切换)低(无 OS 切换)低(事件驱动)

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「java-enterprise」更多文章

  1. 限流算法深度解析:令牌桶、漏桶与滑动窗口计数
  2. Java 代码质量:SonarQube、Checkstyle 与 SpotBugs 工程化实践
  3. Spring IoC 容器与依赖注入原理深度剖析