批量处理相似任务是后端开发里非常典型的一类场景:一次请求要调用几十个外部接口、定时任务要同步上千条数据、消息消费要批量落库等等。如果简单地写个 for 循环串行执行,每个任务耗时 200 毫秒,100 个任务就是 20 秒,接口基本没法用;反过来如果每个任务都 new 一个线程,任务量上来之后线程创建销毁的开销和内存占用又会拖垮整个服务。这篇文章就来梳理 Java 中处理这类问题的几种主流方案,从写法到原理都过一遍,顺便聊聊各自的适用场景和容易踩的坑。

一、线程池 + CountDownLatch:最经典的批量方案
这是很多老项目里最常见的写法。核心思路是把所有任务一次性提交给线程池,主线程用 CountDownLatch 等待全部完成。它的好处是代码直观、可控性强,任务数量和线程数量都能明确掌握。
List<Callable<Result>> tasks = new ArrayList<>();
for (Item item : items) {
tasks.add(() -> processOne(item));
}
ExecutorService pool = new ThreadPoolExecutor(
8, 16, 60, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(500),
new ThreadPoolExecutor.CallerRunsPolicy());
List<Future<Result>> futures = pool.invokeAll(tasks);
for (Future<Result> f : futures) {
results.add(f.get());
}用 invokeAll 比 submit 加 CountDownLatch 更简洁一些,语义也完全一致:阻塞直到所有任务结束。需要注意的一个坑是拒绝策略,默认的 AbortPolicy 会直接抛异常,生产环境建议换成 CallerRunsPolicy,让提交任务的线程自己执行,相当于天然的限流降级,避免任务被丢弃。
这种方案的问题在于它是“整体等待”模式,必须所有任务都跑完才能继续,如果其中某个任务特别慢,整体耗时就被它拖住了。另外拿结果的循环里调用 f.get() 如果不设置超时,一个任务卡死就会让主线程永久阻塞,这点在对接不靠谱的外部接口时要格外小心。
二、CompletableFuture:更灵活的异步编排
JDK 8 之后,CompletableFuture 成了批量任务的主流选择。它支持回调式拿结果、单个任务超时控制(JDK 9 的 orTimeout)、以及 allOf 聚合等能力,写起来更符合异步思维。
List<CompletableFuture<Result>> futures = items.stream()
.map(item -> CompletableFuture
.supplyAsync(() -> processOne(item), pool)
.orTimeout(3, TimeUnit.SECONDS)
.exceptionally(e -> Result.fallback(item, e)))
.collect(Collectors.toList());
CompletableFuture
.allOf(futures.toArray(new CompletableFuture[0]))
.thenRun(() -> {
List<Result> results = futures.stream()
.map(CompletableFuture::join)
.collect(Collectors.toList());
callback.accept(results);
});这段代码里有三个细节值得注意。第一,supplyAsync 一定要显式传入自定义线程池,否则默认用的是 ForkJoinPool.commonPool,它是全 JVM 共享的,批量任务把它占满后,其他使用并行流的代码都会被波及,这是线上事故的高发点。第二,exceptionally 保证了单个任务失败不会导致整体失败,可以逐条做降级或兜底。第三,orTimeout 给每个任务加上了独立的超时上限,慢任务不会拖垮整批。
和线程池方案相比,CompletableFuture 的优势在于控制粒度更细,可以做到“失败隔离、单任务超时、非阻塞回调”。如果任务是 IO 密集型的(比如 HTTP 调用),它几乎是首选;如果任务之间还有依赖关系,比如先查再算再写,它的链式组合能力也能把逻辑写得比嵌套回调清晰得多。
三、分片批处理:应对超大任务量
当任务量达到几万甚至几十万级别时,一次性提交所有任务会带来两个问题:内存里要同时持有全部任务对象,线程池队列也可能被撑爆。这时候更稳妥的做法是分片(分批)处理,每批固定大小,处理完一批再提交下一批。
int batchSize = 200;
List<List<Item>> partitions = Lists.partition(items, batchSize);
for (List<Item> batch : partitions) {
List<CompletableFuture<Result>> futures = batch.stream()
.map(item -> CompletableFuture.supplyAsync(() -> processOne(item), pool))
.collect(Collectors.toList());
List<Result> batchResults = futures.stream()
.map(CompletableFuture::join)
.collect(Collectors.toList());
saveBatch(batchResults); // 每批落库,及时释放内存
}分片的核心价值不只是控制并发数,更重要的是把“整体成功或整体失败”拆成了“每批独立提交”,中间任何一批出问题都可以单独重试或记录,不影响已完成的批次。对于批量写数据库的场景,配合 MyBatis 的批量插入或 JDBC 的 rewriteBatchedStatements=true 参数,性能往往能提升一个数量级。
批次大小的选择没有万能值,一般经验是:外部接口类任务每批 50 到 200 比较合适(避免触发对方限流),数据库批量写入可以到 500 到 1000。建议通过压测确定,观察线程池活跃线程数、队列堆积情况和下游响应时间,找到系统吞吐的拐点。
四、生产者消费者模式:应对持续流入的任务
前面几种方案都是“有一批已知任务,处理完就结束”的模式。如果任务是持续不断产生的,比如消息队列消费、日志采集上报,更适合用生产者消费者模型,用一个有界阻塞队列做缓冲,固定数量的消费者线程循环取任务。
BlockingQueue<Task> queue = new LinkedBlockingQueue<>(1000);
// 消费者:固定线程数,持续拉取
for (int i = 0; i < 8; i++) {
pool.submit(() -> {
while (!stopped.get()) {
try {
Task task = queue.poll(1, TimeUnit.SECONDS);
if (task != null) {
processOne(task);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
});
}这种模式的好处是削峰填谷:任务产生的速度波动会被队列吸收,消费者始终以稳定的速率处理。有界队列是关键,容量满了之后生产者被阻塞,天然起到背压作用,防止内存无限增长。消费线程里一定要正确处理 InterruptedException,恢复中断标志并退出循环,否则优雅停机会失效,进程关不掉或者任务处理到一半被强杀。
五、选型建议与常见坑
总结一下怎么选:任务量小且固定,直接 parallelStream 或线程池 invokeAll 就够;任务量大、需要失败隔离和超时控制,用 CompletableFuture;任务量超大或涉及批量落库,加分片;任务是持续流入的长驻场景,用生产者消费者加阻塞队列。
几个高频踩坑点也列一下:一是线程池参数按 CPU 核数或任务类型区分,IO 密集型可以开到核数的两倍以上,CPU 密集型则接近核数即可,盲目照搬公式不如压测;二是线程池一定要命名并监控,方便排查问题时看线程堆栈;三是批量任务里共享的可变状态要特别注意,优先用不可变对象或 ConcurrentHashMap,避免隐式的竞态条件;四是所有对外调用的任务都要设置超时,无论是 HTTP 客户端的连接超时还是 Future.get 的等待超时,“没有超时的远程调用”是稳定性事故里最常见的原因之一。
把这几套方案和对应的细节掌握扎实,大部分批量任务场景都能写出既快又稳的代码。真正的功力不在于用了多新的 API,而在于对任务特征的理解和对应并发模型的匹配。
Java并发编程线程池CompletableFuture修改时间:2026-09-13 09:12:35