一个聚合接口需要依次查询用户资料、最近订单和积分余额。三个下游各耗时 200 毫秒时,串行调用至少需要 600 毫秒;如果三者互不依赖,并行执行后的理论耗时接近最慢的那个请求。

但把三行代码套进 supplyAsync 并不等于完成了异步治理。我们还需要回答:任务在哪个线程池运行?某个下游超时后是否拖垮整个接口?异常是继续传播还是返回降级值?请求链路中的 TraceId 能否传到异步线程?

本文示例基于 Java 21 + Spring Boot 3.x


一、先区分并发、异步与并行

这三个概念经常被混在一起:

  • 并发:多个任务在同一时间段内推进;
  • 并行:多个任务在同一时刻由不同计算资源执行;
  • 异步:调用方提交任务后,不必原地等待任务结束。

CompletableFuture 负责描述异步任务及其依赖关系,但实际是否并行,取决于使用的 Executor 以及可用线程数。

下面的写法虽然返回 CompletableFuture,却没有产生异步执行:

CompletableFuture<String> future =
        CompletableFuture.completedFuture("already done");

supplyAsync 会将有返回值的任务提交给执行器:

CompletableFuture<String> future = CompletableFuture.supplyAsync(
        () -> remoteClient.query(),
        executor
);

Oracle 的 Java 21 文档明确说明:未显式传入 Executor 的异步方法默认使用 ForkJoinPool.commonPool()。这对简单演示很方便,对生产服务却通常缺少隔离、容量控制和清晰的线程命名。


二、为什么不要直接依赖 commonPool?

以下代码能运行,但不适合作为生产默认方案:

CompletableFuture<UserProfile> future =
        CompletableFuture.supplyAsync(() -> userClient.query(userId));

问题在于:

  1. JVM 内其他代码也可能使用公共池,容易互相争抢线程;
  2. 阻塞式 HTTP、JDBC 调用会长期占用工作线程;
  3. 无法针对某类下游单独设置队列、拒绝策略和监控指标;
  4. 默认线程名难以快速关联具体业务。

为聚合查询创建独立线程池:

@Configuration
public class AsyncExecutorConfig {

    @Bean("dashboardExecutor")
    public ThreadPoolTaskExecutor dashboardExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(8);
        executor.setMaxPoolSize(16);
        executor.setQueueCapacity(200);
        executor.setThreadNamePrefix("dashboard-");
        executor.setWaitForTasksToCompleteOnShutdown(true);
        executor.setAwaitTerminationSeconds(20);
        executor.setRejectedExecutionHandler(
                new ThreadPoolExecutor.CallerRunsPolicy()
        );
        executor.initialize();
        return executor;
    }
}

线程池参数不能照抄固定公式,需要从目标并发量、平均耗时、下游容量和允许排队时间反推。尤其要注意:扩大线程池只能增加并发,不能扩大数据库连接池或下游接口的真实容量。

CallerRunsPolicy 会让提交任务的线程执行被拒绝的任务,从而形成一定反压,但它也会抬高当前请求延迟。对于必须快速失败的接口,可以改用 AbortPolicy,并将 RejectedExecutionException 映射为明确的降级响应。


三、示例:并行聚合三个下游结果

先定义返回模型:

public record UserProfile(Long userId, String nickname) {
}

public record OrderSummary(int count, BigDecimal totalAmount) {
    public static OrderSummary empty() {
        return new OrderSummary(0, BigDecimal.ZERO);
    }
}

public record PointsBalance(long available) {
    public static PointsBalance unavailable() {
        return new PointsBalance(-1);
    }
}

public record DashboardView(
        UserProfile profile,
        OrderSummary orders,
        PointsBalance points
) {
}

聚合服务如下:

@Service
public class DashboardService {

    private static final Logger log =
            LoggerFactory.getLogger(DashboardService.class);

    private final UserClient userClient;
    private final OrderClient orderClient;
    private final PointsClient pointsClient;
    private final Executor dashboardExecutor;

