使用Semaphore同步多阶段线程时第三阶段无法运行问题求助
问题分析与解决方案
咱们先拆解你遇到的两个核心问题:第三阶段无法执行,以及信号量许可异常,这俩其实是连锁反应,根源出在线程间的通信可见性和同步逻辑上。
1. 第三阶段不执行的核心原因:可见性缺失与无同步的空轮询
你的Fase2和Fase3里用了普通的boolean receivedMessage标记是否收到消息,但这个变量没有被volatile修饰,而且线程之间没有任何同步屏障。根据Java内存模型,一个线程修改的普通变量,其他线程可能看不到最新值——也就是说,Fase2的receiveMessage方法把receivedMessage设为true后,Fase2的run方法可能一直读取到初始的false,根本不会进入处理逻辑,自然不会把消息传给Fase3,Fase3也就一直处于闲置状态。
另外,这种while(true)里只判断if(receivedMessage)的空轮询,非常浪费CPU资源,也不是线程间通信的正确方式。
2. 信号量许可异常的潜在原因:不合理的许可申请与异常吞吃
你看Fase1每次申请2个信号量0和1的许可,但初始信号量1只有3个许可——第一次申请后还剩1个,第二次再申请2个的话,就会直接阻塞在resources[1].acquire(2)这里,导致Fase1无法继续生成消息,后续流程全卡壳。而且你的代码直接吞掉了InterruptedException,如果线程在acquire之后、release之前被中断,就会导致许可无法释放,进而引发整个信号量池的异常。
修复后的完整代码
咱们一步步修正这些问题:
- 用
volatile修饰共享变量,保证线程间的可见性; - 替换空轮询为
wait/notify机制,更高效地实现线程间通信; - 调整信号量的许可申请数量,适配初始许可数避免阻塞;
- 不要吞掉
InterruptedException,至少打印日志便于排查; - 添加优雅终止逻辑,避免无限循环。
import java.util.concurrent.Semaphore; class Fase1 extends Thread { private int i = 0; private final Semaphore[] resources; private final Fase2 recipient; private volatile boolean running = true; public Fase1(Semaphore[] res, Fase2 fase2) { recipient = fase2; resources = res; } @Override public void run() { try { // 限制生成10条消息,避免无限循环 while (running && i < 10) { // 调整为每次申请1个许可,适配初始信号量的许可数 resources[0].acquire(1); resources[1].acquire(1); synchronized (recipient) { recipient.receiveMessage(i); recipient.notify(); // 主动通知Fase2有新消息 } i++; Thread.sleep(200); resources[1].release(1); resources[0].release(1); } // 通知Fase2结束流程 synchronized (recipient) { running = false; recipient.notify(); } } catch (InterruptedException e) { System.err.println("Fase1被中断: " + e.getMessage()); Thread.currentThread().interrupt(); // 恢复中断状态 } } } class Fase2 extends Thread { private final Semaphore[] resources; private final Fase3 recipient; private volatile boolean receivedMessage = false; private volatile int message = 0; private volatile boolean running = true; public Fase2(Semaphore[] res, Fase3 fase3) { recipient = fase3; resources = res; } @Override public void run() { try { while (running) { synchronized (this) { // 没有消息时进入等待,避免空轮询浪费CPU while (!receivedMessage && running) { wait(); } if (!running) break; resources[0].acquire(1); resources[1].acquire(1); resources[2].acquire(1); synchronized (recipient) { recipient.receiveMessage(message * 2); recipient.notify(); // 通知Fase3处理消息 } receivedMessage = false; Thread.sleep(200); resources[2].release(1); resources[1].release(1); resources[0].release(1); } } // 通知Fase3结束流程 synchronized (recipient) { running = false; recipient.notify(); } } catch (InterruptedException e) { System.err.println("Fase2被中断: " + e.getMessage()); Thread.currentThread().interrupt(); } } public synchronized void receiveMessage(int msg) { message = msg; receivedMessage = true; } } class Fase3 extends Thread { private final Semaphore[] resources; private volatile boolean receivedMessage = false; private volatile int message = 0; private volatile boolean running = true; public Fase3(Semaphore[] res) { resources = res; } @Override public void run() { try { while (running) { synchronized (this) { while (!receivedMessage && running) { wait(); } if (!running) break; resources[1].acquire(1); resources[2].acquire(1); System.out.println("最终输出: " + (message + 1)); receivedMessage = false; Thread.sleep(200); resources[2].release(1); resources[1].release(1); } } } catch (InterruptedException e) { System.err.println("Fase3被中断: " + e.getMessage()); Thread.currentThread().interrupt(); } } public synchronized void receiveMessage(int msg) { message = msg; receivedMessage = true; } } // 测试主类 public class MultiPhaseSystem { public static void main(String[] args) { Semaphore[] resources = new Semaphore[3]; // 初始化公平信号量,保持你原来的许可数设置 resources[0] = new Semaphore(4, true); resources[1] = new Semaphore(3, true); resources[2] = new Semaphore(2, true); Fase3 fase3 = new Fase3(resources); Fase2 fase2 = new Fase2(resources, fase3); Fase1 fase1 = new Fase1(resources, fase2); // 按从后到前的顺序启动线程,避免消息丢失 fase3.start(); fase2.start(); fase1.start(); } }
关键修复点说明
- 可见性保障:用
volatile修饰共享变量,确保一个线程修改后其他线程能立即看到最新值; - 高效线程通信:用
synchronized配合wait/notify替代空轮询,减少CPU消耗,保证消息传递的及时性; - 信号量逻辑优化:把每次申请的许可数从2调整为1,适配你初始的信号量许可数量,避免不必要的阻塞;
- 异常处理:不再吞掉
InterruptedException,而是打印日志并恢复中断状态,便于排查问题; - 终止逻辑:添加
running标记,让线程可以优雅终止,避免无限循环。
内容的提问来源于stack exchange,提问作者Federico Chiesa
相关产品推荐
相关产品推荐

