RocketMQ消息消费者幂等设计:重复消息四种消除方案与代码实现

分布式系统中消息重复投递无法完全避免——网络抖动导致消费者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/

(0)
小编小编
上一篇 9小时前
下一篇 9小时前

相关推荐