Administrator
发布于 2020-07-05 / 3754 阅读
44

RocketMQ 消息重复消费与幂等处理

运营找上门:同一笔积分被发了两次

七月的一个下午,运营拿着工单来找我:"用户 88372 投诉,说下单后积分到账两次,多领了 300 积分。"

我第一反应是代码里有循环或者被调用了两次。翻了消费日志,同一个 msgId 确实处理了两遍:

2020-07-05 14:22:31.117 [ConsumeMessageThread_3] INFO  PointConsumer - 开始处理积分, msgId=0A1B2C3D00002A9F0000000000012E45, orderNo=SO202007051422310088
2020-07-05 14:22:31.240 [ConsumeMessageThread_3] INFO  PointConsumer - 积分发放成功, userId=88372, point=300
2020-07-05 14:22:33.502 [ConsumeMessageThread_7] INFO  PointConsumer - 开始处理积分, msgId=0A1B2C3D00002A9F0000000000012E45, orderNo=SO202007051422310088
2020-07-05 14:22:33.615 [ConsumeMessageThread_7] INFO  PointConsumer - 积分发放成功, userId=88372, point=300

两条日志间隔 2.4 秒,msgId 一模一样。这就是 RocketMQ 的 at least once 语义。

为什么一定会重复

RocketMQ 4.7 的投递模型是至少一次:消息必须被消费成功,否则会重试。但"消费成功"这个信号从 Consumer 回到 Broker 有可能丢,或者 Consumer 处理完了但还没来得及返回就挂了。Broker 等不到响应,就认为没消费,再投一次。

具体有这几种触发场景,我基本都踩过:

  • 消费超时consumeTimeout 默认 15 分钟,但这期间消息在 Broker 看来是"处理中"。Consumer 重启、或者消费线程被卡住,消息会被重新投递。
  • 返回 RECONSUME_LATER:业务异常后我手动返回了重试,消息进入 %RETRY% 队列,延迟一段时间再投。但"业务已经做了一半"这个事实不会回滚。
  • 消费者 rebalance:扩容或者某台机器 GC 卡顿被踢出消费组,队列重新分配。新消费者不知道上一条消息处理到哪了,会从上次提交的 offset 开始,把中间没提交的部分重新拉一遍。
  • 生产者重发:发送超时(比如网络闪断)时 Producer 会重试,Broker 上可能已经落了两条内容相同但 msgId 不同的消息。

所以靠 msgId 去重是不够的——生产端重发产生的两条消息 msgId 不同,但业务内容一样。去重必须落在业务键上。

方案一:数据库唯一索引兜底(最简单)

先建一张去重表,把业务键做成唯一索引:

CREATE TABLE t_point_record (
    id          BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
    order_no    VARCHAR(32) NOT NULL COMMENT '订单号,幂等键',
    user_id     BIGINT      NOT NULL,
    point       INT         NOT NULL,
    create_time DATETIME    NOT NULL DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (id),
    UNIQUE KEY uk_order_no (order_no)
) ENGINE = InnoDB;

消费时先插去重表,插入成功才发积分,两步放在同一个事务里:

@Transactional(rollbackFor = Exception.class)
public void handle(String orderNo, Long userId, int point) {
    try {
        pointRecordMapper.insert(new PointRecord(orderNo, userId, point));
    } catch (DuplicateKeyException e) {
        log.warn("重复消息,已忽略, orderNo={}", orderNo);
        return;                       // 直接返回成功,让 Broker 别再重试
    }
    accountMapper.addPoint(userId, point);
    pointDetailMapper.insert(...);
}

这个方案的关键是去重记录和业务操作必须在同一个本地事务。如果先插去重表再单独发积分,中间宕机就会出现"去重表有记录但积分没到账",而重试又被去重表挡掉,变成消息丢失。我们第一次改的时候就是这个 bug,测试环境跑了两天才发现。

方案二:Redis SETNX 去重(性能更好)

去重表每次消费多一次 DB 写入,积分这块量不大(日均 40 万条)还扛得住。但另一个优惠券发放的场景 QPS 高,我换成了 Redis:

private static final String DEDUP_KEY = "mq:dedup:point:";

