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

如何在CompletableFuture流程中实现并发去重并执行后续阶段?

解决CompletableFuture去重时的并发阻塞问题

咱们先拆解下你现在代码里的核心问题:你在filter环节调用了CompletableFuture::get,这个方法会阻塞当前线程直到Future完成。而Stream的操作是按顺序处理每个元素的,这就导致你的并发任务被迫变成了串行——必须等前一个Future执行完、拿到结果、判断去重后,才会处理下一个Future,完全浪费了supplyAsync带来的并发优势。

接下来给你两个可行的方案,你可以根据自己的业务场景选择:

方案一:先完成所有第一阶段任务,再去重并发执行第二阶段

这个方案逻辑最直观,先让所有获取Pod的任务并发跑完,等全部结果出来后再做去重,最后再并发启动获取Instance的任务。

// 第一步:并发发起所有获取Pod的请求,收集Future
List<CompletableFuture<PodDo>> podFutures = podNames.stream()
        .map(podName -> CompletableFuture.supplyAsync(RethrowExceptionUtil.rethrowSupplier(() -> {
            PodDo podDo = getPodRetriever().getPod(envId, podName);
            podDoList.add(podDo);
            return podDo;
        }), executor))
        .collect(Collectors.toList());

// 等待所有Pod任务完成(不会阻塞主线程,这里只是创建一个完成标记Future)
CompletableFuture<Void> allPodsDone = CompletableFuture.allOf(podFutures.toArray(new CompletableFuture[0]));

// 所有Pod完成后,收集结果、去重,再并发获取Instance
CompletableFuture<List<InstanceDo>> finalResultFuture = allPodsDone.thenApply(ignored -> 
        podFutures.stream()
                .map(CompletableFuture::join) // 因为已经完成,join不会阻塞
                .map(PodDo::getInstanceName)
                .distinct() // 用Java自带的去重,如果你需要自定义key,也可以用你自己的distinctByKey逻辑
                .map(instanceName -> CompletableFuture.supplyAsync(() -> get(envId, instanceName), executor))
                .collect(Collectors.toList())
).thenCompose(allInstanceFutures -> 
        // 等待所有Instance任务完成,再收集最终结果
        CompletableFuture.allOf(allInstanceFutures.toArray(new CompletableFuture[0]))
                .thenApply(ignored -> allInstanceFutures.stream()
                        .map(CompletableFuture::join)
                        .collect(Collectors.toList()))
);

// 最终获取结果(如果需要同步等待的话)
List<InstanceDo> instanceDos = finalResultFuture.join();

这个方案的好处是逻辑简单,去重操作在所有结果都就绪后执行,没有并发判断的风险;缺点是必须等所有Pod任务都完成,才能启动Instance任务,整体耗时是「Pod任务最长耗时」+「Instance任务最长耗时」。

方案二:异步实时去重,无需等待所有第一阶段完成

如果想尽可能压缩整体耗时,让Instance任务在Pod任务完成后立即启动(只要没重复),可以用并发安全的集合来实时记录已处理的instanceName,在异步流程中判断是否需要执行第二阶段。

import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
import java.util.Set;

// 并发安全的集合,用来记录已经处理过的instanceName
Set<String> processedInstances = ConcurrentHashMap.newKeySet();

List<CompletableFuture<InstanceDo>> instanceFutures = podNames.stream()
        .map(podName -> CompletableFuture.supplyAsync(RethrowExceptionUtil.rethrowSupplier(() -> {
            PodDo podDo = getPodRetriever().getPod(envId, podName);
            podDoList.add(podDo);
            return podDo;
        }), executor))
        .map(podFuture -> podFuture.thenCompose(podDo -> {
            String instanceName = podDo.getInstanceName();
            // 原子性判断是否已处理:add返回true表示之前没处理过
            if (processedInstances.add(instanceName)) {
                // 没处理过,启动获取Instance的任务
                return CompletableFuture.supplyAsync(() -> get(envId, instanceName), executor);
            } else {
                // 已处理过,返回一个完成的空Future,后续过滤掉即可
                return CompletableFuture.completedFuture(null);
            }
        }))
        .collect(Collectors.toList());

// 收集最终结果,过滤掉重复导致的null
CompletableFuture<List<InstanceDo>> finalResultFuture = CompletableFuture.allOf(instanceFutures.toArray(new CompletableFuture[0]))
        .thenApply(ignored -> instanceFutures.stream()
                .map(CompletableFuture::join)
                .filter(Objects::nonNull)
                .collect(Collectors.toList()));

List<InstanceDo> instanceDos = finalResultFuture.join();

这个方案的优势是:只要某个Pod任务完成,且对应的instanceName没被处理过,就会立即启动对应的Instance任务,整体耗时更接近「max(Pod任务最长耗时, Instance任务最长耗时)」;缺点是需要维护并发安全集合,还要处理空结果,逻辑稍微复杂一点。

总结

  • 如果业务允许等待所有Pod任务完成再启动Instance任务,优先选方案一,简单可靠;
  • 如果追求极致的并发效率,不想浪费等待时间,可以选方案二;
  • 绝对不要在Stream的filter/map等同步操作中调用Future.get(),这会直接毁掉并发特性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 09:26:15