一、先从一个让人抓狂的现象说起

你肯定遇到这样的事:一个订单系统,后端处理完支付后往ActiveMQ里发一条“支付成功”的消息。发送方代码里明明用了事务,也调用了提交方法。消息也到了消息队列。可是消费端在处理完订单之后,服务凌晨重启了一下,第二天发现订单被重复处理了两次。这时候你的第一反应多半是——消息队列是不是出bug了?还是消息中间件丢了状态?其实都不是,真正的原因多半是搞混了“JMS事务”和“本地事务”的边界。通俗点说,你以为生产端提交了事务,消息就不会重复了;但消费端那一侧的确认动作可能根本没完成,消息自然会被重新投递。

接下来我用大白话拆开揉碎了讲,保证你从没接触过消息事务也能听懂。我还会给出完整的Java示例,把复现过程一步步跑出来。

二、JMS事务到底管什么

2.1 事务不是“只发一次”的意思

JMS事务是一个严格意义上的“原子性”工具。举个例子,你在银行转账,既要扣付款人的钱,又要加收款人的钱,这两件事必须同时成功或同时失败。JMS事务解决的是这类问题:生产端在同一个事务里发消息和更新数据库,要么都成功,要么都失败。它解决的是“原子性”,不是“消息不会重复”。这俩概念差别太大了。消息队列的网络环境很复杂,一个消息从生产端传到Broker,再从Broker传到消费端,每一段都可能发生断线、超时、重连。实现“只传一次”的真实成本非常高,所以大部分消息中间件默认提供“至少一次”的语义。也就是说,重复是完全可能出现的,事务并不负责消除重复。

2.2 本地事务的边界在哪

所谓“本地事务”,指的是只在一个会话内部生效的事务。生产端有一个会话,消费端有自己的会话,这俩会话没有直接关系。生产端的事务管的是“消息是否成功发送到Broker”,消费端的事务管的是“消息是否被正确消费掉”。你用生产端的事务去约束消费端的行为,就像你在北京店门口贴了张“本店已付款”的条子,结果上海店照样给你发货,两个店根本不搭边。所以“生产端提交了事务,消费端不应该重复”这个直觉是完全站不住脚的。

三、ActiveMQ的投递保障机制

3.1 AUTO_ACKNOWLEDGE 到底确认了什么

ActiveMQ里使用JMS接口时,你会遇到一个“确认模式”参数。最常见的是AUTO_ACKNOWLEDGE,意思是“自动确认”。但自动确认并不等于“每次都确认”。在同步receive()的场景下,消息返回给你之后会立刻自动确认;但在异步onMessage()的场景下,要等监听器方法返回后才会确认。假如你的业务方法抛了异常,或者进程突然挂掉,这个确认动作就没完成,消息会被ActiveMQ重新放回队列,等下一位消费者再收。这就是重复的源头之一。

3.2 客户端手动确认的坑

还有另一种模式叫CLIENT_ACKNOWLEDGE,就是“手动确认”。你收到消息后,必须主动调用消息的acknowledge()方法。很多新手容易踩坑:收完消息,处理完业务,但忘了确认。结果服务一重启,所有处理过但没有确认的消息,全部重新来一遍。这不能怪ActiveMQ,要怪没有遵守它的游戏规则。所以在讨论“重复”的时候,先别急着怀疑中间件,先看看自己消费端的确认逻辑是不是完整的。

四、那个“重复”到底是谁造成的

4.1 生产端事务与消费端事务是两回事

我们来理一理一条消息完整的一生。生产端把消息发给Broker,Broker存好,然后找到消费端,把消息推过去。这可以拆成两段独立的事务边界。第一段:生产端到Broker,由生产端的Session事务负责。第二段:Broker到消费端,由消费端的Session事务或确认模式负责。这两段是串联的,A段成功,B段失败,消息不会消失,只会在B段重新投递。你看到的“重复”,绝大多数发生在第二段。也就是说,问题不在生产端,而在消费端没有把事务或确认做完。