public void handle(String orderNo, Long userId, int point) {
    String key = DEDUP_KEY + orderNo;
    // 24 小时足够覆盖 RocketMQ 的最大重试周期(16 次,最长约 4 小时 46 分)
    Boolean ok = stringRedisTemplate.opsForValue()
            .setIfAbsent(key, "1", 24, TimeUnit.HOURS);
    if (!Boolean.TRUE.equals(ok)) {
        log.warn("重复消息, orderNo={}", orderNo);
        return;
    }
    try {
        doSendPoint(userId, point);       // 真正的业务
    } catch (Exception e) {
        stringRedisTemplate.delete(key);  // 业务失败要清掉标记,允许重试
        throw e;
    }
}

这里有个细节要注意:业务失败时必须把去重标记删掉,否则消息重试时会被自己的去重挡住,直接丢消息。但删除本身可能失败,所以更稳妥的做法是不删,而是把值设成一个"处理中"状态,重试时判断上次是不是完成了。我们的做法更简单——业务失败就抛异常让 RocketMQ 重试,重试次数用完后进死信队列人工处理。

方案三:update 型业务用状态机

还有一类场景根本不需要去重表:订单状态流转。这种"改状态"的操作可以用带条件的 update 天然幂等。

<update id="paySuccess">
    UPDATE t_order
    SET status = 'PAID', pay_time = #{payTime}, update_time = NOW()
    WHERE order_no = #{orderNo}
      AND status = 'WAIT_PAY'      /* 只有待支付才能变成已支付 */
</update>

第一条消息把状态改成 PAID,返回影响行数 1;重复的第二条再执行,status = 'WAIT_PAY' 不成立,影响行数 0。消费端判断一下:

int rows = orderMapper.paySuccess(orderNo, payTime);
if (rows == 0) {
    log.info("订单状态已流转,忽略重复消息, orderNo={}", orderNo);
}
// 无论 rows 是 0 还是 1,都返回 CONSUME_SUCCESS

这个"乐观锁 + 状态前置条件"的写法是最省事的,适合所有 update 场景。但它对 insert 型业务(发积分、发券、加流水)无效,那种还是得靠去重表。

关于消费端幂等的一些配置

顺便把我们 RocketMQ 4.7 消费端的配置贴一下,有几个参数是踩坑之后才调的:

@Bean
public DefaultRocketMQListenerContainer pointContainer() {
    DefaultRocketMQListenerContainer container = new DefaultRocketMQListenerContainer();
    container.setConsumerGroup("point-consumer-group");
    container.setTopic("ORDER_TOPIC");
    container.setSelectorExpression("PAID || FINISHED");   // tag 过滤
    container.setConsumeThreadMax(20);
    container.setConsumeThreadMin(8);
    container.setMessageModel(MessageModel.CLUSTERING);    // 集群消费
    container.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
    return container;
}

消费逻辑里别吞异常,返回 RECONSUME_LATER 让 Broker 重新投:

@RocketMQMessageListener(
        topic = "ORDER_TOPIC",
        selectorExpression = "PAID",
        consumerGroup = "point-consumer-group")
public class PointConsumer implements RocketMQListener<MessageExt> {
    @Override
    public void onMessage(MessageExt msg) {
        String orderNo = new String(msg.getBody(), StandardCharsets.UTF_8);
        pointService.handle(orderNo, ...);
        // 抛异常会自动返回 RECONSUME_LATER
    }
}

重试次数用默认 16 次,间隔是 10s、30s、1m、2m、3m、4m……逐级拉长,最后一次间隔 2 小时。超过 16 次进 %DLQ%point-consumer-group 死信队列,我们在控制台配了告警,死信队列有消息就发钉钉。

小结

三句话:

  1. RocketMQ 只保证至少一次,重复是常态不是异常。别指望配置能关掉它,只能在消费端做幂等。
  2. 幂等键必须是业务键(订单号、流水号),不能是 msgId——生产端重发会绕过它。
  3. 去重表和业务操作要同库同事务,否则会造出"去重成功但业务没做"的消息丢失。

那次投诉最后处理方式:写了个脚本扫全表,按 (user_id, order_no) 分组找出重复记录,把多发的 300 积分冲正回来。一共查出 47 笔,全是 7 月 3 号那次 Broker 主备切换期间产生的。

参考