    public DashboardService(
            UserClient userClient,
            OrderClient orderClient,
            PointsClient pointsClient,
            @Qualifier("dashboardExecutor") Executor dashboardExecutor
    ) {
        this.userClient = userClient;
        this.orderClient = orderClient;
        this.pointsClient = pointsClient;
        this.dashboardExecutor = dashboardExecutor;
    }

    public DashboardView query(Long userId) {
        CompletableFuture<UserProfile> profileFuture =
                CompletableFuture.supplyAsync(
                                () -> userClient.query(userId),
                                dashboardExecutor
                        )
                        .orTimeout(500, TimeUnit.MILLISECONDS);

        CompletableFuture<OrderSummary> ordersFuture =
                CompletableFuture.supplyAsync(
                                () -> orderClient.summary(userId),
                                dashboardExecutor
                        )
                        .completeOnTimeout(
                                OrderSummary.empty(),
                                800,
                                TimeUnit.MILLISECONDS
                        )
                        .exceptionally(ex -> {
                            log.warn("query order summary failed, userId={}",
                                    userId, unwrap(ex));
                            return OrderSummary.empty();
                        });

        CompletableFuture<PointsBalance> pointsFuture =
                CompletableFuture.supplyAsync(
                                () -> pointsClient.balance(userId),
                                dashboardExecutor
                        )
                        .orTimeout(300, TimeUnit.MILLISECONDS)
                        .exceptionally(ex -> {
                            log.warn("query points failed, userId={}",
                                    userId, unwrap(ex));
                            return PointsBalance.unavailable();
                        });

        CompletableFuture.allOf(
                profileFuture,
                ordersFuture,
                pointsFuture
        ).join();

        return new DashboardView(
                profileFuture.join(),
                ordersFuture.join(),
                pointsFuture.join()
        );
    }

    private static Throwable unwrap(Throwable throwable) {
        if (throwable instanceof CompletionException
                && throwable.getCause() != null) {
            return throwable.getCause();
        }
        return throwable;
    }
}

这段代码刻意区分了关键数据和非关键数据:

  • 用户资料是页面主体,失败时让异常继续传播;
  • 订单汇总超时或失败时返回空汇总;
  • 积分不可用时返回 -1,由前端展示“暂不可用”。

降级值应当在业务语义上可区分。“查询失败”不能悄悄伪装成“真实余额为 0”,否则系统会把可用性问题变成数据正确性问题。


四、allOf 为什么不能直接拿到所有结果?

CompletableFuture.allOf 的返回类型是 CompletableFuture<Void>。它只表达“这些任务全部结束”,不会自动收集结果:

CompletableFuture<Void> all = CompletableFuture.allOf(f1, f2, f3);
all.join();

任务类型一致时,可以封装一个通用方法:

public static <T> CompletableFuture<List<T>> sequence(
        List<CompletableFuture<T>> futures
) {
    CompletableFuture<Void> all = CompletableFuture.allOf(
            futures.toArray(CompletableFuture[]::new)
    );

    return all.thenApply(ignored -> futures.stream()
            .map(CompletableFuture::join)
            .toList());
}

使用方式:

List<CompletableFuture<Product>> futures = productIds.stream()
        .map(id -> CompletableFuture.supplyAsync(
                () -> productClient.query(id), executor))
        .toList();

List<Product> products = sequence(futures).join();

需要注意,allOf 中任意任务异常完成,组合结果也会异常完成。如果希望“成功几个返回几个”,必须先为每个子任务转换结果,例如包装成 TaskResult<T>,而不是在最外层统一吞掉异常。

public record TaskResult<T>(T value, Throwable error) {
    public static <T> TaskResult<T> success(T value) {
        return new TaskResult<>(value, null);
    }

    public static <T> TaskResult<T> failure(Throwable error) {
        return new TaskResult<>(null, error);
    }
}

CompletableFuture<TaskResult<Product>> safeFuture =
        CompletableFuture.supplyAsync(
                        () -> productClient.query(productId), executor)
                .handle((value, error) -> error == null
                        ? TaskResult.success(value)
                        : TaskResult.failure(error));