4.2 消费端事务回滚之后会发生什么

如果消费端开启了事务,那么只有调用了session.commit(),消息才算真正被消费掉。要是你调用了session.rollback(),或者事务没提交连接就断了,ActiveMQ会认为这个消息还存在,于是继续投递。更让人难受的是,如果这个过程中消费端已经把业务数据写入了数据库,那就变成了“业务执行了,但消息没确认”的尴尬局面。等消息重新投递时,业务逻辑又被执行一遍,于是出现了重复。所以消费端的事务必须非常小心,一定要保证“先提交事务,再完成外部操作”,或者反过来做幂等。

五、实战示例:复现“重复消费”

技术栈:Java(使用 ActiveMQ Client 5.18.2 + Maven)。

先看项目依赖,只需要一个客户端包。在pom.xml中加入下面的内容。

<!-- pom.xml 依赖片段 -->
<dependency>
    <groupId>org.apache.activemq</groupId>
    <artifactId>activemq-client</artifactId>
    <version>5.18.2</version>
</dependency>

接下来写生产端。注意我开启了会话事务,并在发送完所有消息后调用commit()。这一步非常重要。

import org.apache.activemq.ActiveMQConnectionFactory;

import javax.jms.*;

/**
 * 生产端:用事务来发送消息
 */
public class TxProducer {

    public static void main(String[] args) throws Exception {
        // 连接本机 ActiveMQ
        ActiveMQConnectionFactory factory =
                new ActiveMQConnectionFactory("tcp://127.0.0.1:61616");
        Connection connection = factory.createConnection();
        connection.start();

        // 第一个参数 true 表示开启事务
        // 第二个参数就随意了,因为事务会话不使用确认模式
        Session session = connection.createSession(true, Session.AUTO_ACKNOWLEDGE);
        Queue queue = session.createQueue("TEST.QUEUE");
        MessageProducer producer = session.createProducer(queue);

        for (int i = 1; i <= 3; i++) {
            TextMessage message = session.createTextMessage("订单支付成功-" + i);
            producer.send(message);
            System.out.println("发送消息: " + message.getText());
        }

        // 注意:只有 commit 之后,Broker 才会真正收到这批消息
        session.commit();
        System.out.println("事务已提交,消息已经进入 Broker");

        session.close();
        connection.close();
    }
}

现在写消费端,故意“不提交事务”。这是复现重复的关键。

import org.apache.activemq.ActiveMQConnectionFactory;

import javax.jms.*;

/**
 * 消费端:开启事务,但故意不提交,模拟消费端“假成功”
 */
public class DuplicateConsumer {

    public static void main(String[] args) throws Exception {
        ActiveMQConnectionFactory factory =
                new ActiveMQConnectionFactory("tcp://127.0.0.1:61616");
        Connection connection = factory.createConnection();
        connection.start();

        // 开启事务会话
        Session session = connection.createSession(true, Session.AUTO_ACKNOWLEDGE);
        Queue queue = session.createQueue("TEST.QUEUE");
        MessageConsumer consumer = session.createConsumer(queue);

        System.out.println("开始接收消息,接收后不提交事务...");

        // 用同步接收,方便演示
        for (int i = 0; i < 3; i++) {
            Message message = consumer.receive(2000);
            if (message == null) {
                continue;
            }
            TextMessage text = (TextMessage) message;
            System.out.println("收到消息: " + text.getText());

            // 注释掉 commit,假装业务处理完了,但事务没有提交
            // session.commit();
        }

        // 直接关闭连接,未提交的消息会被 ActiveMQ 重新投递
        System.out.println("没有提交事务,直接关闭连接,消息会回到队列");
        session.close();
        connection.close();
    }
}

运行结果很清晰:第一次启动消费者,会打印收到三条消息,但因为没有提交事务,这三条消息重新回到队列。第二次启动消费者,又会收到完全相同三条消息。在我们这个例子里,生产端的事务明明已经提交了,但消费者依然收到重复。这就证明了重复的根源在消费端,而不在发送端。

