导读:本期聚焦于胡建平创作的《怎么利用 Phaser 灵活处理参与者数量动态变化的阶段性并发计算任务》,敬请观看详情。阶段性并发计算中最棘手的情况之一,是任务运行到某个同步点时,参与线程的数量发生了变化。固定参与方的 CyclicBarrier 和一次性 CountDownLatch 在这种情况下会变得非常别扭:要么提前注册一批根本不会执行的线程,要么在任务中途重新构建屏障。Phaser 提供了更自然的解决方案,它把同步点建模为阶段,并通过 register、bulkRegister、arriveAndDeregister 等方法允许参与方动态增减。线程可以在任意阶段调用 register 注册进入,也可以在完成自己的阶段任务后通过 arriveAndDeregister 退出,不会阻塞其他线程。本文以一个动态工人加入和退出的分阶段计算场景为例,说明 Phaser 的核心 API 用法、阶段推进机制以及若干容易踩到的坑,帮助你在多阶段并发任务中灵活组织线程同步。

在 Java 并发编程里,Phaser 是一个经常被忽略但能力很强的同步工具。CountDownLatch 只能倒计时一次,CyclicBarrier 虽然可以复用,但参与方数量在创建后无法调整。Phaser 的出现正是为了解决这类限制。它把一次同步看作一个阶段,参与线程可以在阶段之间动态注册或注销,因此特别适合那些任务数量会随阶段变化的计算场景,比如分批处理、流水线作业、并行搜索中的剪枝任务等。理解 Phaser 的机制之后,你就能少写很多为了控制线程同步而额外维护的计数器或锁。

怎么利用 Phaser 灵活处理参与者数量动态变化的阶段性并发计算任务

Phaser 与固定参与方屏障的本质区别

CountDownLatch 的计数在初始化时确定,每个线程调用 countDown 使计数减一,主线程调用 await 等待计数归零。这个结构适合一次性的启动或完成通知,但无法应对多次同步。CyclicBarrier 引入了可复用的屏障,所有线程到达后一起释放,但它要求 parties 数量必须固定。如果某个线程中途退出,其他线程会一直等待;如果想加入新线程,则必须等下一轮并且重新计算 parties,非常不便。

Phaser 的核心模型是阶段和参与方。一个 Phaser 内部维护一个当前阶段号,初始为 0。调用 arriveAndAwaitAdvance 的线程表示自己已经到达当前阶段,并阻塞等待其他参与方到达。当所有已注册的参与方都到达后,阶段号加一,所有等待线程继续执行。与 CyclicBarrier 不同,Phaser 的参与方数量是可变的。通过 register 可以增加一方,通过 arriveAndDeregister 可以在到达的同时注销一方。这个能力意味着下一阶段的等待方数量会自动减少,不需要外部干预。

还有一个重要方法叫 bulkRegister,可以一次注册多个参与方。Phaser 的构造器可以传入初始参与方数量,也可以传入一个父 Phaser 形成层级结构。不过在大多数单独使用场景中,我们更关注 register、arriveAndAwaitAdvance、arriveAndDeregister 以及 awaitAdvance 这几个方法。它们共同构成了动态阶段性同步的基础。

动态加入与退出:register 和 arriveAndDeregister 的配合

一个常见做法是把主线程也注册到 Phaser 中。例如使用 new Phaser(1) 创建时先注册一个参与方,通常称为控制方或主控线程。这样做的好处是可以避免所有工作线程都注销后 Phaser 提前进入终止状态。主线程可以在适当的时候调用 arriveAndAwaitAdvance 参与阶段同步,或者调用 arriveAndDeregister 结束整个 Phaser。

新的任务线程启动前调用 register,这样 Phaser 的已注册参与方数量加一。线程任务执行到一个阶段边界时,调用 arriveAndAwaitAdvance 等待其他线程。如果某个线程的任务较轻,提前完成并且不再参与后续阶段,可以调用 arriveAndDeregister,表示到达当前阶段并注销。Phaser 会在内部减少参与方数量,下一阶段不再等待这个线程。