这样既能保留部分成功结果,也不会丢失失败原因。


五、thenApply、thenCompose 与 thenCombine 怎么选?

1. thenApply:同步转换结果

当后续步骤只需要转换上一步结果时使用:

CompletableFuture<String> nicknameFuture = profileFuture
        .thenApply(UserProfile::nickname);

它类似 Stream.map。如果转换函数返回的本身也是 CompletableFuture,就会出现嵌套:

CompletableFuture<CompletableFuture<Address>> nested = profileFuture
        .thenApply(profile -> addressService.queryAsync(profile.userId()));

2. thenCompose:串联两个异步任务

后一异步任务依赖前一任务结果时,使用 thenCompose 展平:

CompletableFuture<Address> addressFuture = profileFuture
        .thenCompose(profile ->
                addressService.queryAsync(profile.userId()));

它类似 flatMap

3. thenCombine:合并两个互相独立的结果

CompletableFuture<AccountView> accountFuture =
        profileFuture.thenCombine(
                pointsFuture,
                (profile, points) -> new AccountView(profile, points)
        );

一个实用判断方法是:

  • 只转换一个结果:thenApply
  • 下一个异步任务依赖当前结果:thenCompose
  • 两个独立任务完成后合并:thenCombine
  • 等待任意一个任务完成:applyToEither
  • 等待一组任务全部完成:allOf

六、超时不是取消:最容易被忽略的陷阱

Java 9 以后提供了两个常用超时 API:

future.orTimeout(500, TimeUnit.MILLISECONDS);
future.completeOnTimeout(fallback, 500, TimeUnit.MILLISECONDS);

它们的差别是:

  • orTimeout:到期后以 TimeoutException 异常完成;
  • completeOnTimeout:到期后用指定默认值正常完成。

Future 超时不代表底层 I/O 已停止。例如 HTTP 请求仍可能占用连接和线程,随后才真正返回。CompletableFuture.cancel(true) 也不能被理解成可靠地中断所有底层操作。

因此必须同时配置资源自身的超时:

HttpClient client = HttpClient.newBuilder()
        .connectTimeout(Duration.ofMillis(300))
        .build();

HttpRequest request = HttpRequest.newBuilder(uri)
        .timeout(Duration.ofMillis(500))
        .GET()
        .build();

数据库查询也应设置连接获取、事务和语句执行超时。合理的超时层级通常满足:

下游单次请求超时 < 聚合任务超时 < Web 接口总超时 < 网关超时

如果内层超时反而比网关更长,请求即使已经被客户端放弃,服务端仍会继续消耗资源。


七、exceptionally、handle、whenComplete 的区别

三个方法都能观察异常,但语义不同:

方法 正常时执行 异常时执行 能否改变结果 典型用途
exceptionally 异常降级
handle 将成功和失败统一转换
whenComplete 通常不改变 日志、指标、清理

exceptionally 适合提供同类型降级值:

future.exceptionally(ex -> fallbackValue);

handle 适合将结果转换为统一包装:

future.handle((value, ex) -> ex == null
        ? ApiResult.success(value)
        : ApiResult.failure(unwrap(ex).getMessage()));

whenComplete 适合记录耗时和异常:

long start = System.nanoTime();

future.whenComplete((value, ex) -> {
    long costMs = TimeUnit.NANOSECONDS.toMillis(
            System.nanoTime() - start
    );
    metrics.record("points", costMs, ex == null);
});

不要在 whenComplete 中做可能失败的业务补偿。观察逻辑一旦抛出新异常,可能覆盖或干扰原任务结果,让排障更加困难。


八、异步线程中的 TraceId 为什么丢了?

日志框架的 MDC 通常基于 ThreadLocal。任务切换到线程池后,原请求线程中的 TraceId 不会自动出现。

Spring 的 TaskDecorator 可以在提交任务时复制上下文:

public class MdcTaskDecorator implements TaskDecorator {

