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

如何顺序链式调用多个CompletionStage并收集结果为列表?

问题描述

我有一个返回类型为CompletionStage<Integer>的长耗时业务任务,希望通过for循环多次运行这些任务并将结果收集为列表,但必须顺序执行(不能并发)。

示例process函数:

private CompletionStage<Integer> process(int a) {
    return CompletableFuture.supplyAsync(() -> {
        System.err.printf("%s dispatch %d\n", LocalDateTime.now(), a);

        // 模拟长耗时业务逻辑
        return a + 10;
    }).whenCompleteAsync((e, t) -> {
        if (t != null)
            System.err.printf("!!! error processing '%d' !!!\n", a);

        System.err.printf("%s finish %d\n", LocalDateTime.now(), e);
    });
}

第一种实现(可行但有缺陷)

通过thenApplyAsync结合join()实现顺序执行,但在CompletableFuture内部调用join()会导致线程阻塞,可能占用额外线程资源:

// First approach
List<Integer> arr = IntStream.range(1, 10).boxed().collect(Collectors.toList());

CompletionStage<List<Integer>> result = CompletableFuture.completedFuture(new ArrayList<>());
for (Integer element: arr) {
    result = result.thenApplyAsync((ret) -> {
        Integer a = process(element).toCompletableFuture().join();

        ret.add(a);

        return ret;
    });
}

List<Integer> computeResult = result.toCompletableFuture().join();

执行日志显示任务按顺序执行:

2022-11-01T10:43:24.571573 dispatch 1
2022-11-01T10:43:24.571999 finish 11
2022-11-01T10:43:24.572414 dispatch 2
2022-11-01T10:43:24.572629 finish 12
2022-11-01T10:43:24.572825 dispatch 3
2022-11-01T10:43:24.572984 finish 13
2022-11-01T10:43:24.573097 dispatch 4
2022-11-01T10:43:24.573227 finish 14
2022-11-01T10:43:24.573354 dispatch 5
2022-11-01T10:43:24.573541 finish 15
2022-11-01T10:43:24.573657 dispatch 6
2022-11-01T10:43:24.573813 finish 16
2022-11-01T10:43:24.573929 dispatch 7
2022-11-01T10:43:24.574055 finish 17
2022-11-01T10:43:24.574168 dispatch 8
2022-11-01T10:43:24.574326 finish 18
2022-11-01T10:43:24.574428 dispatch 9
2022-11-01T10:43:24.574589 finish 19

第二种实现(失败,并发执行)

尝试用thenCombineAsync重构,但该方法会同时触发两个CompletionStage执行,导致任务并发,不符合顺序要求:

// Second approach
List<Integer> arr = IntStream.range(1, 10).boxed().collect(Collectors.toList());

CompletionStage<List<Integer>> result = CompletableFuture.completedFuture(new ArrayList<>());
for (Integer element : arr) {
    result = result.thenCombineAsync(process(element), (array, ret) -> { array.add(ret); return array; });
}

result.toCompletableFuture().join();

执行日志显示任务乱序执行:

2022-11-01T10:44:36.875930 dispatch 1
2022-11-01T10:44:36.876438 finish 11
2022-11-01T10:44:36.876461 dispatch 2
2022-11-01T10:44:36.876832 dispatch 4
2022-11-01T10:44:36.876847 finish 12
2022-11-01T10:44:36.876908 dispatch 3
2022-11-01T10:44:36.876968 dispatch 5
2022-11-01T10:44:36.877108 dispatch 8
2022-11-01T10:44:36.877042 dispatch 6
2022-11-01T10:44:36.876988 finish 14
2022-11-01T10:44:36.877123 dispatch 9
2022-11-01T10:44:36.877275 finish 18
2022-11-01T10:44:36.877195 finish 15
2022-11-01T10:44:36.877371 finish 19
2022-11-01T10:44:36.877262 dispatch 7
2022-11-01T10:44:36.877316 finish 16
2022-11-01T10:44:36.877191 finish 13
2022-11-01T10:44:36.877553 finish 17
最优解决方案:使用thenComposeAsync链式调用

thenComposeAsync的作用是等待前一个CompletionStage完成后,再启动下一个CompletionStage,完美符合顺序执行要求,同时避免线程阻塞问题。

实现代码:

List<Integer> arr = IntStream.range(1, 10).boxed().collect(Collectors.toList());

// 初始化已完成的Future,持有空列表
CompletionStage<List<Integer>> result = CompletableFuture.completedFuture(new ArrayList<>());

for (Integer element : arr) {
    // 使用thenComposeAsync确保前一个任务完成后再执行当前process
    result = result.thenComposeAsync(currentList -> 
        // 执行process任务,完成后将结果加入列表
        process(element).thenApplyAsync(resultValue -> {
            currentList.add(resultValue);
            return currentList;
        })
    );
}

// 等待所有任务完成,获取最终结果
List<Integer> computeResult = result.toCompletableFuture().join();

方案说明

  • thenComposeAsync会等待前一个CompletionStage<List<Integer>>完成后,才会执行传入的函数,而函数内部会启动当前的process(element)任务,严格保证任务顺序。
  • 全程无阻塞调用(如join()),所有操作都是异步链式执行,不会浪费线程资源。
  • 执行日志与第一种方法一致,任务按顺序dispatch和finish。

可选优化:函数式无副作用实现

如果想避免修改同一个列表的副作用,可以每次返回新列表,更符合函数式编程风格:

result = result.thenComposeAsync(currentList -> 
    process(element).thenApplyAsync(resultValue -> {
        List<Integer> newList = new ArrayList<>(currentList);
        newList.add(resultValue);
        return newList;
    })
);

这种方式会有少量列表复制开销,可根据业务场景选择。

内容的提问来源于stack exchange,提问作者Suh. Junmin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 01:15:45