WebSphere J2EE迁移Spring Boot后JMS监听器重复消费问题
JMS消息重复消费问题排查与修复
核心问题分析
当前配置存在事务与消息确认机制不匹配的问题,直接导致消息在事务未提交时就被MQ视为已处理,进而分发给其他并发消费者:
DefaultMessageListenerContainer同时配置了JTA事务管理器,但设置了sessionTransacted=false和Session.AUTO_ACKNOWLEDGE,这会让JMS自动在消息接收后立即确认,而非等待事务提交完成。- Bitronix连接工厂的
allowLocalTransactions=true可能引发XA事务与本地事务混用,破坏MQ的消息锁定逻辑。 - 并发消费者设置为2,但指定了固定
clientId,队列场景下多个消费者共用同一ClientId会干扰MQ的消息分发锁定策略。
修复步骤
1. 调整消息监听器容器的事务与确认配置
修改wagenstatusMessageListenerContainer配置,让容器通过事务管理消息的提交/回滚:
@Bean("wgstatusML") public DefaultMessageListenerContainer wagenstatusMessageListenerContainer( ConnectionFactory jmsXAConnectionFactory, PlatformTransactionManager jtaTransactionManager, @Qualifier("wagenstatusBean") WagenstatusBean wagenstatusBean) { DefaultMessageListenerContainer container = new DefaultMessageListenerContainer(); container.setConnectionFactory(jmsXAConnectionFactory); container.setTransactionManager(jtaTransactionManager); container.setDestinationName(WAGENSTATUS_QUEUE); container.setMessageListener(wagenstatusBean); container.setAutoStartup(false); container.setConcurrentConsumers(2); // 移除固定clientId,队列场景不需要,并发消费者不能共用同一clientId // container.setClientId("wgstatListener"); // 开启会话事务,由事务管理器控制消息的提交/回滚 container.setSessionTransacted(true); // 当使用事务管理器时,确认模式会被容器忽略,统一使用事务控制 container.setSessionAcknowledgeMode(Session.SESSION_TRANSACTED); return container; }
2. 修正Bitronix连接工厂的XA配置
禁用本地事务,确保所有操作都参与XA事务:
@Bean ConnectionFactory jmsXAConnectionFactory() { PoolingConnectionFactory connectionFactory = new PoolingConnectionFactory(); connectionFactory.setClassName("com.ibm.mq.jms.MQXAQueueConnectionFactory"); connectionFactory.setUniqueName("mq-xa-" + appName); // 禁用本地事务,强制使用XA事务,避免事务边界混乱 connectionFactory.setAllowLocalTransactions(false); connectionFactory.setTestConnections(false); connectionFactory.setUser(user); connectionFactory.setPassword(password); connectionFactory.setMaxIdleTime(1800); connectionFactory.setMinPoolSize(1); connectionFactory.setMaxPoolSize(25); connectionFactory.setAcquisitionTimeout(60); connectionFactory.setAutomaticEnlistingEnabled(true); connectionFactory.setDeferConnectionRelease(true); connectionFactory.setShareTransactionConnections(false); Properties driverProperties = connectionFactory.getDriverProperties(); driverProperties.setProperty("queueManager", queueManager); driverProperties.setProperty("hostName", connName); driverProperties.setProperty("port", "1414"); driverProperties.setProperty("channel", channel); driverProperties.setProperty("transportType", "1"); driverProperties.setProperty("messageRetention", "1"); return connectionFactory; }
3. 优化消息监听器的事务逻辑
确保onMessage方法的事务控制与容器一致,避免手动操作的冲突:
@Service("wagenstatusBean") @Scope(SCOPE_PROTOTYPE) public class WagenstatusBean extends AbstractMDB { @Transactional(propagation = Propagation.REQUIRED) public void onMessage(javax.jms.Message msg) { String localMessageText = null; try { try { localMessageText = ((TextMessage) msg).getText(); } catch (JMSException e) { LOGGER.error("Failed to get message text", e); } String errmsg = null; readableMessageID = null; try { verarbeiteMeldung(msg); } catch (InvalidMessageException ime) { errmsg = ime.getMessage(); // 触发异常时标记事务回滚 TransactionAspectSupport.currentTransactionStatus().setRollbackOnly(); } if (sendMessageToErrorQueue) { try { logBusinessData(localMessageText, BusinessLogger.STATUS_ERROR); } catch (Exception e) { LOGGER.error("Failed to log business data", e); } if (localMessageText != null) { localMessageText = this.addErrorMessageToXML(localMessageText, errmsg); } DispatcherServiceLocator.getDispatcherBean().sendToDestination( QueueNames.WAGENSTATUS_ERROR_QUEUE, localMessageText); } } catch (ConsistencyException ex) { TransactionAspectSupport.currentTransactionStatus().setRollbackOnly(); try { logBusinessData(localMessageText, BusinessLogger.STATUS_ERROR); } catch (Exception e) { LOGGER.error("Failed to log business data", e); } LOGGER.error("ConsistencyException in WagenStatus-onMessage", ex); } catch (RuntimeException ex) { TransactionAspectSupport.currentTransactionStatus().setRollbackOnly(); try { logBusinessData(localMessageText, BusinessLogger.STATUS_ERROR); } catch (Exception e) { LOGGER.error("Failed to log business data", e); } LOGGER.error("RuntimeException in WagenStatus-onMessage", ex); } } }
额外注意事项
- 检查IBM MQ队列的
DEFPSIST属性是否设置为YES,确保消息持久化,避免重启丢失或重复分发。 - 根据MQ服务器的承载能力调整
concurrentConsumers数量,避免过度并发导致的资源竞争。
内容的提问来源于stack exchange,提问作者myborobudur
相关产品推荐
相关产品推荐

