找回密码
立即注册
搜索
发回帖 发新帖

6284

积分

0

好友

793

主题
发表于 16 小时前 | 查看: 4| 回复: 0

一、一个真实的生产事故

某电商平台上线了“下单送积分”功能:用户支付成功后,发送 MQ 消息,积分服务收到消息后给用户加积分。

上线第三天,客服接到投诉:“我买了1单,怎么加了2次积分?”

排查发现,MQ 消息被重复消费了。原因有两个:

原因一:生产者重试导致重复投递

订单服务发送消息 → 网络抖动,Broker未确认
    ↓
订单服务重试 → 消息被投递两次

原因二:消费者 ACK 丢失导致重投

消费者处理成功 → 调用ACK → 网络抖动,ACK丢失
    ↓
Broker认为消息未处理 → 重新投递
    ↓
消费者再次处理 → 重复加积分

事故复盘:

问题 说明
消息重复不可避免 MQ 只保证 At Least Once,不保证 Exactly Once
无幂等设计 消费逻辑没有做去重,重复消息被重复执行
无唯一约束 数据库没有唯一索引兜底
无监控 重复消费发生了3天才被发现

二、问题根源:为什么消息一定会重复?

2.1 MQ的三种投递语义

语义 含义 实现难度 适用场景
At Most Once 最多一次,可能丢 简单 日志采集等允许丢失的场景
At Least Once 至少一次,可能重复 中等 绝大多数业务场景
Exactly Once 恰好一次 极难 金融级场景(通常通过幂等实现)

主流 MQ(RabbitMQ、RocketMQ、Kafka)都只能保证 At Least Once。

2.2 消息重复的三个来源

来源一:生产者重试

生产者发送消息 → 未收到Broker确认 → 重试发送
→ 消息可能已经在Broker中,重试导致重复

来源二:消费者 ACK 丢失

消费者处理成功 → 发送ACK → 网络抖动,ACK丢失
→ Broker认为消息未处理 → 重新投递
→ 消费者再次收到消息

来源三:消费者处理超时

消费者拉取消息 → 处理时间超过超时时间
→ Broker认为消费者已死 → 重新投递消息给其他消费者
→ 消息被两个消费者同时处理

2.3 幂等性 vs 去重的区别

概念 说明
去重 通过某种方式识别重复消息,丢弃重复的
幂等 保证同一操作执行多次的结果与执行一次相同

去重是手段,幂等是目标。但很多场景下,去重就等于幂等。

2.4 幂等消费的核心原则

不要试图消除重复,而要容忍重复。

既然重复不可避免,正确的做法是让消费逻辑本身具有幂等性。

三、核心技术原理

3.1 幂等消费的四种方案

方案 实现方式 优点 缺点 适用场景
Redis SET NX 消息ID写入 Redis,重复则跳过 性能高 依赖 Redis 高并发场景
数据库唯一索引 业务唯一ID插入唯一表 强一致 性能略低 强一致场景
状态机 检查业务状态,已处理则跳过 无额外存储 需要状态设计 状态流转场景
乐观锁 版本号控制 无锁 需要版本字段 更新场景

3.2 Redis SET NX 的原理

// 核心操作:SET message:id:{msgId} 1 NX EX 86400
Boolean success = redisTemplate.opsForValue()
        .setIfAbsent(key, "1", 24, TimeUnit.HOURS);

if (Boolean.FALSE.equals(success)) {
    // key已存在,说明消息已被处理,跳过
    return;
}
// key不存在,首次处理

三个关键点:

  1. 原子性:SET NX 是原子操作,并发安全
  2. 过期时间:设置 TTL,防止 Redis 内存无限增长
  3. 唯一Key:通常用 消息ID 或 业务ID + 操作类型

3.3 数据库唯一索引的原理

CREATE TABLE message_consume_log (
    id BIGINT AUTO_INCREMENT,
    message_id VARCHAR(64) NOT NULL,
    consume_time DATETIME,
    PRIMARY KEY (id),
    UNIQUE KEY uk_message_id (message_id)  -- 唯一索引
);
try {
    // 插入消费日志,如果messageId已存在则抛异常
    jdbcTemplate.update(
        "INSERT INTO message_consume_log (message_id, consume_time) VALUES (?, NOW())",
        messageId);
    // 插入成功,执行业务逻辑
    processBusiness();
} catch (DuplicateKeyException e) {
    // 唯一索引冲突,说明已消费过
    log.warn("消息已消费: messageId={}", messageId);
}

