异步连接IBM MQ时无法消费消息或仅能检索到单条消息
问题描述
使用Java编写IBM MQ异步消息消费逻辑时,运行存在异常:部分场景下仅消费1条消息后进程就停止运行,部分场景下完全无法消费任何消息。代码中已经通过setMessageListener()方法实现异步消息监听,但无法定位错误点,原实现代码如下:
public class MQConnectivity implements MessageListener { MQQueueConnectionFactory factory = new MQQueueConnectionFactory(); QueueConnection con; Queue queue; QueueSession session; QueueReceiver receiver; public void init() { try { factory.setHostName("******"); factory.setPort(123); factory.setChannel("Channel1"); factory.setQueueManager("queue_manager"); factory.setTransportType(new Integer(WMQConstants.WMQ_CM_CLIENT)); con=factory.createQueueConnection("user1"," password1"); queue= new MQQueue("myQueue"); session=con.createQueueSession(false, Session.AUTO_ACKNOWLEDGE); receiver=session.createReceiver(queue); receiver.setMessageListener(this); System.out.println("=======connection happened========="); con.start(); System.out.println("=======connection started========="); }catch(Exception ex) { System.out.println("exception is : "+ex); } } public static void main(String[] args) { new MQConnectivity().init(); } @Override public void onMessage(Message msg) { try { System.out.println("Inside onMessage method"); TextMessage txtMsg=(TextMessage)msg; System.out.println(txtMsg.getText()); }catch(Exception ex) { System.out.println("error in onMessage"+ex); } } }
根因分析
代码存在4个问题,直接导致上述异常现象:
- 核心问题:主线程执行完
init()方法后直接终止,JVM没有存活的非守护线程会直接退出。IBM MQ客户端内部用于异步回调onMessage的工作线程是守护线程,不会阻止JVM关停。如果main线程执行完con.start()立刻结束,JVM直接退出就表现为完全收不到消息;如果JVM退出前刚好有一条消息投递触发回调,执行完onMessage后JVM依旧会终止,就表现为仅消费1条消息就停止。 - 认证参数错误:创建连接时传入的密码参数前多了一个前置空格,即
" password1",会导致MQ服务端认证失败,连接被主动断开,无法正常消费消息。 - 缺少异常监听:没有为MQ连接注册异常回调,连接因为网络波动、权限不足、认证失败等原因断开时,不会打印任何错误日志,只会静默失败,排查难度极高。
- 消息处理逻辑隐患:
onMessage中直接将消息强转为TextMessage,如果队列中存在字节消息、Map消息等其他类型的消息,会直接抛出类型转换异常,导致消费失败。
修复方案
针对上述问题逐点修复即可:
- 在
con.start()之后增加主线程阻塞逻辑,保证JVM持续存活,不要让main方法执行完直接退出。 - 移除密码参数前的多余空格,保证认证参数正确。
- 为连接注册
ExceptionListener,捕获连接异常,方便排查问题和实现自定义重连逻辑。 - 在
onMessage中增加消息类型判断,避免非文本消息触发强转错误。 - 异常捕获处增加堆栈打印,不要只打印异常对象,方便定位问题。
修复后的可运行代码示例:
public class MQConnectivity implements MessageListener, ExceptionListener { MQQueueConnectionFactory factory = new MQQueueConnectionFactory(); QueueConnection con; Queue queue; QueueSession session; QueueReceiver receiver; public void init() { try { factory.setHostName("******"); factory.setPort(123); factory.setChannel("Channel1"); factory.setQueueManager("queue_manager"); factory.setTransportType(WMQConstants.WMQ_CM_CLIENT); // 修复密码前置空格问题 con = factory.createQueueConnection("user1", "password1"); // 注册连接异常监听器 con.setExceptionListener(this); queue = new MQQueue("myQueue"); session = con.createQueueSession(false, Session.AUTO_ACKNOWLEDGE); receiver = session.createReceiver(queue); receiver.setMessageListener(this); System.out.println("=======connection created========="); con.start(); System.out.println("=======connection started, waiting for messages========="); // 阻塞主线程,防止JVM退出 synchronized (this) { this.wait(); } } catch (Exception ex) { System.out.println("init exception: " + ex); ex.printStackTrace(); } } public static void main(String[] args) { new MQConnectivity().init(); } @Override public void onMessage(Message msg) { try { System.out.println("received message, start processing"); // 增加消息类型判断,避免强转报错 if (msg instanceof TextMessage) { TextMessage txtMsg = (TextMessage) msg; System.out.println("message content: " + txtMsg.getText()); } else { System.out.println("skip non-text message, type: " + msg.getJMSType()); } } catch (Exception ex) { System.out.println("process message error: " + ex); ex.printStackTrace(); } } @Override public void onException(JMSException e) { System.out.println("mq connection error: " + e); e.printStackTrace(); // 可在此处添加自定义重连逻辑 } }
内容的提问来源于stack exchange,提问作者Prasad
相关产品推荐
相关产品推荐

