分布式系统中消息重复投递无法完全避免——网络抖动导致消费者ACK丢失、Broker端重试机制、消费者端宕机重启都会产生重复消息。消费者端必须实现幂等性处理,保证同一消息被消费多次时结果一致。本文对比四种幂等方案,给出RocketMQ场景下的完整实现代码。
RocketMQ重复消息的产生场景与消费模型
RocketMQ默认保证At-Least-Once(至少一次)语义,同一条消息可能被投递多次。常见重复场景:
消费者处理消息耗时超过 consumeTimeout(默认15分钟),Broker认为消费失败重新投递。消费者端处理完成后提交offset前宕机,重启后从上次offset重新拉取。消费者本地重试逻辑 maxReconsumeTimes 触发后消息进入死信队列前被重复消费。
RocketMQ消费模式分两种:集群消费(CLUSTERING,默认)下同一消息只会被一个消费者实例消费;广播消费(BROADCASTING)下每个消费者实例都消费全量消息。幂等问题主要集中在集群消费模式。
方案一:数据库唯一索引去重
利用数据库唯一索引约束,将消息ID作为唯一键。插入成功表示首次消费,插入失败表示重复消息直接跳过:
@Component
@RocketMQMessageListener(
topic = "order_topic",
consumerGroup = "order_consumer_group",
messageModel = MessageModel.CLUSTERING
)
public class OrderConsumer implements RocketMQListener<OrderMessage> {
@Autowired
private JdbcTemplate jdbcTemplate;
@Override
@Transactional
public void onMessage(OrderMessage message) {
// Step 1: 插入消费记录,唯一键约束防止重复
try {
jdbcTemplate.update(
"INSERT INTO msg_consume_log (msg_id, consumer_group, consume_time) " +
"VALUES (?, ?, NOW()) " +
"ON DUPLICATE KEY UPDATE consume_time = consume_time",
message.getMsgId(), "order_consumer_group"
);
} catch (DuplicateKeyException e) {
// 已消费过,直接返回
return;
}
// Step 2: 执行业务逻辑
processOrder(message);
}
private void processOrder(OrderMessage message) {
// 订单处理逻辑
jdbcTemplate.update(
"UPDATE orders SET status = ? WHERE order_id = ?",
message.getStatus(), message.getOrderId()
);
}
}
优点:实现简单,与数据库事务天然整合。缺点:每条消息插入一条记录,高频消费场景下数据库写入压力大。
方案二:Redis SETNX标记去重
消费前在Redis中设置消息ID标记,利用SETNX(SET IF NOT EXISTS)的原子性判断是否已消费:
@Component
public class OrderConsumerRedis implements RocketMQListener<OrderMessage> {
@Autowired
private StringRedisTemplate redisTemplate;
@Autowired
private OrderService orderService;
private static final String CONSUME_KEY_PREFIX = "mq:consumed:";
private static final long EXPIRE_HOURS = 48L;
@Override
public void onMessage(OrderMessage message) {
String key = CONSUME_KEY_PREFIX + message.getMsgId();
Boolean firstConsume = redisTemplate.opsForValue()
.setIfAbsent(key, "1", EXPIRE_HOURS, TimeUnit.HOURS);
if (Boolean.FALSE.equals(firstConsume)) {
// 已消费过,直接跳过
return;
}
try {
orderService.processOrder(message);
} catch (Exception e) {
// 业务处理失败,删除标记以便重试
redisTemplate.delete(key);
throw e; // 抛出异常触发RocketMQ重试
}
}
}
关键点:设置过期时间(48小时)避免Redis内存无限增长。业务处理失败时删除标记,确保下次重试能正常消费。过期时间需大于RocketMQ的最大重试间隔,否则过期后重试消息会被当作首次消费。
方案三:业务状态机幂等
不依赖外部存储,直接利用业务状态流转的不可逆性。订单状态从”待支付”到”已支付”是单向的,重复消息尝试从”已支付”改为”已支付”时SQL影响行数为0:
@Service
public class OrderService {
@Autowired
private OrderMapper orderMapper;
public void processOrder(OrderMessage message) {
// 状态机校验:只有"待支付"状态才能改为"已支付"
int affectedRows = orderMapper.updateStatusWithCondition(
message.getOrderId(),
OrderStatus.PAID, // 目标状态
OrderStatus.PENDING // 前置状态(条件)
);
if (affectedRows == 0) {
// 状态不满足条件(可能已处理过),直接返回
log.info("Order {} already processed, skip", message.getOrderId());
return;
}
// 执行后续逻辑(扣库存、发通知等)
deductInventory(message.getOrderId());
}
}
对应Mapper XML的条件更新:
<update id="updateStatusWithCondition">
UPDATE orders
SET status = #{targetStatus}, update_time = NOW()
WHERE order_id = #{orderId} AND status = #{sourceStatus}
</update>
状态机幂等不需要维护额外的消费日志表或Redis标记,但要求业务流程具有明确的状态流转模型。
方案四:分布式锁串行化处理
同一业务ID的消息通过分布式锁串行化执行,锁内完成判断和处理:
@Component
public class OrderConsumerLock implements RocketMQListener<OrderMessage> {
@Autowired
private RedissonClient redissonClient;
@Autowired
private OrderService orderService;
@Override
public void onMessage(OrderMessage message) {
String lockKey = "lock:order:" + message.getOrderId();
RLock lock = redissonClient.getLock(lockKey);
try {
// 尝试加锁,等待3秒,自动释放10秒
boolean locked = lock.tryLock(3, 10, TimeUnit.SECONDS);
if (!locked) {
throw new RuntimeException("获取锁失败,稍后重试");
}
// 锁内查询业务状态
Order order = orderService.getById(message.getOrderId());
if (order.getStatus() != OrderStatus.PENDING) {
return; // 已处理过
}
orderService.processOrder(message);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
if (lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
}
}
四种方案对比与选型建议
数据库唯一索引:适合消费频率低、已有数据库的场景,零额外基础设施。 Redis SETNX:适合高并发消费,Redis作为基础设施已存在时优先考虑。 状态机幂等:适合业务流程清晰、状态流转明确的核心交易场景,无副作用且最可靠。 分布式锁:适合同一业务ID需要严格串行处理的场景,如库存扣减,但锁竞争是性能瓶颈。
生产环境中通常组合使用:核心交易链路用状态机幂等作为兜底保障,Redis标记减少无效处理,数据库唯一索引作为最终防线。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/rocketmq-xiao-xi-xiao-fei-zhe-mi-deng-she-ji-chong-fu-xiao/