上一篇的CyclicBarrier适合同一组工作者反复会合。再往前走一步:如果准备阶段需要三位工作者,后续加工只需要其中两位,怎样让队伍在阶段之间自然调整?今天认识Phaser,把注册、到达和退出放进一个完整例子里。🧩
本文以JDK21为准,Phaser自Java7提供;官方资料核对日期为2026年9月23日。例子使用虚构的批处理流程,不涉及真实业务数据。
一、先看结构:阶段与参与份额
把一次批处理拆成准备、加工两个阶段。每个仍在队伍中的工作者,本阶段完成后报到;当前阶段需要的份额全部报到,阶段才会推进。阶段编号从0开始,而注册份额可以随着流程变化。JDK21的Phaser说明
Phaser只管理计数,并不保存“哪一个线程已经注册”的名单。因此,代码必须自己保证注册和到达次数一致。把一个份额理解成一项持续参与的责任,比简单地数线程更容易检查流程。
本文安排A、B、C三个工作者:A和B准备完后等待;C只负责准备,完成后退出。此时第一阶段齐了,A和B可以开始加工,后续不用再等待C。
图中的卡片代表参与份额,猫是讲解者,不表示线程身份。C退出的同时也完成了本阶段的报到。
二、把几个动作分清楚
| 动作 | 方法 | 理解方式 |
|---|---|---|
| 增加参与份额 | register() | 当前流程又多了一项到达责任 |
| 到达后等待 | arriveAndAwaitAdvance() | 本阶段工作完成,等大家齐了再继续 |
| 到达但不等待 | arrive() | 只报到,是否等待另行安排 |
| 到达并退出 | arriveAndDeregister() | 本阶段报到,后续不再参加 |
| 按阶段等待 | awaitAdvance(phase) | 等指定阶段发生变化 |
这些方法的语义依据JDK21 API。实际使用时,建议先画出每个工作者在哪一阶段离开,再写代码。尤其不要把arrive()和随后一次arriveAndAwaitAdvance()连着调用来“补一个等待”:后者会再次报到,计数就乱了。
三、完整示例:三人准备,两人加工
这里先一次性注册三个份额,再启动三个线程。示例故意不使用随机休眠,让协调逻辑保持清晰。
importjava.util.concurrent.Phaser;publicclassPhaserDemo{publicstaticvoidmain(String[]args)throwsInterruptedException{Phaserphaser=newPhaser(3);Threada=worker("A",false,phaser);Threadb=worker("B",false,phaser);Threadc=worker("C",true,phaser);a.start();b.start();c.start();a.join();b.join();c.join();System.out.println("全部结束,已终止="+phaser.isTerminated());}privatestaticThreadworker(Stringname,booleanpreparationOnly,Phaserphaser){returnnewThread(()->{System.out.println(name+":准备完成");if(preparationOnly){phaser.arriveAndDeregister();return;}intnextPhase=phaser.arriveAndAwaitAdvance();if(nextPhase<0){return;}System.out.println(name+":加工完成");phaser.arriveAndDeregister();},"worker-"+name);}}保存为PhaserDemo.java,可用javac PhaserDemo.java编译,再运行java PhaserDemo。按代码逻辑,一种可能的输出如下;同阶段内部的先后顺序并不固定:
A:准备完成 C:准备完成 B:准备完成 B:加工完成 A:加工完成 全部结束,已终止=true观察输出时抓住两个约束即可:两条“加工完成”前,三条“准备完成”都已经出现;最后一行在三个线程结束后出现。这里主线程靠**join()**等待线程生命周期结束,工作者之间靠Phaser协调阶段,两个等待点各有职责。
四、动态注册时,先留住协调份额
如果任务列表是运行时才确定的,可以让协调线程先占一个份额,避免前面的任务过早结束、整个协调器已经终止,后面的任务才来注册。
Phaserphaser=newPhaser(1);// 协调者先占一个份额for(Runnabletask:tasks){phaser.register();// 注册在提交前;提交失败时也必须处理这份责任try{executor.execute(()->{try{task.run();}finally{phaser.arriveAndDeregister();}});}catch(java.util.concurrent.RejectedExecutionExceptionex){phaser.arriveAndDeregister();// 示例选择记录并继续;生产代码应明确失败策略System.err.println("任务提交失败:"+ex.getMessage());}}phaser.arriveAndAwaitAdvance();phaser.arriveAndDeregister();这段是局部用法,tasks和executor由调用方提供。它等待已提交任务结束,却没有汇总业务成功与否。需要结果时,可额外保存Future或失败记录;“到齐”与“全部成功”应分别判断。
五、容易忽略的注意事项
**每阶段只报到一次。**如果已经用arrive()报到,后面需要的是等待操作;异常清理时也要结合当前阶段责任判断,不能无条件再减一次。
**超时不替你移除参与者。**需要可中断或限时等待时,可使用awaitAdvanceInterruptibly。等待异常不会自动破坏Phaser;forceTermination()能释放等待者,但不会自动取消业务任务。默认情况下,全部份额退出后会终止,之后重新注册无法恢复。官方等待与终止语义
**线程池容量要覆盖等待依赖。**如果先运行的任务占满线程池并等待尚在队列里的参与者,流程就可能无法推进。动态计数解决不了执行资源不足的问题。
六、思维导图
总结要点
Phaser适合参与数量会变化的阶段协作。先明确每个工作者承担几轮责任,再把注册、到达和退出对应到流程里,代码会更容易核对。
到达次数正确,是阶段正常推进的基础。对动态任务,要同时处理提交失败、执行失败和等待超时,避免留下无人履行的份额。
协调工具只负责推进时机。业务是否成功、线程是否结束、失败后如何取消,仍然需要各自清楚的处理方式。
下一篇继续认识Semaphore,看看怎样限制同时进入某项资源的任务数量。
👉如果你觉得这篇文章对你有所帮助,欢迎点赞、收藏、分享!😊