下面给出一个纠正版的消费端,在业务成功后提交事务。

import org.apache.activemq.ActiveMQConnectionFactory;

import javax.jms.*;

/**
 * 正确消费端:成功处理之后提交事务
 */
public class CorrectConsumer {

    public static void main(String[] args) throws Exception {
        ActiveMQConnectionFactory factory =
                new ActiveMQConnectionFactory("tcp://127.0.0.1:61616");
        Connection connection = factory.createConnection();
        connection.start();

        Session session = connection.createSession(true, Session.AUTO_ACKNOWLEDGE);
        Queue queue = session.createQueue("TEST.QUEUE");
        MessageConsumer consumer = session.createConsumer(queue);

        for (int i = 0; i < 3; i++) {
            Message message = consumer.receive(2000);
            if (message == null) {
                continue;
            }
            TextMessage text = (TextMessage) message;
            System.out.println("处理消息: " + text.getText());

            // 业务处理成功后,提交事务
            session.commit();
            System.out.println("事务已提交,消息被正确消费");
        }

        session.close();
        connection.close();
    }
}

这个版本再跑,消息就不会重复出现了。注意:提交事务的时机要放在业务成功之后,如果业务抛异常,应该调用session.rollback()。

六、怎么避免这种重复

6.1 让消费端真正提交事务

第一步就是检查消费端代码,确保一定调用commit()或者acknowledge()。尤其是用异步监听器的时候,别把异常吞了。如果onMessage()里抛出异常,事务会话要回滚;如果你在catch块里把异常吞掉,然后继续执行,那消息可能被错误地确认或错误地重投,都非常麻烦。推荐的写法是:业务成功,commit;业务失败,rollback。不要什么都不做。

6.2 开启幂等控制

不过,即便你正确地提交了事务,也依然有极小概率重复。因为网络抖动、Broker崩溃、消费端在提交事务之后但还没告诉Broker之前宕机,都可能导致消息被重投。所以最稳妥的做法是,让消费端天生支持重复。怎么支持?就是幂等。幂等的意思很简单:同一个业务操作,执行一次和执行一百次,结果一样。比如扣库存时,不是简单把库存减一,而是“先把订单状态改成已处理,再减库存”,并利用唯一约束保证同一个订单只处理一次。这样重复消息来了也无所谓。

6.3 合理配置重发策略

ActiveMQ提供了重发策略,可以限制一条消息最多重发几次,以及重发间隔。设置一个合理的值,比如最多重发3次,间隔2秒。如果连续失败,就转到死信队列,人工处理。这样可以避免一条坏消息无限重发,把系统拖垮。在代码里可以这样配置:

import org.apache.activemq.ActiveMQConnectionFactory;
import org.apache.activemq.RedeliveryPolicy;

public class RedeliveryConfig {

    public static ActiveMQConnectionFactory createFactory() {
        ActiveMQConnectionFactory factory =
                new ActiveMQConnectionFactory("tcp://127.0.0.1:61616");

        RedeliveryPolicy policy = new RedeliveryPolicy();
        policy.setMaximumRedeliveries(3);       // 最多重发 3 次
        policy.setInitialRedeliveryDelay(1000);  // 第一次重发延迟 1 秒
        policy.setUseExponentialBackOff(true);   // 是否使用指数退避
        policy.setBackOffMultiplier(2);          // 退避倍数,例如 1s、2s、4s

        factory.setRedeliveryPolicy(policy);
        return factory;
    }
}

注意,重发策略只对“同一条消息被重新投递”有效,它不能防止重复,只是限制重复的次数,给系统一个止损的边界。

七、关联技术:消息幂等性与去重表

7.1 用业务ID去重