下面的代码展示了动态注册和注销的基本流程。初始只有一个主线程注册在 Phaser 中,随后三个工作线程依次注册并开始执行阶段任务。每个工作线程完成阶段一后直接注销,主线程等待它们全部到达后结束。

import java.util.concurrent.Phaser;

public class DynamicRegisterDemo {
    public static void main(String[] args) {
        Phaser phaser = new Phaser(1); // 初始注册主线程,避免提前终止
        System.out.println("初始阶段: " + phaser.getPhase());

        for (int i = 0; i < 3; i++) {
            phaser.register(); // 每个工作线程注册一方
            final int workerId = i;
            new Thread(() -> {
                System.out.println("Worker " + workerId + " 进入阶段 0");
                phaser.arriveAndAwaitAdvance(); // 到达并等待其他线程
                System.out.println("Worker " + workerId + " 通过阶段 0");
                phaser.arriveAndDeregister(); // 到达并注销,不参与后续阶段
            }).start();
        }

        // 主线程也到达阶段 0,所有参与方到齐后阶段推进
        int phase = phaser.arriveAndAwaitAdvance();
        System.out.println("所有 Worker 已通过阶段 0,当前阶段: " + phase);

        // 主线程注销,Phaser 没有参与方后终止
        phaser.arriveAndDeregister();
        System.out.println("Phaser 是否终止: " + phaser.isTerminated());
    }
}

这个例子中,Phaser 初始参与方为 1,即主线程。三个工作线程注册后参与方变成 4。四个线程都调用 arriveAndAwaitAdvance 后,阶段从 0 推进到 1,所有线程继续。接着工作线程调用 arriveAndDeregister,主线程也调用一次 arriveAndDeregister,Phaser 的参与方变为 0,进入终止状态。如果把主线程的初始注册去掉,三个工作线程注销后 Phaser 可能提前终止,导致主线程的等待调用行为不符合预期。

阶段性并发计算任务的完整示例

假设有一个计算任务需要分两阶段执行:第一阶段每个工作线程处理自己的数据分片并计算局部结果;第二阶段由一个汇总线程根据所有局部结果生成最终结果。同时任务队列中还会动态补充新的数据分片,因此新的工作线程会在第一阶段结束后注册进来,继续下一轮第一阶段。旧的工作线程如果已经处理完自己的分片,则可以在第一阶段结束时注销,不再参与汇总。

这种结构很适合用 Phaser 表达。所有工作线程和主控线程都注册到同一个 Phaser。第一阶段结束时,所有线程都调用 arriveAndAwaitAdvance,确保局部结果全部计算完成。随后,需要退出的线程调用 arriveAndDeregister,但如果是新来的线程,它们在进入下一轮之前先 register。需要注意 register 的时机:如果某个线程在其他线程已经到达阶段边界后才注册,Phaser 的参与方数量会在当前阶段变化,但已经等待的线程可能不会重新检查该数量。一般建议在阶段开始前完成注册,避免阶段推进过程中的竞态。

import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Phaser;

public class StagedComputeDemo {
    private static final List<Integer> partialResults = new ArrayList<>();

    public static void main(String[] args) throws InterruptedException {
        Phaser phaser = new Phaser(1); // 主控线程注册

        for (int round = 0; round < 2; round++) {
            int workerCount = 2 + round; // 每轮增加一个工作线程
            for (int i = 0; i < workerCount; i++) {
                phaser.register();
                final int workerId = i;
                new Thread(() -> {
                    int part = workerId * 10;
                    synchronized (partialResults) {
                        partialResults.add(part);
                    }
                    System.out.println("Worker " + workerId + " 完成阶段一,局部结果: " + part);
                    phaser.arriveAndAwaitAdvance(); // 等待所有工作线程完成阶段一
                    phaser.arriveAndDeregister();   // 本线程不再参与下一阶段
                }).start();
            }

            // 主控线程等待所有工作线程完成阶段一并汇总
            phaser.arriveAndAwaitAdvance();
            int total = partialResults.stream().mapToInt(Integer::intValue).sum();
            System.out.println("第 " + round + " 轮汇总结果: " + total);
            partialResults.clear();

            // 主控线程继续保留注册,进入下一轮前仍然是一方
        }

        phaser.arriveAndDeregister();
        System.out.println("任务完成,Phaser 是否终止: " + phaser.isTerminated());
    }
}