3.4 组合方案:Redis + 数据库

单一方案的问题:

方案 问题
只用 Redis Redis 故障后无法去重
只用数据库 高并发下数据库压力大

组合方案的优势:

第一层:Redis SET NX → 快速拦截99%的重复消息
第二层:数据库唯一索引 → 兜底,即使Redis故障也能保证幂等

3.5 状态机的幂等设计

很多业务场景天然支持幂等——通过状态检查:

待支付 → 已支付 → 已发货 → 已完成

消费"支付成功"消息时:
  当前状态 = 待支付 → 更新为已支付 ✅
  当前状态 = 已支付 → 跳过(已处理)✅
  当前状态 = 已完成 → 跳过 ✅

核心SQL:

UPDATE orders 
SET status = 1, pay_time = NOW() 
WHERE id = ? AND status = 0;  -- 关键:WHERE条件限制状态

如果订单已经是 status=1,则 WHERE status=0 不匹配,更新影响行数为0,天然幂等。

四、实现方案及场景

4.1 本文采用的方案

三层幂等防护:

层次 实现 作用
第一层 Redis SET NX(消息ID去重) 快速拦截重复消息
第二层 业务状态检查(条件更新) 业务层面幂等
第三层 数据库唯一索引(消费日志) 最终兜底

4.2 整体流程

消费消息
    ↓
① Redis SET NX message:id:{msgId}
    ├── 已存在 → 跳过
    └── 不存在 → 继续
    ↓
② 检查业务状态
    ├── 已处理 → 跳过
    └── 未处理 → 继续
    ↓
③ 执行业务(条件更新 WHERE status = 预期值)
    ↓
④ 插入消费日志(唯一索引兜底)
    ↓
⑤ ACK

4.3 适用场景

场景 说明
订单支付通知 支付成功后更新订单状态
库存扣减 扣减库存需要幂等
积分发放 发放积分需要防重复
优惠券发放 优惠券只能发放一次

五、完整实现

5.1 项目结构

heyou-mq-idempotent/
├── pom.xml
└── src/main/java/com/heyou/idempotent/
    ├── IdempotentApplication.java
    ├── config/
    │   └── RabbitMQConfig.java
    ├── consumer/
    │   ├── OrderPayConsumer.java           # 支付消息消费者
    │   └── PointsConsumer.java             # 积分消息消费者
    ├── service/
    │   └── IdempotentService.java          # 幂等服务
    ├── mapper/
    │   └── OrderMapper.java
    └── entity/
        ├── Order.java
        └── MessageConsumeLog.java

5.2 Maven依赖

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-data-redis</artifactId>
    </dependency>
    <dependency>
        <groupId>org.apache.commons</groupId>
        <artifactId>commons-pool2</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-jdbc</artifactId>
    </dependency>
    <dependency>
        <groupId>com.mysql</groupId>
        <artifactId>mysql-connector-j</artifactId>
        <scope>runtime</scope>
    </dependency>
    <dependency>
        <groupId>org.projectlombok</groupId>
        <artifactId>lombok</artifactId>
        <optional>true</optional>
    </dependency>
</dependencies>

5.3 数据库表