最常见且可靠的幂等方案是“业务唯一ID + 去重表”。比如订单支付成功后,生产端在消息头里塞一个唯一的业务ID,消费端在处理前先查一下“消息去重表”,如果已经存在就说明处理过了,直接跳过。这个去重表必须有唯一索引,防止并发下两条一模一样的消息同时进来。

7.2 用Redis实现简单幂等

如果不想引入额外的数据库表,用Redis也能做。把业务ID作为key,用SETNX命令设置成功就说明第一次处理,设置失败说明重复。下面是一个简单示例。

技术栈:Java(Redis 客户端使用 Jedis)。

import redis.clients.jedis.Jedis;

/**
 * 消息去重工具
 */
public class MsgDedup {

    private static final String DEDUP_KEY_PREFIX = "msg:";

    /**
     * 尝试标记一个业务ID为已处理
     * @param businessId 业务唯一ID
     * @return true 表示第一次出现,可以处理;false 表示重复
     */
    public boolean tryProcess(String businessId) {
        try (Jedis jedis = new Jedis("127.0.0.1", 6379)) {
            String key = DEDUP_KEY_PREFIX + businessId;
            Long result = jedis.setnx(key, "1");
            if (result == 1L) {
                // 设置过期时间,避免key永久堆积
                jedis.expire(key, 3600);
                return true;
            }
            return false;
        }
    }
}

使用方式:在消费端收到消息时,取出消息里的业务ID,先调用tryProcess(),如果返回true就执行业务,否则直接忽略。注意,这个方案有一个小缺陷:如果业务执行失败,但ID已经写进Redis了,后续重试就会被误判为重复。所以更严谨的做法是在业务成功后再标记,或者把标记和业务放在同一个本地事务里。这也是上面说的“边界陷阱”的另一种体现。

八、应用场景、技术优缺点与注意事项

8.1 什么时候该用JMS事务

如果你的业务里,写数据库和发消息必须保持原子性,那就应该用事务。比如支付成功后,既要更新订单库里的支付状态,又要发一条“支付成功”消息。如果数据库更新成功但消息没发出去,用户就收不到通知;如果消息发出了但数据库更新失败,用户可能以为自己没支付成功。用JMS事务可以把这两个动作绑定在一起。它的优点是强一致,缺点也明显:事务会锁会话、增加网络开销,吞吐量会下降。所以不要滥用。

8.2 什么时候不该用事务

如果业务对实时性、性能要求很高,并且能接受偶尔重复,那么推荐使用非事务加AUTO_ACKNOWLEDGE,再配合幂等方案。这样做吞吐量大很多,结构也简单。实际上在绝大多数业务场景里,“幂等”比“事务”更能解决问题。

8.3 注意事项清单

我这里再强调几个必须记住的点:

  • JMS事务分为生产端事务和消费端事务,两者互相独立。
  • 生产端提交事务,不代表消费者不会收到重复消息。
  • 消费端事务中,必须显式提交或回滚,不能既不开枪也不卸弹。
  • 只要消息队列提供的是“至少一次”语义,就要默认消息可能重复。
  • 不要把“去重”的希望寄托在消息队列上,要在自己的业务系统里做幂等。
  • 配置重发策略和死信队列,可以防止错误消息无限循环。
  • 使用Redis做幂等时,注意业务失败后要清理标记,或者用更可靠的方式记录状态。

九、文章总结

回到最初的问题:为什么ActiveMQ事务消息提交后,消费者仍然收到重复?根本原因不是ActiveMQ坏了,也不是你用了假事务,而是你只看到了生产端的事务,没意识到消费端还有一条独立的确认边界。生产端提交事务只代表消息成功送达Broker,消费者那侧该确认还是要确认,该提交还是要提交。只有消费端正确提交事务或确认消息,才不会再收到同一条消息。即使如此,由于分布式环境的复杂性和消息队列“至少一次”的投递语义,重复依然可能发生。因此,防御重复的最终武器永远是幂等设计。记住了这个,以后再遇到类似问题,你就可以很淡定地告诉同事:不是消息队列的事,是我们的边界又踩坑了。