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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 17:09:33