运营找上门:同一笔积分被发了两次
七月的一个下午,运营拿着工单来找我:"用户 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 死信队列,我们在控制台配了告警,死信队列有消息就发钉钉。
小结
三句话:
- RocketMQ 只保证至少一次,重复是常态不是异常。别指望配置能关掉它,只能在消费端做幂等。
- 幂等键必须是业务键(订单号、流水号),不能是 msgId——生产端重发会绕过它。
- 去重表和业务操作要同库同事务,否则会造出"去重成功但业务没做"的消息丢失。
那次投诉最后处理方式:写了个脚本扫全表,按 (user_id, order_no) 分组找出重复记录,把多发的 300 积分冲正回来。一共查出 47 笔,全是 7 月 3 号那次 Broker 主备切换期间产生的。