You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 09:06:39