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()操作的是队列对象的锁等待池。两者的等待池完全独立,发布者的通知无法传递到订阅者,导致订阅者永远无法被唤醒。
具体修改点
修改发布者的锁与通知对象
去掉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(); // 通知等待队列锁的订阅者线程 } } }调整订阅者的持续监听逻辑
原订阅者方法仅遍历一次订阅主题就结束,处理完一次消息后会退出方法,无法接收后续消息。需在外层添加循环,保持持续监听(直到线程被中断):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(); } } } } }移除订阅者方法的synchronized修饰
订阅者方法原本的synchronized会锁定订阅者实例,与队列锁无关,属于冗余锁,直接移除即可。
内容的提问来源于stack exchange,提问作者LennartPfeiler
相关产品推荐
相关产品推荐

