导读:本期聚焦于吴凌云创作的《如何在 Java 中使用 ExecutorCompletionService 按照异步任务完成顺序获取返回结果》,敬请观看详情。提交了一批异步任务之后,怎样才能让先执行完的任务结果先被处理,而不是按照提交顺序死等?ExecutorCompletionService 正是解决这个问题的利器。它内部通过一个阻塞队列把已完成任务的 Future 缓存起来,调用 take 方法时总能拿到最先完成任务的结果。本文详细讲解它的底层实现原理、take 与 poll 的区别、配合线程池的典型用法,并对比 Future 的 get 方式与 invokeAll 方案的差异,同时给出批量任务超时控制和异常处理的完整示例,帮助你写出更高吞吐、更低等待浪费的并发代码。

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

如何在 Java 中使用 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

免责声明:​ 已尽一切努力确保本网站所含信息的准确性。网站内容多为原创整理与精心编撰,观点力求客观中立。本站旨在免费分享,内容仅供个人学习、研究或参考使用。若引用了第三方作品,版权归原作者所有。如内容涉及您的权益,请联系我们处理。
内容垂直聚焦
专注技术核心技术栏目,确保每篇文章深度聚焦于实用技能。从代码技巧到架构设计,为用户提供无干扰的纯技术知识沉淀,精准满足专业提升需求。
知识结构清晰
覆盖从开发到部署的全链路。AI、前端、编程、数据库、服务器、建站、系统层层递进,构建清晰学习路径,帮助用户系统化掌握开发与运维所需的核心技术。
深度技术解析
拒绝泛泛而谈,深入技术细节与实践难点。无论是数据库优化还是服务器配置,均结合真实场景与代码示例进行剖析,致力于提供可直接应用于工作的解决方案。
专业领域覆盖
精准对应开发生命周期。从前端界面到后端编程,从数据库操作到服务器运维,形成完整闭环,一站式满足全栈工程师和运维人员的技术需求。
即学即用高效
内容强调实操性,步骤清晰、代码完整。用户可根据教程直接复现和应用于自身项目,显著缩短从学习到实践的距离,快速解决开发中的具体问题。
持续更新保障
专注既定技术方向进行长期、稳定的内容输出。确保各栏目技术文章持续更新迭代,紧跟主流技术发展趋势,为用户提供经久不衰的学习价值。