CompletionService 通过 ExecutorService 与 BlockingQueue 结合实现任务提交与结果获取解耦,支持按完成顺序处理异步任务;常用实现为 ExecutorCompletionService,需注意异常处理、及时 shutdown、优先使用 Callable、结果处理异步化及采用 poll 超时机制。

CompletionService 是 Java 并发编程中一个非常实用的工具,它把任务提交和结果获取解耦,让你能按“谁先完成谁先处理”的顺序获取异步任务结果,特别适合批量提交耗时差异大的任务(比如调用多个外部 API、查询不同数据库等)。
核心思路:用 ExecutorService + BlockingQueue 包装任务结果
CompletionService 本身是个接口,常用实现是 ExecutorCompletionService。它内部持有一个线程池(ExecutorService)和一个阻塞队列(通常是 LinkedBlockingQueue),当任务执行完,会把 Future 封装成 Future<T> 放入队列。你只需从队列里 take() 或 poll(),就能拿到最先完成的任务结果,无需轮询或按提交顺序等待。
基本使用步骤
- 创建线程池(推荐用
Executors.newFixedThreadPool(n)或更可控的ThreadPoolExecutor) - 用该线程池构造
ExecutorCompletionService<String>(泛型为任务返回类型) - 循环提交
Callable任务(用submit()),不关心返回的 Future - 调用
take()阻塞获取首个完成结果,或poll(timeout, unit)带超时获取 - 对每个取到的
Future调用get()拿实际返回值(此时已确定完成,不会阻塞)
一个典型示例
假设要并发请求 5 个不同响应时间的 HTTP 接口:
ExecutorService executor = Executors.newFixedThreadPool(3);
CompletionService<String> cs = new ExecutorCompletionService<>(executor);
<p>// 提交任务
for (int i = 0; i < 5; i++) {
final int taskId = i;
cs.submit(() -> {
Thread.sleep((long) (Math.random() * 2000 + 500)); // 模拟不同耗时
return "Result-" + taskId;
});
}</p><p>// 按完成顺序获取结果
for (int i = 0; i < 5; i++) {
try {
String result = cs.take().get(); // take() 阻塞直到有结果
System.out.println("Got: " + result);
} catch (InterruptedException | ExecutionException e) {
Thread.currentThread().interrupt();
break;
}
}
executor.shutdown();注意事项与优化点
-
不要忽略异常:
Future.get()可能抛出ExecutionException,需捕获并检查getCause() -
及时 shutdown:任务全部取完后记得调用
executor.shutdown(),避免线程泄漏 - 慎用 submit(Runnable, T):如果用 Runnable 提交,返回值固定,无法体现任务差异;优先用 Callable
- 结果处理可并行化:取到结果后若还需耗时处理(如写 DB、发消息),建议另起线程或用独立线程池,避免阻塞 CompletionService 的消费线程
-
超时控制更安全:生产环境建议用
poll(long, TimeUnit)替代take(),防止因某个任务卡死导致整体阻塞


