    @Override
    public Runnable decorate(Runnable runnable) {
        Map<String, String> callerContext = MDC.getCopyOfContextMap();

        return () -> {
            Map<String, String> workerContext = MDC.getCopyOfContextMap();
            try {
                if (callerContext != null) {
                    MDC.setContextMap(callerContext);
                } else {
                    MDC.clear();
                }
                runnable.run();
            } finally {
                if (workerContext != null) {
                    MDC.setContextMap(workerContext);
                } else {
                    MDC.clear();
                }
            }
        };
    }
}

注册到线程池:

executor.setTaskDecorator(new MdcTaskDecorator());

finally 中恢复或清理上下文非常重要。线程池会复用线程,如果只设置不清理,下一个请求可能继承上一个请求的 TraceId。


九、CompletableFuture 与 @Async 如何选择?

@Async 适合将某个 Spring Bean 方法声明为异步入口,CompletableFuture 更适合在方法内部表达复杂编排。两者可以组合,但需要注意 Spring AOP 的代理边界:同一个对象内部的自调用不会经过代理,因此 @Async 可能失效。

@Service
public class ReportService {

    @Async("dashboardExecutor")
    public CompletableFuture<Report> generate(Long userId) {
        return CompletableFuture.completedFuture(doGenerate(userId));
    }
}

调用 generate 的代码应来自另一个 Spring Bean。对于聚合接口,直接注入 Executor 并显式使用 supplyAsync,依赖关系通常更直观,也更容易测试。

此外,返回 void@Async 方法无法通过 Future 把异常交还调用方,只能依赖 AsyncUncaughtExceptionHandler。涉及可靠业务结果时,应优先返回 CompletableFuture<T>,或者使用消息队列承载真正的后台任务。


十、如何测试异步编排?

业务测试不应该依赖真实等待。可以给服务注入一个当前线程执行器,使异步代码同步运行:

Executor directExecutor = Runnable::run;

针对超时和异常路径,再使用可控的 Stub:

@Test
void should_fallback_when_points_service_fails() {
    PointsClient pointsClient = userId -> {
        throw new IllegalStateException("points unavailable");
    };

    DashboardService service = new DashboardService(
            userClient,
            orderClient,
            pointsClient,
            Runnable::run
    );

    DashboardView view = service.query(1L);

    assertThat(view.points().available()).isEqualTo(-1);
}

集成测试还应覆盖:

  • 一个关键任务失败时,接口是否返回预期错误;
  • 非关键任务失败时,降级结果是否可识别;
  • 线程池队列满时,系统是反压、降级还是快速失败;
  • 总耗时是否接近最慢子任务,而不是所有任务耗时之和;
  • TraceId 是否存在串号或丢失。

十一、生产环境检查清单

上线前至少确认以下事项:

  1. 每个 *Async 调用是否显式使用了合适的执行器;
  2. 阻塞任务与 CPU 密集任务是否使用不同线程池;
  3. 线程池大小是否受到下游连接池和限流能力约束;
  4. 队列是否有界,拒绝策略是否符合接口语义;
  5. 每个远程调用是否配置连接超时和读取超时;
  6. 关键数据与非关键数据是否采用不同失败策略;
  7. 降级值是否会和真实业务数据混淆;
  8. 是否监控活跃线程、队列长度、拒绝次数和任务耗时;
  9. MDC、安全上下文等 ThreadLocal 数据是否正确传递和清理;
  10. 异常是否保留原始 cause,而不是只记录 CompletionException

十二、总结

高质量的 CompletableFuture 代码,核心不在于链式调用写得多漂亮,而在于四个边界是否清晰:

  • 执行边界:任务在哪个线程池运行;
  • 时间边界:任务和底层资源何时超时;
  • 失败边界:哪些错误传播,哪些结果允许降级;
  • 容量边界:线程、队列和下游连接数能承受多少并发。

当任务之间确实存在依赖或并行聚合需求时,CompletableFuture 很有表达力;如果只是把长任务扔到后台且要求可靠执行,消息队列、任务表和调度系统通常比进程内 Future 更合适。


参考资料

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