示例中 partialResults 使用同步块保护,因为多个工作线程可能同时写入。Phaser 本身只负责线程的到达和等待,不提供数据可见性保证,所以共享数据的同步仍然需要额外的锁或并发容器。这一点需要和 CyclicBarrier 类似,屏障释放只能说明所有线程都已到达,但并不自动保证每个线程写入的数据对其他线程可见。好在屏障本身的 happens-before 关系可以帮助建立可见性,但前提是没有数据竞争。

每一轮完成阶段一后,主控线程汇总。工作线程通过 arriveAndDeregister 退出,下一轮开始前主控线程再注册新的 worker。由于主控线程一直保留注册,Phaser 不会因为中间没有工作线程而终止。这个模式可以持续执行多轮,每轮的参与方数量都可能不同。

阶段推进与线程等待的底层行为

Phaser 内部通过一个状态变量同时保存阶段号和参与方数量,使用 CAS 操作保证并发修改的正确性。每次调用 arrive 或 arriveAndDeregister 都会使未到达参与方计数减少。当未到达计数减到 0 时,说明当前阶段所有参与方都已经到达,Phaser 会推进阶段号,并唤醒所有在 awaitAdvance 或 arriveAndAwaitAdvance 上等待的线程。

这里有一个容易被误解的地方:register 可以在任何阶段调用,但它增加的是未到达计数,而不是已到达计数。如果在一个阶段推进过程中注册,新的参与方会被计入当前阶段或下一阶段,具体取决于注册时是否所有老参与方都已经到达。Phaser 的 API 并没有像 CyclicBarrier 那样在突破屏障后就自动重置参与方,而是保留注册状态。所以如果希望在下一阶段减少等待方,必须调用 arriveAndDeregister 或主动 deregister。

awaitAdvance 系列方法可以指定要等待的阶段号。当线程调用 awaitAdvance(phase) 时,如果当前阶段不是传入的 phase,会立即返回;否则会阻塞直到阶段号变化。这个机制允许一些控制线程不直接参与任务,而是观察某个阶段是否完成。比如主控线程想在阶段 0 完成后再注册新的 worker,就可以调用 int p = phaser.arriveAndAwaitAdvance() 等待阶段推进。

使用 Phaser 的常见注意事项

第一,要小心 Phaser 提前终止。Phaser 在参与方数量为零时不会自动终止,只有当 arriveAndDeregister 导致参与方数量降到 0 时,才会进入终止状态。如果所有工作线程都注销了,主控线程没有注册,Phaser 会终止,后续 register 会直接返回负数,无法再恢复。因此建议创建 Phaser 时先注册一个控制方,确保持续存在。

第二,register 与 arrive 之间存在竞态窗口。虽然 Phaser 是线程安全的,但如果一个线程在阶段推进后立即 register,而主控线程假设新线程会在当前阶段参与同步,可能会造成阶段计数错误。最好把 register 放在阶段开始之前,并且避免在多个线程同时修改参与方数量时依赖具体的阶段号。

第三,Phaser 原生不提供带超时的等待,但提供了 awaitAdvanceInterruptibly 和 awaitAdvance,如果需要超时,可以用 awaitAdvance 配合循环检查或使用类似 Future 的方式。Phaser 也不像 CyclicBarrier 那样支持一个每轮执行的屏障动作,不过可以重写 onAdvance 方法,在阶段推进前执行自定义逻辑。onAdvance 的返回值决定 Phaser 是否终止,默认在参与方为零时终止。

第四,任务数量动态变化时,最好把注册和注销动作集中在一个控制线程中处理,或者使用 ExecutorService 配合明确的提交边界,避免工作线程自己随意注册引发难以追踪的阶段错位。Phaser 虽然灵活,但越灵活的工具越需要清晰的生命周期管理。

Phaser并发编程阶段性任务修改时间:2026-09-19 10:52:32

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