You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何CompletableFuture的thenApplyAsync有时在主线程执行execute方法?

问题描述

我重写了ThreadPoolExecutor中的java.util.concurrent.Executor.execute方法,新实现仅对Runnable做装饰后调用原execute方法。当使用两个该类的执行器时,执行以下代码:

supplyAsync(() -> foo(), firstExecutor).thenApplyAsync(firstResult -> bar(), secondExecutor)

会触发两次execute调用,通常分别由main线程和firstExecutor执行,但有时两次都由main线程执行。这是否与supplyAsync中Supplier的执行耗时有关?

最小复现示例

重复执行10000次测试,约3次会抛出java.lang.AssertionError: Unexpected second decorator: main错误:

package com.foo;

import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.RepeatedTest;

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

class DecorationTest {

    record WhoCalled(String decorator, String runnable) {}

    static class DecoratedExecutor extends ThreadPoolExecutor{

        private final List<WhoCalled> callers;

        public DecoratedExecutor(List<WhoCalled> callers, String threadName) {
            super(1, 1, 1, TimeUnit.MINUTES, new LinkedBlockingQueue<>(), runnable -> new Thread(runnable, threadName));
            this.callers = callers;
        }

        @Override
        public void execute(final Runnable command) {
            String decoratingThread = Thread.currentThread().getName();
            Runnable decorated = () -> {
                String runningThread = Thread.currentThread().getName();
                callers.add(new WhoCalled(decoratingThread, runningThread));
                command.run();
            };
            super.execute(decorated);
        }
    }

    List<WhoCalled> callers;
    ExecutorService firstExecutor;
    ExecutorService secondExecutor;

    @BeforeEach
    void beforeEach() {
        callers = new ArrayList<>();
        firstExecutor = new DecoratedExecutor(callers, "firstExecutor");
        secondExecutor = new DecoratedExecutor(callers, "secondExecutor");
    }

    @AfterEach
    void afterEach() {
        firstExecutor.shutdown();
        secondExecutor.shutdown();
    }


    @RepeatedTest(10_000)
    void testWhoCalled() throws Exception {
        Integer result = CompletableFuture.supplyAsync(() -> 1, firstExecutor)
                .thenApplyAsync(supplyResult -> supplyResult, secondExecutor)
                .get();

        assert result == 1;

        WhoCalled firstCallers = callers.get(0);
        assert firstCallers.decorator().equals("main");
        assert firstCallers.runnable().equals("firstExecutor");

        WhoCalled secondCallers = callers.get(1);
        assert secondCallers.decorator().equals("firstExecutor") : "Unexpected second decorator: " + secondCallers.decorator;
        assert secondCallers.runnable().equals("secondExecutor");
    }
}
原因分析与结论

这和supplyAsync中Supplier的**执行耗时极短(或为0)**直接相关,本质是CompletableFuture的快速路径优化导致的。

具体逻辑

  1. supplyAsync将Supplier任务提交到firstExecutor后,main线程立即返回CompletableFuture并进入get()等待阶段
  2. 由于测试中的Supplier只是直接返回1,执行耗时极短,可能在firstExecutor的线程还没完全启动时就已经完成,甚至可能在supplyAsync的提交逻辑中就被执行(取决于线程池调度时机)
  3. 当CompletableFuture完成时,如果当前线程(main)正在等待get(),CompletableFuture会触发快速路径优化:直接在当前线程提交后续的thenApplyAsync任务到secondExecutor,而非由完成Supplier的firstExecutor线程提交

验证方式

如果给Supplier增加少量耗时,比如:

supplyAsync(() -> {
    Thread.sleep(10);
    return 1;
}, firstExecutor)

你会发现几乎不会再出现second decorator: main的情况——因为此时Supplier的执行耗时足够让firstExecutor线程真正完成任务,并由它来触发后续的thenApplyAsync提交逻辑。

补充说明

CompletableFuture的这种优化是为了减少线程切换开销,当任务可以快速完成时,尽量在当前线程处理后续依赖逻辑。这是符合其设计初衷的,但会导致你观察到的这种“非预期”的线程提交情况。

内容的提问来源于stack exchange,提问作者nluk

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 05:09:52