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

使用带Executor的CompletableFuture线程安全收集异常Foo名称

优化方案:结合ExecutorService线程安全收集失败Foo列表

核心解决思路

原代码依赖CompletableFuture.supplyAsync()默认的ForkJoinPool,面对超大规模列表时易出现线程资源耗尽问题。我们需要自定义可控线程池,同时用线程安全集合收集失败的Foo名称,规避并发写入的线程安全风险。

优化后完整代码

import java.util.List;
import java.util.Objects;
import java.util.concurrent.*;
import java.util.stream.Collectors;

public class FooProcessor {
    // 按业务场景配置线程池参数:IO密集型任务可适当提高核心线程数,避免资源浪费
    private static final ExecutorService EXECUTOR_SERVICE = new ThreadPoolExecutor(
            Runtime.getRuntime().availableProcessors() * 2,
            Runtime.getRuntime().availableProcessors() * 4,
            60L,
            TimeUnit.SECONDS,
            new LinkedBlockingQueue<>(1000),
            Executors.defaultThreadFactory(),
            new ThreadPoolExecutor.CallerRunsPolicy() // 队列满时让调用线程兜底处理,防止任务被拒绝
    );

    private final BarService barService;
    private final BooleanResponseClient booleanResponseClient;

    public FooProcessor(BarService barService, BooleanResponseClient booleanResponseClient) {
        this.barService = barService;
        this.booleanResponseClient = booleanResponseClient;
    }

    public void startProcess(List<Foo> fooDetails) throws Exception {
        // 用ConcurrentLinkedQueue做高并发写入场景下的线程安全收集容器,比CopyOnWriteArrayList性能更优
        ConcurrentLinkedQueue<String> failedFooNames = new ConcurrentLinkedQueue<>();

        List<CompletableFuture<Void>> futures = fooDetails.stream()
                .map(fooDetail -> processFoo(fooDetail, failedFooNames))
                .collect(Collectors.toList());

        // 等待所有任务完成,超时时间根据业务接口响应速度调整
        CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).get(30, TimeUnit.MINUTES);

        // 优雅关闭线程池:先停止接受新任务,等待现有任务完成;超时未完成则强制终止
        EXECUTOR_SERVICE.shutdown();
        if (!EXECUTOR_SERVICE.awaitTermination(1, TimeUnit.MINUTES)) {
            EXECUTOR_SERVICE.shutdownNow();
        }

        if (!failedFooNames.isEmpty()) {
            throw new SyncException("Foo同步失败,失败条目: " + failedFooNames);
        }
    }

    private CompletableFuture<Void> processFoo(Foo fooDetail, ConcurrentLinkedQueue<String> failedFooNames) {
        return CompletableFuture.supplyAsync(() -> {
            // 修复原代码拼写错误:foodetail -> fooDetail
            Bar request = barService.createBarFromDB(Collections.singletonList(fooDetail));
            boolean pushFailed = booleanResponseClient.updateBar(request);
            return pushFailed ? fooDetail.getName() : null;
        }, EXECUTOR_SERVICE)
        // 处理业务逻辑失败的情况
        .thenAcceptAsync(failedName -> {
            if (Objects.nonNull(failedName)) {
                failedFooNames.add(failedName);
            }
        }, EXECUTOR_SERVICE)
        // 处理执行过程中抛出的所有异常
        .exceptionally(ex -> {
            failedFooNames.add(fooDetail.getName());
            // 可在此添加异常日志,方便排查具体失败原因
            // log.error("处理Foo[{}]失败", fooDetail.getName(), ex);
            return null;
        });
    }
}

关键优化点说明

  • 自定义线程池:通过ThreadPoolExecutor手动配置核心线程数、最大线程数、任务队列等参数,避免默认线程池在大列表场景下的资源失控问题,CallerRunsPolicy策略能防止任务被拒绝。
  • 线程安全收集:使用ConcurrentLinkedQueue存储失败名称,适配高并发写入场景,性能优于CopyOnWriteArrayList。
  • 统一结果与异常处理:用thenAcceptAsync处理业务失败、exceptionally捕获执行异常,统一收集失败名称,替代原代码中繁琐的future.get()异常捕获逻辑。
  • 线程池优雅关闭:任务完成后按步骤关闭线程池,避免线程资源泄漏。
  • 代码细节修正:修复原代码中变量拼写错误,优化集合创建写法。

额外提示

  • 若使用Java 8及以下版本,将Collections.singletonList()保留即可(Java 9+推荐用List.of())。
  • 线程池参数需根据实际业务的QPS、接口响应时间调整,IO密集型任务可适当增加核心线程数。
  • 建议在异常处理分支添加日志,方便定位具体失败的Foo条目及原因。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 15:18:17