消息重复投递如何根治?activemq4cj的ActiveMQMessageAudit幂等审核机制
【免费下载链接】activemq4cj仓颉语言实现的ActiveMQ客户端SDK。遵行JMS规范,支持OpenWire协议,支持点对点和发布订阅模式,支持失效转移。当前main分支适配仓颉1.0.0 LTS版本,分支develop适配仓颉0.53.4 Beta版本,分支Branch_cj0.60.5适配仓颉0.60.5 Beta版本。项目地址: https://gitcode.com/Cangjie-TPC/activemq4cj
做消息系统的人都遇到过同一个噩梦:网络抖动、Broker 故障转移、ACK 丢失……消息被投递了两次,业务逻辑被执行了两次,订单重复扣款、日志重复入账。activemq4cj 是用仓颉语言实现的 ActiveMQ 客户端 SDK,遵循 JMS 规范并支持 OpenWire 协议。它在消费端内置了ActiveMQMessageAudit 幂等审核机制——通过“生产者种子 + 序列号 + 滑动位窗口”的方式识别并抑制重复投递,从客户端层面把消息重复投递问题治得干干净净。
为什么消息必然会被重复投递?
先理解问题本质。消息中间件普遍提供at-least-once(至少一次)投递保证:只要消息在 Broker 确认之前连接断开,重连后 Broker 就会把“可能没送达”的消息再发一遍。
在 activemq4cj 中,这种情况尤其常见,因为它支持failover 失效转移(见 failover_transport.cj)——主节点宕机自动切备用节点时,正在途的消息天然存在“重放”风险。也就是说:重复投递不是异常,而是常态,客户端必须有能力识别它。
审核的核心原理:种子 + 序列号 + 滑动窗口
📌 每条消息都有一个全局唯一 ID,格式为种子:序列号。IdGenerator 生成种子(包含主机名、随机数、时间戳,保证每个生产者唯一),序列号是生产者内的自增计数。
ActiveMQMessageAuditNoSync 的isDuplicate逻辑非常精巧:
- 按生产者分桶:用
MessageId中的producerId作为键,查一张 LRUCache 缓存,每个生产者对应一个位图容器; - 位图打标:取出该生产者的 BitSetBin(默认窗口大小2048个序列号),对
producerSequenceId对应的位执行setBit; - 判定重复:
setBit会返回该位原先是否已被置位——原来就是 1,说明这条消息(或同序号消息)已经处理过,直接判定为重复; - 窗口滑动:序列号超出窗口时,
BitSetBin自动丢弃最旧的一组位、追加新的组,内存占用恒定可控。
这个设计最妙的地方在于:它不需要存储消息本身,每个生产者仅消耗几 KB 的位图,就能在海量消息流中做到 O(1) 的重复检测。
线程安全方面,SDK 提供两个版本:ActiveMQMessageAuditNoSync是无锁版本(单线程消费场景),ActiveMQMessageAudit 在其上加了Mutex同步,多消费线程共享时也能保证判定正确。
队列与 Topic 的差异化组织:ConnectionAudit
ConnectionAudit 是审核机制的连接级管理者,它对两种目的地采用不同策略(对应 JMS 语义差异):
- Queue(点对点):按目的地维护一份
ActiveMQMessageAudit。同一条消息无论被多少个消费者竞争,全队列只应处理一次; - Topic(发布订阅):按dispatcher(消费者)维护审核窗口。订阅者各自独立,A 消费者收到过不算 B 消费者的重复;
- 两类审核窗口各自缓存在 LRU 中(上限 1000 份),目的地/消费者太多时自动淘汰最久未用的,防止内存膨胀。
重复消息被拦截后的处理
审核判定只是第一步,ActiveMQMessageConsumer 在dispatch分发时会拦截重复消息并按场景分流:
| 场景 | 处理方式 |
|---|---|
| 重复投递发生在当前事务中 | 记录为事务内重投递,正常提交后生效 |
| 存在竞争中的挂起事务 | rollbackDuplicate回滚审核位 + 重新分发,由事务裁决 |
| 普通连接上的重复 | poisonAck毒确认,直接抑制该投递并告警日志 |
同时 SDK 提供rollback回滚接口:当消息处理失败触发重投递时,先把审核位清零,让同序号消息能重新通过审核——这保证了“失败重投”与“网络重放”能被精确区分,互不干扰。
开关与调参:何时开启审核
审核机制并非无条件运行,而是与传输层的faultTolerant(容错)属性联动:在 activemq_connection.cj 中,connectionAudit.checkForDuplicates直接绑定传输层是否容错——failover 场景自动开启,普通直连场景可通过 activemq_connection_factory.cj 的checkForDuplicates属性手动控制,兼顾性能与安全。
两个可调参数,按需伸缩:
auditDepth:滑动窗口深度,默认 2048。队列积压越大、重放距离越远,应调大;auditMaximumProducerNumber:同时追踪的生产者数,默认 64。多生产者并发写入同一目的地时建议调大。
小结
activemq4cj 的 ActiveMQMessageAudit 用不到百行的核心逻辑,把“消息重复投递”这个分布式经典难题收敛为一次位图查询:
- 按生产者分桶 + 滑动位窗口:O(1) 判重,内存恒定;
- Queue/Topic 双策略审核:精准匹配 JMS 投递语义;
- 审核回滚 + 毒确认:失败重投与网络重放互不冲突;
- 与 failover 自动联动:容错连接开箱即防重复。
如果你在用仓颉语言构建消息消费端,这套幂等审核机制值得直接借鉴。更多消息类型示例(文本、对象、二进制等)可参考 samples/ 目录与 ActiveMQ_SDK_User_Guide.md,单元测试 activemq_message_audit_test.cj 演示了判重与回滚的完整行为。
【免费下载链接】activemq4cj仓颉语言实现的ActiveMQ客户端SDK。遵行JMS规范,支持OpenWire协议,支持点对点和发布订阅模式,支持失效转移。当前main分支适配仓颉1.0.0 LTS版本,分支develop适配仓颉0.53.4 Beta版本,分支Branch_cj0.60.5适配仓颉0.60.5 Beta版本。项目地址: https://gitcode.com/Cangjie-TPC/activemq4cj
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考