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

处理共享列表并提交结果的并发构造选型咨询

针对该场景的优化并发方案

首先先修正你提供的原始代码里的明显错误(变量名混淆、类型不匹配、方法调用错误),修正后的基础代码如下:

class Foo {
    private int len;
    public Foo(int l) { this.len = l; } // 修正构造函数的变量赋值错误
    public void process(List<Bar> list) {
        int start = 0;
        while(start < list.size()) {
            int end = Math.min(start + len, list.size()); // 修正原代码中N的混淆问题
            List<Bar> sublist = list.subList(start, end);
            handle(sublist); // 原代码误写为process,改为调用handle
            start = end;
        }
    }
  
    private void handle(List<Bar> sublist) {
        // 耗时处理逻辑
    }
}

针对并行化需求,推荐以下几种更简洁高效的并发构造:

1. 使用CompletableFuture实现异步并行处理

CompletableFuture可以更优雅地处理异步任务的提交与结果聚合,无需手动管理线程等待,还能灵活指定线程池:

import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.stream.Collectors;
import java.util.stream.IntStream;

class Foo {
    private int len;
    private final ExecutorService executor;

    public Foo(int l, ExecutorService executor) {
        this.len = l;
        this.executor = executor;
    }

    public void process(List<Bar> list) {
        // 拆分列表并提交异步任务
        List<CompletableFuture<Void>> futures = IntStream.range(0, (list.size() + len - 1) / len)
                .mapToObj(i -> {
                    int start = i * len;
                    int end = Math.min(start + len, list.size());
                    List<Bar> sublist = list.subList(start, end);
                    return CompletableFuture.runAsync(() -> handle(sublist), executor);
                })
                .collect(Collectors.toList());

        // 等待所有任务完成
        CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
    }

    private void handle(List<Bar> sublist) {
        // 耗时处理逻辑
    }
}

如果需要收集处理结果,可改用supplyAsync()替代runAsync(),后续通过allOf()统一获取结果,避免使用需要同步的共享列表。

2. 使用并行流(Stream.parallel())

如果不需要自定义线程池,且任务间无状态依赖,并行流是最简洁的方案:

import java.util.List;
import java.util.stream.IntStream;

class Foo {
    private int len;

    public Foo(int l) {
        this.len = l;
    }

    public void process(List<Bar> list) {
        // 拆分列表为子列表,并行处理
        IntStream.range(0, (list.size() + len - 1) / len)
                .mapToObj(i -> {
                    int start = i * len;
                    int end = Math.min(start + len, list.size());
                    return list.subList(start, end);
                })
                .parallel() // 启用并行处理
                .forEach(this::handle);
    }

    private void handle(List<Bar> sublist) {
        // 耗时处理逻辑
    }
}

注意:并行流默认使用ForkJoinPool.commonPool(),若应用有其他并行任务,可能互相影响,此时优先选择自定义线程池的方案。

3. ExecutorService+CountDownLatch(你的原始方案优化版)

如果坚持用Executor提交任务的思路,CountDownLatch比手动等待更可靠,且能避免任务异常导致的等待阻塞:

import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;

class Foo {
    private int len;
    private final ExecutorService executor;

    public Foo(int l, ExecutorService executor) {
        this.len = l;
        this.executor = executor;
    }

    public void process(List<Bar> list) throws InterruptedException {
        int batchCount = (list.size() + len - 1) / len;
        CountDownLatch latch = new CountDownLatch(batchCount);

        for (int i = 0; i < batchCount; i++) {
            int start = i * len;
            int end = Math.min(start + len, list.size());
            List<Bar> sublist = list.subList(start, end);
            
            executor.submit(() -> {
                try {
                    handle(sublist);
                } finally {
                    latch.countDown(); // 无论任务成功失败,都确保计数减一
                }
            });
        }

        latch.await(); // 等待所有任务完成
    }

    private void handle(List<Bar> sublist) {
        // 耗时处理逻辑
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 18:10:26