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

ActiveMQ单条消息被多消费者重复消费问题及解决方案咨询

问题分析与解决方案

核心原因

JMS队列本身是点对点模型,理论上单条消息只会被一个消费者消费。出现重复消费的常见诱因:

  • 消息确认模式配置错误:如果用了AUTO_ACKNOWLEDGE,消费者处理中出现异常或Broker判定消费者断开时,会重发消息;DUPS_OK_ACKNOWLEDGE模式本身允许重复投递以换取性能。
  • Prefetch Size设置过大:ActiveMQ默认预取1000条消息,并发场景下,消费者预取消息后处理过慢,Broker会认为消费者挂掉,将未确认的消息重新分配给其他消费者。
  • 事务处理不当:若使用事务但提交失败,消息会回滚到队列,触发重新投递。
  • 业务无幂等校验:即使消息重复投递,也没有机制识别并跳过重复处理。

针对性解决步骤

1. 切换为手动消息确认模式

将确认模式改为CLIENT_ACKNOWLEDGE,确保只有业务逻辑完全执行成功后才手动确认消息,避免异常导致的重发。

修改消费代码:

@Autowired
ReconMessageProcessor reconMessageProcessor;
@Value("${recondataProcessorRunEnable}")
public Boolean recondataProcessorRunEnable = true;

public void onMessage(Message pMessage) {
    if (recondataProcessorRunEnable) {
        try {
            logger.info("RetryReconDataInsertMessageListener pMessage.getJMSCorrelationID() -->  " + pMessage.getJMSCorrelationID());
            logger.info("RetryReconDataInsertMessageListener pMessage.getJMSMessageID() -->  " + pMessage.getJMSMessageID());
            MessageEntity message = (MessageEntity)((ActiveMQObjectMessage) pMessage).getObject();
            String jsonMessage = message.getJsonString();
            reconMessageProcessor.processMessage(jsonMessage, "RetryReconDataInsertMessageListener");
            
            // 业务执行成功后手动确认消息
            pMessage.acknowledge();
        } catch (Exception e) {
            logger.error("Exception has occured while executing ReconMessageListener. Error details is : " + e.getMessage());
            e.printStackTrace();
            // 不可恢复异常可考虑转入死信队列,避免无限重发
        }
    } else {
        logger.info(" Value of recondataProcessorRunEnable is " + recondataProcessorRunEnable + " at property file. Please set it to true to enable recondataProcessor run");
    }
}

同时在监听器配置中指定确认模式(以Spring配置为例):

<jms:listener-container connection-factory="connectionFactory" acknowledge="client">
    <jms:listener destination="update" ref="retryReconDataInsertMessageListener" />
</jms:listener-container>

2. 缩小Prefetch Size

将update队列的预取数设为1,确保每个消费者每次仅处理一条消息,处理完成确认后再取下一条,避免消息被重复分配。

队列配置方式:

// 消费者端声明队列时设置
ActiveMQQueue updateQueue = new ActiveMQQueue("update?consumer.prefetchSize=1");

或在ActiveMQ的activemq.xml全局配置:

<policyEntry queue=">">
    <policyMap>
        <entry queue="update" value="prefetchSize=1"/>
    </policyMap>
</policyEntry>

3. 启用事务机制(可选)

如果业务涉及数据库等操作,开启JMS事务确保业务操作与消息确认的原子性:事务提交则消息被Broker移除,事务回滚则消息重新投递。

修改监听器配置:

<jms:listener-container connection-factory="connectionFactory" transaction-manager="transactionManager" acknowledge="transacted">
    <jms:listener destination="update" ref="retryReconDataInsertMessageListener" />
</jms:listener-container>

代码无需手动确认,事务提交时自动完成消息确认:

public void onMessage(Message pMessage) {
    if (recondataProcessorRunEnable) {
        try {
            logger.info("RetryReconDataInsertMessageListener pMessage.getJMSCorrelationID() -->  " + pMessage.getJMSCorrelationID());
            logger.info("RetryReconDataInsertMessageListener pMessage.getJMSMessageID() -->  " + pMessage.getJMSMessageID());
            MessageEntity message = (MessageEntity)((ActiveMQObjectMessage) pMessage).getObject();
            String jsonMessage = message.getJsonString();
            reconMessageProcessor.processMessage(jsonMessage, "RetryReconDataInsertMessageListener");
            // 事务自动提交,消息确认
        } catch (Exception e) {
            logger.error("Exception has occured while executing ReconMessageListener. Error details is : " + e.getMessage());
            e.printStackTrace();
            // 抛出异常触发事务回滚,消息重新投递
            throw new RuntimeException(e);
        }
    } else {
        logger.info(" Value of recondataProcessorRunEnable is " + recondataProcessorRunEnable + " at property file. Please set it to true to enable recondataProcessor run");
    }
}

4. 实现业务幂等性(兜底方案)

无论消息是否重复投递,通过唯一标识确保业务逻辑不会重复执行。比如用消息的JMSMessageID或业务唯一ID做校验:

// 在ReconMessageProcessor的processMessage方法中添加幂等校验
public void processMessage(String jsonMessage, String listenerName) {
    // 解析消息中的唯一标识(比如业务ID或JMSMessageID)
    String uniqueMsgId = extractUniqueId(jsonMessage);
    // 检查数据库/缓存中是否已有处理记录
    if (msgProcessRecordRepo.existsById(uniqueMsgId)) {
        logger.info("消息已处理,跳过:" + uniqueMsgId);
        return;
    }
    // 执行更新业务逻辑
    updateReconRecord(jsonMessage);
    // 记录处理完成状态
    msgProcessRecordRepo.save(new MsgProcessRecord(uniqueMsgId));
}

总结

优先通过调整消息确认模式和Prefetch Size解决重复消费的触发问题,同时必须在业务层实现幂等性校验,这是保障业务最终一致性的兜底手段。

内容的提问来源于stack exchange,提问作者Hus Mukh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 20:20:39