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

Wait()与NotifyAll()在消息中间件订阅者方法中失效问题排查

问题:消息中间件订阅者无法从等待状态被唤醒

我正在开发一款面向消息的中间件,发布者与订阅者通过中间类MessageBroker中的队列进行通信。发布者需在主题队列未满时向队列投递消息,订阅者需在队列非空时获取订阅主题的消息。目前发布者端功能正常,但订阅者调用receiveMessage()方法时无法从等待状态被唤醒,无任何操作输出。

现有代码

发布者方法

public synchronized void sendMessage() {
        BlockingQueue<Message> topicQueue = mb.getQueue(topic);
        if (topicQueue != null) {
            try {
                // 发送消息前检查队列是否有剩余空间
                while (topicQueue.remainingCapacity() == 0) {
                    wait();
                }
            } catch (InterruptedException e) {
                e.printStackTrace();
            }

            Message m = topic.generateMessage();
            mb.publish(m);
            if (!gotActive) {
                mb.increasePublisherCounter(1);
                gotActive = true;
            }
            System.out.println(m.getContent() + " 已添加至队列");
            notifyAll(); // 通知所有线程
        }
    }

订阅者方法

public synchronized void receiveMessage() {
        for (Topic topic : topics) {
            BlockingQueue<Message> queue = mb.getQueue(topic);
            synchronized (queue) {
                try {
                    // 等待队列非空
                    while (queue.isEmpty()) {
                        queue.wait(); // 队列为空时等待通知
                    }

                    // 从队列获取消息
                    Message message = queue.peek();
                    if (message != null) {
                        System.out.println(name + " 收到消息: " + message.getContent());
                        incrementProcessedCounter(topic);
                        if (topic.getSubscriberCount() == getProcessedCounter(topic)) {
                            queue.remove();
                            resetProcessedCounter(topic);
                            queue.notifyAll(); // 通知其他等待线程
                        }
                    }
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        }
    }

问题原因与修复方案

核心问题:锁对象不匹配

发布者的synchronized方法是锁定当前发布者实例,wait()和notifyAll()操作的是发布者实例的锁等待池;而订阅者是锁定队列对象,queue.wait()操作的是队列对象的锁等待池。两者的等待池完全独立,发布者的通知无法传递到订阅者,导致订阅者永远无法被唤醒。

具体修改点

  1. 修改发布者的锁与通知对象
    去掉sendMessage方法的synchronized修饰,改为对对应的队列对象加锁,所有wait()和notifyAll()操作都绑定到队列对象:

    public void sendMessage() {
        BlockingQueue<Message> topicQueue = mb.getQueue(topic);
        if (topicQueue != null) {
            synchronized (topicQueue) { // 锁定队列对象
                try {
                    // 等待队列有剩余空间
                    while (topicQueue.remainingCapacity() == 0) {
                        topicQueue.wait(); // 在队列锁上等待
                    }
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
    
                Message m = topic.generateMessage();
                mb.publish(m); // 确保此方法是向topicQueue中添加消息
                if (!gotActive) {
                    mb.increasePublisherCounter(1);
                    gotActive = true;
                }
                System.out.println(m.getContent() + " 已添加至队列");
                topicQueue.notifyAll(); // 通知等待队列锁的订阅者线程
            }
        }
    }
    
  2. 调整订阅者的持续监听逻辑
    原订阅者方法仅遍历一次订阅主题就结束,处理完一次消息后会退出方法,无法接收后续消息。需在外层添加循环,保持持续监听(直到线程被中断):

    public void receiveMessage() {
        // 持续监听直到线程被中断
        while (!Thread.currentThread().isInterrupted()) {
            for (Topic topic : topics) {
                BlockingQueue<Message> queue = mb.getQueue(topic);
                if (queue == null) continue;
                synchronized (queue) {
                    try {
                        while (queue.isEmpty()) {
                            queue.wait();
                        }
    
                        Message message = queue.peek();
                        if (message != null) {
                            System.out.println(name + " 收到消息: " + message.getContent());
                            incrementProcessedCounter(topic);
                            if (topic.getSubscriberCount() == getProcessedCounter(topic)) {
                                queue.remove();
                                resetProcessedCounter(topic);
                            }
                            queue.notifyAll(); // 通知其他线程(包括发布者)
                        }
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt(); // 保留线程中断状态
                        e.printStackTrace();
                    }
                }
            }
        }
    }
    
  3. 移除订阅者方法的synchronized修饰
    订阅者方法原本的synchronized会锁定订阅者实例,与队列锁无关,属于冗余锁,直接移除即可。

内容的提问来源于stack exchange,提问作者LennartPfeiler

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 15:48:10