在并发编程中,我们经常需要向线程池提交一批异步任务,然后逐个获取它们的结果。如果直接使用 Future.get(),你会发现一个明显的问题:结果是按照调用 get 的顺序返回的,而不是按照任务实际完成的顺序。假设第一个任务需要执行 10 秒,后面的任务各只需 1 秒,那么即使后面九个任务早就执行完了,你也必须干等第一个任务结束才能拿到第一个结果,白白浪费了大量等待时间。JDK 提供的 ExecutorCompletionService 就是专门用来解决这个问题的,它能把先完成的任务结果先交到你手上,实现真正意义上的按完成顺序处理结果。

ExecutorCompletionService 的核心原理
ExecutorCompletionService 位于 java.util.concurrent 包下,它实现了 CompletionService 接口。它的设计思路非常巧妙:内部维护了一个阻塞队列,每当一个任务完成时,就把代表该任务的 Future 对象放入这个队列中。这样一来,队列头的元素永远是先完成的任务,调用者只需要从队列中取结果即可。
具体来看,ExecutorCompletionService 在提交任务时并不会直接把原始任务丢给线程池,而是把它包装成一个 QueueingFuture 对象。这个包装类重写了 done() 方法,在任务执行结束(无论成功还是异常)时,自动把当前 Future 塞入完成队列。这个细节很关键:即使是抛出异常的任务,同样会进入队列,因此调用方在 take() 之后调用 get() 时仍可能拿到 ExecutionException,需要做好异常捕获。
它的构造函数有两个重载:
// 使用指定的线程池,内部创建一个无界的 LinkedBlockingQueue
ExecutorCompletionService<T> ecs =
new ExecutorCompletionService<>(executor);
// 使用指定的线程池和自定义的完成队列
ExecutorCompletionService<T> ecs =
new ExecutorCompletionService<>(executor, new LinkedBlockingQueue<>());第二个重载允许你自定义队列,例如换成有界队列或者 PriorityBlockingQueue,在某些特殊场景下非常有用,比如希望按任务的优先级而非完成时间来消费结果。
基本用法:按完成顺序处理结果
下面通过一个完整的例子演示它的典型用法。我们模拟三个耗时不相同的任务,耗时最短的任务最先提交结果,主线程按照完成顺序依次打印,而不是按照提交顺序。
import java.util.concurrent.*;
public class CompletionServiceDemo {
public static void main(String[] args) throws Exception {
ExecutorService executor = Executors.newFixedThreadPool(3);
ExecutorCompletionService<String> completionService =
new ExecutorCompletionService<>(executor);
// 提交三个任务,耗时分别为 3 秒、1 秒、2 秒
completionService.submit(() -> {
Thread.sleep(3000);
return "任务A(耗时3秒)";
});
completionService.submit(() -> {
Thread.sleep(1000);
return "任务B(耗时1秒)";
});
completionService.submit(() -> {
Thread.sleep(2000);
return "任务C(耗时2秒)";
});
// 按完成顺序获取结果:B -> C -> A
for (int i = 0; i < 3; i++) {
Future<String> future = completionService.take();
String result = future.get();
System.out.println("收到结果:" + result);
}
executor.shutdown();
}
}运行上面的代码,输出顺序一定是任务B、任务C、任务A,与提交顺序无关。这里有两个关键方法需要区分清楚:take() 和 poll()。take() 是阻塞式的,如果当前没有已完成的任务,它会一直等待,直到有任务完成;poll() 是非阻塞式的,队列 为空时立即返回 null。此外 poll(long timeout, TimeUnit unit) 提供了带超时的等待方式,在需要对整体处理时间做限制的场景中非常实用。
需要注意的一点是,ExecutorCompletionService 本身并不管理线程池的生命周期,它只是线程池的一个装饰器。所以 shutdown() 依然要调用原始的 ExecutorService 对象,这一点初学者容易忽略。
与 Future.get 和 invokeAll 的对比
最传统的方式是保存所有 Future 然后依次调用 get()。这种方式的问题在于它强制按提交顺序等待结果。假设有 100 个任务,第一个任务卡住 30 秒,其余任务 1 秒内全部完成,那么后续结果的消费会被整体拖慢 30 秒,这在响应延迟敏感的系统中是不可接受的。
另一种常见方案是 invokeAll。它会阻塞到所有任务完成才返回,虽然拿到的是完成状态的 Future 列表,但无法做到边完成边处理,整体吞吐依然受最慢任务的制约。而 ExecutorCompletionService 的优势恰恰在于流式处理:任务一完成就能立刻被消费,整体处理时间接近最快任务的时间加上消费开销,而不是最慢任务的时间。
三者的对比如下表所示:
| 方案 | 结果顺序 | 是否支持流式处理 | 适用场景 |
|---|---|---|---|
| Future.get 依次获取 | 提交顺序 | 否 | 任务耗时均匀且数量少 |
| invokeAll | 提交顺序 | 否 | 必须等全部结果到齐才能继续 |
| ExecutorCompletionService | 完成顺序 | 是 | 批量异步任务、结果实时消费 |
进阶实践:超时控制与异常处理
在实际项目中,批量任务往往需要设置整体超时,而且个别任务失败不应该影响其他任务的结果处理。下面的例子演示了如何结合 poll 的超时能力和异常捕获来构建一个健壮的批量处理逻辑。
import java.util.concurrent.*;
import java.util.*;
public class RobustBatchDemo {
public static void main(String[] args) throws Exception {
ExecutorService executor = Executors.newFixedThreadPool(4);
ExecutorCompletionService<Integer> cs =
new ExecutorCompletionService<>(executor);
List<Future<Integer>> futures = new ArrayList<>();
for (int i = 1; i <= 6; i++) {
final int id = i;
futures.add(cs.submit(() -> {
Thread.sleep(1000 + id * 500L);
if (id == 4) {
throw new RuntimeException("任务" + id + "执行失败");
}
return id * 100;
}));
}
int total = 0;
long deadline = System.currentTimeMillis() + 8000; // 整体超时 8 秒
for (int i = 0; i < futures.size(); i++) {
long wait = deadline - System.currentTimeMillis();
Future<Integer> f = cs.poll(wait, TimeUnit.MILLISECONDS);
if (f == null) {
System.out.println("整体超时,剩余任务取消");
futures.forEach(x -> x.cancel(true));
break;
}
try {
total += f.get();
} catch (ExecutionException e) {
System.out.println("某任务异常:" + e.getCause().getMessage());
}
}
System.out.println("成功结果累加值:" + total);
executor.shutdown();
}
}这段代码有几个值得注意的细节。第一,我们保存了所有 Future 的引用,这是为了在整体超时时能够主动 cancel 尚未完成的任务,避免线程池中的工作线程被白白占用。第二,异常处理放在了消费侧:单个任务失败只跳过该结果,其余任务照常累加,保证了批量操作的局部容错能力。第三,超时时间通过 deadline 统一计算,而不是每次固定等待若干秒,这样整体耗时才能被精确控制。
还有一个容易踩的坑:完成队列默认是无界的,如果只提交任务却从不消费结果,队列会不断堆积对象,任务数量极大时可能造成内存压力。因此务必保证消费速率与生产速率匹配,或者在自定义队列时使用有界队列并配合合理的拒绝策略。
总结一下,ExecutorCompletionService 通过完成队列解耦了任务提交与结果消费,让结果按照实际完成顺序流式输出,显著降低了批量异步任务的等待浪费。只要掌握 take 与 poll 的语义差异、注意异常传播和队列消费的平衡,就能在接口聚合、并行计算、批量 IO 等场景中写出高性能且健壮的并发代码。
ExecutorCompletionServiceJava并发编程CompletionService修改时间:2026-09-01 18:34:37