IBM MQ消费者onMessage方法未触发问题:receive正常监听器无响应
IBM MQ 消息监听器
onMessage方法未触发问题 问题描述
使用IBM MQ时,调用receive()方法能正常消费消息,但通过MessageListener的onMessage()方法实现异步消费时,该方法从未被调用,无输出也无错误。
正常工作的同步消费代码
public class IBMMQConnect { private String hostName = "myHost"; private int port = 1631; private String clientID = "clientID"; private String channel = "channel"; private int transportType = WMQConstants.WMQ_CM_CLIENT; private String queueManager = "queueManager"; private String sslCipherSuite = "sslCipherSuite"; MQConnectionFactory mqConnectionFactory; CustomSslContectFactory sslFactory; javax.jms.ConnectionFactory connectionFactory; javax.jms.Connection connection; javax.jms.Session session; javax.jms.Topic topic; javax.jms.MessageConsumer consumer; public IBMMQConnect() throws Exception{ System.setProperty("com.ibm.mq.cfg.useIBMCipherMappings", "false"); this.mqConnectionFactory = new MQConnectionFactory(); this.mqConnectionFactory.setHostName(hostName); this.mqConnectionFactory.setPort(port); this.mqConnectionFactory.setClientID(clientID); this.mqConnectionFactory.setChannel(channel); this.mqConnectionFactory.setTransportType(transportType); this.mqConnectionFactory.setQueueManager(queueManager); this.mqConnectionFactory.setSSLCipherSuite(sslCipherSuite); this.sslFactory = new CustomSslContextFactory(); this.sslFactory.setKeyStore("myKeyStore"); this.sslFactory.setKeyStorePassword("myPassword"); this.sslFactory.setTrustStore("myTrustStore"); this.sslFactory.setTrustStorePassword("password"); this.mqConnectionFactory.setSSLSocketFactory(sslFactory.getSslContext().getSocketFactory()); this.connectionFactory = this.mqConnectionFactory; this.connection = this.connectionFactory.createConnection(); this.session = this.connection.createSession(false, javax.jms.Session.CLIENT_ACKNOWLEDGE); this.topic = this.session.createTopic("myTopic"); this.consumer = this.session.createDurableSubscriber(topic, "myTest"); javax.jms.Message message = null; this.connection.start(); while ((message = consumer.receive(10*1000))!=null){ System.out.println(message); } } } public class Main { public static void main(String[] args) throws Exception{ IBMMQConnect ibmMQ = new IBMMQConnect(); } }
异步消费未触发onMessage的代码
public class IBMMQConnect { private String hostName = "myHost"; private int port = 1631; private String clientID = "clientID"; private String channel = "channel"; private int transportType = WMQConstants.WMQ_CM_CLIENT; private String queueManager = "queueManager"; private String sslCipherSuite = "sslCipherSuite"; MQConnectionFactory mqConnectionFactory; CustomSslContectFactory sslFactory; javax.jms.ConnectionFactory connectionFactory; javax.jms.Connection connection; javax.jms.Session session; javax.jms.Topic topic; javax.jms.MessageConsumer consumer; public IBMMQConnect() throws Exception{ System.setProperty("com.ibm.mq.cfg.useIBMCipherMappings", "false"); this.mqConnectionFactory = new MQConnectionFactory(); this.mqConnectionFactory.setHostName(hostName); this.mqConnectionFactory.setPort(port); this.mqConnectionFactory.setClientID(clientID); this.mqConnectionFactory.setChannel(channel); this.mqConnectionFactory.setTransportType(transportType); this.mqConnectionFactory.setQueueManager(queueManager); this.mqConnectionFactory.setSSLCipherSuite(sslCipherSuite); this.sslFactory = new CustomSslContextFactory(); this.sslFactory.setKeyStore("myKeyStore"); this.sslFactory.setKeyStorePassword("myPassword"); this.sslFactory.setTrustStore("myTrustStore"); this.sslFactory.setTrustStorePassword("password"); this.mqConnectionFactory.setSSLSocketFactory(sslFactory.getSslContext().getSocketFactory()); this.connectionFactory = this.mqConnectionFactory; this.connection = this.connectionFactory.createConnection(); this.session = this.connection.createSession(false, javax.jms.Session.CLIENT_ACKNOWLEDGE); this.topic = this.session.createTopic("myTopic"); this.consumer = this.session.createDurableSubscriber(topic, "myTest"); javax.jms.Message message = null; this.connection.start(); this.consumer.setMessageListener(new MessageListener() { @Override public void onMessage(Message message) { System.out.println(message); } }); new Scanner(System.in).nextLine(); } }
问题原因及解决方案
核心原因:监听器设置顺序错误
在JMS规范及IBM MQ的实现中,必须先为消息消费者设置MessageListener,再启动连接。你的异步消费代码中,先调用了connection.start(),之后才设置监听器,此时连接已经开始分发消息,监听器无法捕获已开始传输的消息,后续消息也无法被正确路由到监听器。
修复步骤
调整代码顺序,将consumer.setMessageListener()移到connection.start()之前:
// 先设置监听器 this.consumer.setMessageListener(new MessageListener() { @Override public void onMessage(Message message) { System.out.println(message); } }); // 再启动连接 this.connection.start();
额外检查点
- 持久订阅状态:确保持久订阅的
clientID和订阅名称myTest在MQ服务器上状态正常,未被标记为失效或删除。 - 异常捕获:在
onMessage方法中添加异常捕获,避免因未处理的异常导致监听器静默失效:
@Override public void onMessage(Message message) { try { System.out.println(message); // 若使用CLIENT_ACKNOWLEDGE模式,需手动确认消息 message.acknowledge(); } catch (JMSException e) { e.printStackTrace(); } }
- 线程保持:确保程序主线程不会提前退出,你当前使用
new Scanner(System.in).nextLine()的方式是有效的,但需注意不要阻塞JMS客户端的消息分发线程。
内容的提问来源于stack exchange,提问作者JPA
相关产品推荐
相关产品推荐