-- 订单表
CREATE TABLE `orders` (
    `id` BIGINT NOT NULL AUTO_INCREMENT,
    `order_no` VARCHAR(64) NOT NULL,
    `user_id` BIGINT NOT NULL,
    `amount` DECIMAL(10,2) NOT NULL,
    `status` TINYINT DEFAULT 0 COMMENT '0-待支付 1-已支付 2-已取消',
    `pay_time` DATETIME DEFAULT NULL,
    `create_time` DATETIME DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (`id`),
    UNIQUE KEY `uk_order_no` (`order_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

-- 消息消费日志(唯一索引兜底)
CREATE TABLE `message_consume_log` (
    `id` BIGINT NOT NULL AUTO_INCREMENT,
    `message_id` VARCHAR(64) NOT NULL,
    `biz_type` VARCHAR(50) NOT NULL,
    `consume_time` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
    PRIMARY KEY (`id`),
    UNIQUE KEY `uk_message_id` (`message_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

-- 积分表(用于验证幂等)
CREATE TABLE `user_points` (
    `id` BIGINT NOT NULL AUTO_INCREMENT,
    `user_id` BIGINT NOT NULL,
    `points` INT NOT NULL DEFAULT 0,
    `update_time` DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
    PRIMARY KEY (`id`),
    UNIQUE KEY `uk_user_id` (`user_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

5.4 幂等服务(核心)

package com.example.idempotent.service;

import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;

import java.util.concurrent.TimeUnit;

@Slf4j
@Service
@RequiredArgsConstructor
public class IdempotentService {

    private final StringRedisTemplate redisTemplate;
    private final JdbcTemplate jdbcTemplate;

    private static final String IDEMPOTENT_PREFIX = "mq:idempotent:";
    private static final long EXPIRE_HOURS = 24;

    /**
     * 幂等性检查(第一层:Redis SET NX)
     * 
     * @param messageId 消息唯一ID
     * @return true-首次处理,false-重复消息
     */
    public boolean checkRedisIdempotent(String messageId) {
        String key = IDEMPOTENT_PREFIX + messageId;
        Boolean success = redisTemplate.opsForValue()
                .setIfAbsent(key, "1", EXPIRE_HOURS, TimeUnit.HOURS);

        if (Boolean.FALSE.equals(success)) {
            log.warn("消息已被消费(Redis拦截): messageId={}", messageId);
            return false;
        }
        return true;
    }

    /**
     * 记录消费日志(第三层:数据库唯一索引)
     * 
     * @param messageId 消息唯一ID
     * @param bizType 业务类型
     * @return true-首次消费,false-重复消费
     */
    public boolean recordConsumeLog(String messageId, String bizType) {
        try {
            jdbcTemplate.update(
                "INSERT INTO message_consume_log (message_id, biz_type) VALUES (?, ?)",
                messageId, bizType);
            return true;
        } catch (DuplicateKeyException e) {
            log.warn("消息已被消费(数据库兜底): messageId={}", messageId);
            return false;
        }
    }

    /**
     * 综合幂等检查(三层防护)
     */
    public boolean checkIdempotent(String messageId, String bizType) {
        // 第一层:Redis快速拦截
        if (!checkRedisIdempotent(messageId)) {
            return false;
        }

        // 第三层:数据库兜底
        if (!recordConsumeLog(messageId, bizType)) {
            // Redis通过了但数据库没通过,说明Redis数据过期
            // 此时应该阻止处理,防止重复
            return false;
        }

        return true;
    }
}

5.5 支付消息消费者

package com.example.idempotent.consumer;

import com.example.idempotent.service.IdempotentService;
import com.rabbitmq.client.Channel;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;

import java.io.IOException;

@Slf4j
@Component
@RequiredArgsConstructor
public class OrderPayConsumer {

    private final IdempotentService idempotentService;
    private final JdbcTemplate jdbcTemplate;

    /**
     * 消费支付成功消息
     * 
     * 三层幂等防护:
     * 1. Redis SET NX(消息ID去重)
     * 2. 业务状态检查(条件更新)
     * 3. 数据库唯一索引(消费日志兜底)
     */
    @RabbitListener(queues = "order.pay.queue")
    public void onPaySuccess(Message message, Channel channel) throws IOException {
        String messageId = message.getMessageProperties().getMessageId();
        String content = new String(message.getBody());
        long deliveryTag = message.getMessageProperties().getDeliveryTag();

        log.info("收到支付消息: messageId={}, content={}", messageId, content);

        try {
            // 第一层 + 第三层:综合幂等检查
            if (!idempotentService.checkIdempotent(messageId, "ORDER_PAY")) {
                channel.basicAck(deliveryTag, false);
                return;
            }

            // 解析消息
            Long orderId = Long.parseLong(content.split(":")[0]);

            // 第二层:业务状态检查(条件更新)
            // 关键:WHERE status = 0 保证只有待支付订单才能被更新
            int rows = jdbcTemplate.update(
                "UPDATE orders SET status = 1, pay_time = NOW() " +
                "WHERE id = ? AND status = 0",
                orderId);

            if (rows > 0) {
                log.info("订单支付成功: orderId={}", orderId);
                // 执行后续逻辑(发短信、发积分等)
            } else {
                // 影响行数为0,说明订单状态不是待支付
                log.warn("订单状态不允许支付(可能已处理): orderId={}", orderId);
            }

            // 手动ACK
            channel.basicAck(deliveryTag, false);

        } catch (Exception e) {
            log.error("支付消息处理失败: messageId={}", messageId, e);
            // 处理失败,拒绝并重新入队
            channel.basicNack(deliveryTag, false, true);
        }
    }
}

5.6 积分消息消费者

package com.example.idempotent.consumer;

import com.example.idempotent.service.IdempotentService;
import com.rabbitmq.client.Channel;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;

import java.io.IOException;

@Slf4j
@Component
@RequiredArgsConstructor
public class PointsConsumer {

    private final IdempotentService idempotentService;
    private final JdbcTemplate jdbcTemplate;

    /**
     * 消费积分发放消息
     * 加积分是"非幂等"操作,必须严格去重
     */
    @RabbitListener(queues = "points.queue")
    public void onGrantPoints(Message message, Channel channel) throws IOException {
        String messageId = message.getMessageProperties().getMessageId();
        String content = new String(message.getBody());
        long deliveryTag = message.getMessageProperties().getDeliveryTag();

        log.info("收到积分消息: messageId={}, content={}", messageId, content);

        try {
            // 幂等检查:积分发放绝对不能重复
            if (!idempotentService.checkIdempotent(messageId, "GRANT_POINTS")) {
                log.warn("积分消息重复,跳过: messageId={}", messageId);
                channel.basicAck(deliveryTag, false);
                return;
            }

            // 解析消息:userId:points
            String[] parts = content.split(":");
            Long userId = Long.parseLong(parts[0]);
            int points = Integer.parseInt(parts[1]);

            // 加积分(使用 ON DUPLICATE KEY 保证幂等)
            jdbcTemplate.update(
                "INSERT INTO user_points (user_id, points) VALUES (?, ?) " +
                "ON DUPLICATE KEY UPDATE points = points + ?",
                userId, points, points);

            log.info("积分发放成功: userId={}, points={}", userId, points);

            channel.basicAck(deliveryTag, false);

        } catch (Exception e) {
            log.error("积分消息处理失败: messageId={}", messageId, e);
            channel.basicNack(deliveryTag, false, true);
        }
    }
}

5.7 配置

spring:
  rabbitmq:
    host: localhost
    port: 5672
    username: guest
    password: guest
    listener:
      simple:
        acknowledge-mode: manual
        prefetch: 10
        concurrency: 5
  redis:
    host: localhost
    port: 6379
  datasource:
    url: jdbc:mysql://localhost:3306/test?useUnicode=true&characterEncoding=UTF-8&useSSL=false
    username: root
    password: 123456

5.8 测试Controller

@RestController
@RequestMapping("/api/test")
@RequiredArgsConstructor
public class TestController {

    private final RabbitTemplate rabbitTemplate;

    /**
     * 模拟发送支付消息(可重复发送验证幂等)
     */
    @PostMapping("/send-pay")
    public String sendPayMessage(@RequestParam Long orderId,
                                  @RequestParam String messageId) {
        rabbitTemplate.convertAndSend("order.pay.queue", orderId.toString(),
                message -> {
                    message.getMessageProperties().setMessageId(messageId);
                    return message;
                });
        return "消息已发送: messageId=" + messageId;
    }

    /**
     * 模拟发送积分消息
     */
    @PostMapping("/send-points")
    public String sendPointsMessage(@RequestParam Long userId,
                                     @RequestParam int points,
                                     @RequestParam String messageId) {
        String content = userId + ":" + points;
        rabbitTemplate.convertAndSend("points.queue", content,
                message -> {
                    message.getMessageProperties().setMessageId(messageId);
                    return message;
                });
        return "积分消息已发送: messageId=" + messageId;
    }
}

六、进阶优化

6.1 Redis 数据过期后的兜底

Redis 的 key 设置了24小时过期,24小时后的重复消息 Redis 无法拦截,此时依赖数据库唯一索引兜底。

注意:数据库消费日志会无限增长,需要定期归档:

-- 每月清理3个月前的消费日志
DELETE FROM message_consume_log 
WHERE consume_time < DATE_SUB(NOW(), INTERVAL 3 MONTH);

6.2 业务唯一键替代消息ID

有些场景消息没有唯一ID,可以用业务唯一键替代:

// 用"订单ID + 操作类型"作为幂等key
String idempotentKey = "order:" + orderId + ":pay";

推荐:优先使用业务唯一键,因为它比消息ID更稳定(消息ID可能因为重发而变化)。

6.3 布隆过滤器优化

如果消息量特别大(>1亿),Redis 存储所有消息ID会占用大量内存。可以使用布隆过滤器:

优点:内存占用小(1亿条数据约100MB)
缺点:有假阳性(小概率把未处理的消息判为已处理)

适用:对假阳性容忍度高的场景(如日志去重)。

6.4 分布式锁 + 状态检查

对于并发消费场景,Redis SET NX 可能因为并发时序问题失效。此时需要:

1. 获取分布式锁(按业务ID)
2. 检查业务状态
3. 执行业务
4. 释放锁

其中 分布式锁 的选型与实现,在 Java 高并发面试中也经常被问到,建议结合具体业务深入理解。

6.5 幂等失败的处理

如果第一层 Redis 通过了但业务执行失败,Redis 已经被"占用",消息重试时会被误判为重复。

解决方案:Redis 标记为"处理中",成功后再标记为"已完成":

// 状态1:处理中
redisTemplate.opsForValue().set(key, "PROCESSING", 1, TimeUnit.MINUTES);

try {
    processBusiness();
    // 状态2:已完成(长时间保留)
    redisTemplate.opsForValue().set(key, "COMPLETED", 24, TimeUnit.HOURS);
} catch (Exception e) {
    // 处理失败,删除"处理中"标记,允许重试
    redisTemplate.delete(key);
    throw e;
}

七、踩坑指南

坑 表现 解决方案
幂等Key设置错误 不同消息判为重复 用消息ID或"业务ID+操作类型"
Redis与DB不一致 Redis 过期后重复消息无法拦截 Redis TTL + DB唯一索引双保险
消息ID变化 同一条消息重发ID不同 用业务唯一键替代消息ID
处理失败后重复拦截 处理失败的消息被误判重复 状态标记:处理中/已完成
数据库日志无限增长 磁盘空间被占满 定期归档历史日志
并发消费穿透 两个消费者同时处理 分布式锁
事务与幂等冲突 幂等成功后事务回滚 幂等记录与业务在同一事务

八、生产环境推荐

场景 推荐方案
普通业务 Redis SET NX + 业务状态检查
强一致场景 Redis + DB唯一索引 + 状态检查
高并发场景 Redis SET NX + 分布式锁
海量消息 布隆过滤器 + DB兜底
金融级场景 状态机 + 条件更新 + 唯一索引

参数推荐:

参数 推荐值 说明
Redis TTL 24小时 覆盖消息重试窗口
消费日志保留 3个月 平衡存储和审计
状态标记TTL 1分钟(处理中)/ 24小时(已完成) 区分处理中与完成
最大重试次数 3次 超过则进入死信队列

九、总结与思考

核心三句话:

  1. 消息重复不可避免——MQ 只保证 At Least Once,必须靠幂等消费。
  2. 幂等的核心是用唯一ID去重:Redis SET NX(快)+ 数据库唯一索引(稳)。
  3. 最强的幂等是业务天然幂等——通过状态检查和条件更新实现。

方案速查:

你的需求 推荐方案
快速幂等 Redis SET NX
强一致幂等 Redis + DB唯一索引
天然幂等 状态机 + 条件更新
海量消息 布隆过滤器

附录:完整Maven依赖

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-amqp</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-data-redis</artifactId>
    </dependency>
    <dependency>
        <groupId>org.apache.commons</groupId>
        <artifactId>commons-pool2</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-jdbc</artifactId>
    </dependency>
    <dependency>
        <groupId>com.mysql</groupId>
        <artifactId>mysql-connector-j</artifactId>
        <scope>runtime</scope>
    </dependency>
    <dependency>
        <groupId>org.projectlombok</groupId>
        <artifactId>lombok</artifactId>
        <optional>true</optional>
    </dependency>
</dependencies>

⚠️ 运行前请确认:

  • □ 确保 RabbitMQ、Redis、MySQL 均已启动
  • □ 创建 orders、message_consume_log、user_points 三张表
  • □ 在RabbitMQ中创建 order.pay.queue 和 points.queue
  • □ 使用相同的 messageId 多次调用 POST /api/test/send-pay,验证重复消息被拦截
  • □ 使用相同的 messageId 多次调用 POST /api/test/send-points,验证积分不会重复发放
  • □ 观察日志,验证 Redis 拦截和数据库兜底是否生效



上一篇:自由概率论与最优传输首次统一,熵最优传输方法登上 Invent Math
下一篇:乐元素国庆加班3倍工资:消消乐躺赚,白银之城为何这么赶
您需要登录后才可以回帖 登录 | 立即注册

手机版|小黑屋|网站地图|云栈社区 ( 苏ICP备2022046150号-2 )

GMT+8, 2026-10-5 22:17 , Processed in 0.111294 second(s), 41 queries , Gzip On.

Powered by Discuz! X3.5

© 2025-2026 云栈社区.

快速回复 返回顶部 返回列表