{T}

MQ 幂等重试与一致性排障

概念与背景

消息队列链路真正上生产后,最常见的不是"会不会发消息",而是:

  • 同一条消息会不会重复消费
  • 重试会不会把系统打崩
  • 数据最终为什么没有收敛

这类问题最终都要回到幂等、重试分类和一致性排障

为什么 MQ 场景必须关注这些问题

消息队列的核心特性是异步解耦,这带来了显著优势,但也引入了新的复杂性:

特性优势带来的问题
异步处理提升系统响应速度消息投递时序不确定
削峰填谷保护下游系统消息积压风险
解耦系统独立演进数据一致性难以保证
至少一次投递不丢消息重复消费问题

幂等性设计

为什么 MQ 场景必须幂等

只要系统接受至少一次语义(At-Least-Once),就必须接受重复消息。这不是偶发异常,而是系统默认现实。

消息重复的场景:

java
// 场景1: 生产者重试导致重复
public class OrderProducer {
    public void sendOrderMessage(Order order) {
        try {
            // 网络超时,消息可能已发送成功
            producer.send(orderMessage);
        } catch (TimeoutException e) {
            // 生产者重试,导致消息重复
            producer.send(orderMessage); // 可能发送两次
        }
    }
}

// 场景2: 消费者确认失败导致重复
public class OrderConsumer {
    @RabbitListener(queues = "order.queue")
    public void processOrder(Message message) {
        processBusiness(message);
        // 处理成功但ACK失败,Broker会重新投递
        // throw new RuntimeException("ACK失败"); // 模拟异常
    }
}

// 场景3: Broker重启导致消息重复投递
// Broker在消息投递后但未收到ACK时重启
// 重启后会重新投递该消息

幂等性实现方案

方案一: 数据库唯一约束

适用场景: 业务操作本身需要数据库持久化

java
@Service
public class OrderService {
    
    @Autowired
    private OrderMapper orderMapper;
    
    /**
     * 方案1: 数据库唯一约束实现幂等
     * 优点: 实现简单,利用数据库ACID特性
     * 缺点: 性能受限于数据库,需要业务表设计支持
     */
    @Transactional
    public void createOrder(OrderMessage orderMsg) {
        // 尝试插入订单,订单号唯一约束
        try {
            Order order = new Order();
            order.setOrderNo(orderMsg.getOrderNo()); // 业务单号
            order.setUserId(orderMsg.getUserId());
            order.setAmount(orderMsg.getAmount());
            order.setCreateTime(new Date());
            
            orderMapper.insert(order);
            
            // 其他业务操作...
            
        } catch (DuplicateKeyException e) {
            // 唯一约束冲突,说明订单已存在
            log.warn("订单已存在,幂等拦截: orderNo={}", orderMsg.getOrderNo());
            // 直接返回,不抛异常,让消息正常ACK
        }
    }
}

-- 数据库表设计
CREATE TABLE `t_order` (
  `id` bigint(20) NOT NULL AUTO_INCREMENT,
  `order_no` varchar(64) NOT NULL COMMENT '订单号(唯一约束)',
  `user_id` bigint(20) NOT NULL,
  `amount` decimal(10,2) NOT NULL,
  `status` tinyint(4) NOT NULL DEFAULT '0',
  `create_time` datetime NOT NULL,
  PRIMARY KEY (`id`),
  UNIQUE KEY `uk_order_no` (`order_no`) -- 唯一约束
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

方案二: 独立幂等表

适用场景: 业务操作复杂,不适合在业务表加约束

java
@Service
public class PaymentService {
    
    @Autowired
    private IdempotentMapper idempotentMapper;
    @Autowired
    private PaymentMapper paymentMapper;
    
    /**
     * 方案2: 独立幂等表实现幂等
     * 优点: 不侵入业务表设计,灵活性高
     * 缺点: 需要额外的存储和维护
     */
    @Transactional
    public void processPayment(PaymentMessage paymentMsg) {
        String messageId = paymentMsg.getMessageId();
        
        // 1. 先插入幂等记录(利用唯一约束)
        try {
            IdempotentRecord record = new IdempotentRecord();
            record.setMessageId(messageId);
            record.setBizType("PAYMENT");
            record.setCreateTime(new Date());
            idempotentMapper.insert(record);
        } catch (DuplicateKeyException e) {
            log.warn("消息已处理,幂等拦截: messageId={}", messageId);
            return; // 幂等返回
        }
        
        // 2. 执行业务逻辑
        Payment payment = new Payment();
        payment.setPaymentNo(paymentMsg.getPaymentNo());
        payment.setAmount(paymentMsg.getAmount());
        paymentMapper.insert(payment);
        
        // 3. 其他业务操作...
    }
}

-- 幂等表设计
CREATE TABLE `t_idempotent` (
  `id` bigint(20) NOT NULL AUTO_INCREMENT,
  `message_id` varchar(128) NOT NULL COMMENT '消息ID',
  `biz_type` varchar(64) NOT NULL COMMENT '业务类型',
  `biz_key` varchar(128) DEFAULT NULL COMMENT '业务键',
  `status` tinyint(4) NOT NULL DEFAULT '1' COMMENT '状态:1-处理中,2-成功,3-失败',
  `create_time` datetime NOT NULL,
  `update_time` datetime DEFAULT NULL,
  PRIMARY KEY (`id`),
  UNIQUE KEY `uk_message_id` (`message_id`),
  KEY `idx_biz_type_key` (`biz_type`, `biz_key`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

方案三: Redis SETNX

适用场景: 对性能要求高,可接受极小概率的幂等失效

java
@Service
public class CouponService {
    
    @Autowired
    private StringRedisTemplate redisTemplate;
    @Autowired
    private CouponMapper couponMapper;
    
    /**
     * 方案3: Redis SETNX实现幂等
     * 优点: 性能高,不依赖数据库
     * 缺点: Redis故障可能导致幂等失效,需要配合过期时间
     */
    public void issueCoupon(CouponMessage couponMsg) {
        String messageId = couponMsg.getMessageId();
        String key = "mq:idempotent:coupon:" + messageId;
        
        // 1. SETNX设置幂等键
        Boolean success = redisTemplate.opsForValue()
            .setIfAbsent(key, "1", 7, TimeUnit.DAYS); // 7天过期
        
        if (!success) {
            log.warn("消息已处理,幂等拦截: messageId={}", messageId);
            return; // 幂等返回
        }
        
        try {
            // 2. 执行业务逻辑
            Coupon coupon = new Coupon();
            coupon.setUserId(couponMsg.getUserId());
            coupon.setAmount(couponMsg.getAmount());
            couponMapper.insert(coupon);
            
            // 3. 业务成功,保留Redis键
        } catch (Exception e) {
            // 4. 业务失败,删除幂等键,允许重试
            redisTemplate.delete(key);
            throw e;
        }
    }
}

方案四: 状态机幂等

适用场景: 业务有明确的状态流转

java
@Service
public class DeliveryService {
    
    @Autowired
    private DeliveryMapper deliveryMapper;
    
    /**
     * 方案4: 状态机幂等
     * 优点: 符合业务语义,天然幂等
     * 缺点: 需要设计合理的状态流转
     */
    @Transactional
    public void updateDeliveryStatus(DeliveryMessage deliveryMsg) {
        String deliveryNo = deliveryMsg.getDeliveryNo();
        String targetStatus = deliveryMsg.getTargetStatus();
        
        // 乐观锁更新,只有当前状态符合条件才能更新
        int rows = deliveryMapper.updateStatus(
            deliveryNo, 
            "SHIPPED",  // 当前状态必须是已发货
            targetStatus // 目标状态
        );
        
        if (rows == 0) {
            // 更新失败,可能已更新或状态不符合
            Delivery delivery = deliveryMapper.selectByDeliveryNo(deliveryNo);
            if (targetStatus.equals(delivery.getStatus())) {
                log.info("物流状态已更新,幂等成功: deliveryNo={}", deliveryNo);
                return; // 幂等返回
            }
            throw new BusinessException("物流状态不符合更新条件");
        }
        
        log.info("物流状态更新成功: deliveryNo={}, status={}", 
            deliveryNo, targetStatus);
    }
}

-- 状态机设计示例
-- 待发货 -> 已发货 -> 运输中 -> 已签收
-- 每个状态只能流转到下一个状态
UPDATE t_delivery 
SET status = 'TRANSPORTING', 
    update_time = NOW()
WHERE delivery_no = ? 
  AND status = 'SHIPPED'; -- 只有已发货才能改为运输中

幂等键设计原则

java
/**
 * 幂等键设计原则
 * 
 * 1. 唯一性: 幂等键必须全局唯一
 * 2. 稳定性: 幂等键不应该变化
 * 3. 业务相关性: 最好使用业务单号
 */
public class IdempotentKeyDesign {
    
    // 好的设计: 使用业务单号
    public String buildOrderKey(OrderMessage msg) {
        return msg.getOrderNo(); // 订单号天然唯一
    }
    
    // 好的设计: 组合业务字段
    public String buildPaymentKey(PaymentMessage msg) {
        return msg.getOrderId() + ":" + msg.getPaymentType();
    }
    
    // 好的设计: 使用消息ID
    public String buildMessageKey(Message message) {
        return message.getMessageProperties().getMessageId();
    }
    
    // 不好的设计: 使用时间戳
    public String badDesign() {
        return String.valueOf(System.currentTimeMillis()); // 每次都不同
    }
}

重试机制设计

重试怎么分层

更合理的重试思路通常是:

  • 临时错误: 限次重试
  • 不可恢复错误: 直接转死信
  • 长时间未收敛: 走补偿和人工处理

重试策略详解

1. 固定间隔重试

java
/**
 * 固定间隔重试策略
 * 适用场景: 临时性故障,快速恢复
 */
@Configuration
public class FixedRetryConfig {
    
    @Bean
    public SimpleRabbitListenerContainerFactory rabbitFactory(
            ConnectionFactory connectionFactory) {
        SimpleRabbitListenerContainerFactory factory = 
            new SimpleRabbitListenerContainerFactory();
        factory.setConnectionFactory(connectionFactory);
        
        // 重试配置
        factory.setDefaultRequeueRejected(false); // 不重新入队
        factory.setAdviceChain(RetryInterceptorBuilder
            .stateless()
            .maxAttempts(3) // 最大重试3次
            .backOffOptions(1000, 1.0, 1000) // 固定1秒间隔
            .recoverer(new RejectAndDontRequeueRecoverer()) // 转死信
            .build());
        
        return factory;
    }
}

2. 指数退避重试

java
/**
 * 指数退避重试策略
 * 适用场景: 下游服务压力较大,避免重试风暴
 */
@Configuration
public class ExponentialBackoffConfig {
    
    @Bean
    public SimpleRabbitListenerContainerFactory rabbitFactory(
            ConnectionFactory connectionFactory) {
        SimpleRabbitListenerContainerFactory factory = 
            new SimpleRabbitListenerContainerFactory();
        factory.setConnectionFactory(connectionFactory);
        
        factory.setAdviceChain(RetryInterceptorBuilder
            .stateless()
            .maxAttempts(5) // 最大重试5次
            .backOffOptions(
                1000,   // 初始间隔1秒
                2.0,    // 乘数因子2.0
                30000   // 最大间隔30秒
            )
            .recoverer(new RejectAndDontRequeueRecoverer())
            .build());
        
        return factory;
    }
}

// 重试间隔: 1s -> 2s -> 4s -> 8s -> 16s

3. 带限流的重试

java
/**
 * 带限流的重试策略
 * 适用场景: 防止重试压垮下游系统
 */
@Service
public class RetryWithRateLimit {
    
    @Autowired
    private RateLimiter rateLimiter; // 使用Guava RateLimiter
    
    @Retryable(
        value = {ServiceUnavailableException.class},
        maxAttempts = 3,
        backoff = @Backoff(delay = 1000, multiplier = 2)
    )
    public void processWithRetry(Message message) {
        // 先获取限流许可
        if (!rateLimiter.tryAcquire(1, TimeUnit.SECONDS)) {
            throw new RateLimitException("重试限流");
        }
        
        // 执行业务逻辑
        doProcess(message);
    }
    
    @Recover
    public void recover(ServiceUnavailableException e, Message message) {
        // 重试失败后的恢复逻辑
        log.error("重试失败,转入死信: messageId={}", 
            message.getMessageProperties().getMessageId());
        // 转入死信队列或补偿队列
    }
}

错误分类处理

java
/**
 * 错误分类处理策略
 */
@Service
public class ErrorClassificationConsumer {
    
    @RabbitListener(queues = "order.queue")
    public void processOrder(Message message) {
        try {
            doProcess(message);
            
        } catch (ValidationException e) {
            // 参数校验错误: 不重试,直接转死信
            log.error("参数校验失败,转死信: {}", e.getMessage());
            throw new AmqpRejectAndDontRequeueException(e);
            
        } catch (BusinessException e) {
            // 业务异常: 不重试,记录日志
            log.error("业务处理失败: {}", e.getMessage());
            // 正常ACK,不重试
            return;
            
        } catch (TimeoutException e) {
            // 超时异常: 重试
            log.warn("处理超时,等待重试: {}", e.getMessage());
            throw e; // 抛异常触发重试
            
        } catch (Exception e) {
            // 未知异常: 重试几次后转死信
            log.error("处理异常: {}", e.getMessage(), e);
            throw new AmqpRejectAndDontRequeueException(e);
        }
    }
    
    private void doProcess(Message message) {
        // 业务处理逻辑
    }
}

死信队列配置

java
/**
 * 死信队列配置
 */
@Configuration
public class DeadLetterConfig {
    
    // 死信交换机
    @Bean
    public DirectExchange deadLetterExchange() {
        return new DirectExchange("dlx.exchange");
    }
    
    // 死信队列
    @Bean
    public Queue deadLetterQueue() {
        return QueueBuilder
            .durable("dlx.order.queue")
            .build();
    }
    
    // 绑定
    @Bean
    public Binding dlxBinding() {
        return BindingBuilder
            .bind(deadLetterQueue())
            .to(deadLetterExchange())
            .with("order.dead");
    }
    
    // 业务队列(带死信配置)
    @Bean
    public Queue orderQueue() {
        return QueueBuilder
            .durable("order.queue")
            .deadLetterExchange("dlx.exchange") // 死信交换机
            .deadLetterRoutingKey("order.dead") // 死信路由键
            .ttl(60000) // 消息过期时间60秒
            .build();
    }
}

一致性排障

一致性排障看什么

当业务反馈"消息已经发了,但状态没对上",排查通常要同时看:

  1. 生产端发送记录
  2. Broker 投递情况
  3. 消费端幂等记录
  4. 死信和补偿记录

排查工具和方法

1. 生产端排查

java
/**
 * 生产端排查要点
 */
@Service
public class ProducerTroubleshoot {
    
    @Autowired
    private RabbitTemplate rabbitTemplate;
    
    /**
     * 开启发送确认
     */
    @PostConstruct
    public void init() {
        // 确认回调
        rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> {
            if (ack) {
                log.info("消息发送成功: correlationId={}", 
                    correlationData.getId());
            } else {
                log.error("消息发送失败: correlationId={}, cause={}", 
                    correlationData.getId(), cause);
            }
        });
        
        // 退回回调
        rabbitTemplate.setReturnsCallback(returned -> {
            log.error("消息被退回: exchange={}, routingKey={}, message={}", 
                returned.getExchange(),
                returned.getRoutingKey(),
                new String(returned.getMessage().getBody()));
        });
    }
    
    /**
     * 发送消息并记录日志
     */
    public void sendWithLog(Order order) {
        String messageId = UUID.randomUUID().toString();
        
        // 发送前记录
        log.info("准备发送消息: messageId={}, orderNo={}", 
            messageId, order.getOrderNo());
        
        CorrelationData correlationData = new CorrelationData(messageId);
        correlationData.setReturnedMessage(
            new Message(
                JSON.toJSONString(order).getBytes(),
                new MessageProperties()
            )
        );
        
        rabbitTemplate.convertAndSend(
            "order.exchange",
            "order.create",
            order,
            correlationData
        );
        
        // 发送后记录
        log.info("消息已发送: messageId={}", messageId);
    }
}

// 排查SQL
-- 查询发送日志
SELECT * FROM mq_send_log 
WHERE order_no = 'ORDER123' 
ORDER BY create_time DESC;

2. Broker 排查

bash
# RabbitMQ 管理命令

# 查看队列状态
rabbitmqctl list_queues name messages consumers

# 查看队列详情
rabbitmqctl list_queues name messages_ready messages_unacknowledged

# 查看连接
rabbitmqctl list_connections

# 查看消费者
rabbitmqctl list_consumers

# 查看交换机
rabbitmqctl list_exchanges

# 查看绑定关系
rabbitmqctl list_bindings

# Kafka 命令

# 查看消费者组
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --list

# 查看消费者组详情
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --describe --group order-group

# 查看主题详情
kafka-topics.sh --bootstrap-server localhost:9092 \
  --describe --topic order-topic

# 查看消息积压
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --describe --group order-group | awk '{print $1, $2, $3, $4, $5}'

3. 消费端排查

java
/**
 * 消费端排查要点
 */
@Service
public class ConsumerTroubleshoot {
    
    /**
     * 详细日志记录
     */
    @RabbitListener(queues = "order.queue")
    public void processOrder(Message message) {
        String messageId = message.getMessageProperties().getMessageId();
        
        // 1. 接收消息日志
        log.info("接收消息: messageId={}, timestamp={}", 
            messageId, System.currentTimeMillis());
        
        try {
            // 2. 解析消息
            Order order = JSON.parseObject(
                new String(message.getBody()), 
                Order.class
            );
            log.info("解析消息: messageId={}, orderNo={}", 
                messageId, order.getOrderNo());
            
            // 3. 幂等检查
            if (isProcessed(messageId)) {
                log.warn("消息已处理,幂等拦截: messageId={}", messageId);
                return; // 幂等返回
            }
            
            // 4. 业务处理
            doProcess(order);
            
            // 5. 记录处理成功
            log.info("消息处理成功: messageId={}, orderNo={}", 
                messageId, order.getOrderNo());
            
        } catch (Exception e) {
            // 6. 记录处理失败
            log.error("消息处理失败: messageId={}, error={}", 
                messageId, e.getMessage(), e);
            throw e;
        }
    }
    
    private boolean isProcessed(String messageId) {
        // 查询幂等记录
        return idempotentMapper.selectByMessageId(messageId) != null;
    }
}

-- 查询消费日志
SELECT * FROM mq_consume_log 
WHERE message_id = 'MSG123' 
ORDER BY create_time DESC;

-- 查询幂等记录
SELECT * FROM t_idempotent 
WHERE message_id = 'MSG123';

-- 查询业务数据
SELECT * FROM t_order 
WHERE order_no = 'ORDER123';

4. 死信队列排查

java
/**
 * 死信队列排查
 */
@Service
public class DeadLetterConsumer {
    
    /**
     * 消费死信消息
     */
    @RabbitListener(queues = "dlx.order.queue")
    public void processDeadLetter(Message message) {
        String messageId = message.getMessageProperties().getMessageId();
        
        // 记录死信信息
        log.error("接收到死信消息: messageId={}, reason={}", 
            messageId,
            message.getMessageProperties().getReason()); // 死信原因
        
        // 查看原始信息
        String exchange = message.getMessageProperties()
            .getOriginalExchange();
        String routingKey = message.getMessageProperties()
            .getOriginalRoutingKey();
        
        log.error("原始交换机: {}, 路由键: {}", exchange, routingKey);
        
        // 根据死信原因分类处理
        String reason = message.getMessageProperties().getReason();
        if ("expired".equals(reason)) {
            // 消息过期
            handleExpiredMessage(message);
        } else if ("rejected".equals(reason)) {
            // 消息被拒绝
            handleRejectedMessage(message);
        } else if ("maxlen".equals(reason)) {
            // 队列达到最大长度
            handleMaxLenMessage(message);
        }
    }
}

-- 查询死信记录
SELECT * FROM mq_dead_letter_log 
WHERE message_id = 'MSG123' 
ORDER BY create_time DESC;

端到端链路追踪

java
/**
 * 端到端链路追踪
 */
@Configuration
public class TraceConfig {
    
    @Bean
    public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
        RabbitTemplate template = new RabbitTemplate(connectionFactory);
        
        // 发送前拦截器
        template.setBeforePublishPostProcessors(message -> {
            // 设置追踪ID
            String traceId = MDC.get("traceId");
            if (traceId == null) {
                traceId = UUID.randomUUID().toString();
                MDC.put("traceId", traceId);
            }
            message.getMessageProperties().setHeader("traceId", traceId);
            message.getMessageProperties().setMessageId(traceId);
            
            return message;
        });
        
        return template;
    }
}

@Service
public class TraceConsumer {
    
    @RabbitListener(queues = "order.queue")
    public void process(Message message) {
        // 获取追踪ID
        String traceId = message.getMessageProperties()
            .getHeader("traceId");
        
        MDC.put("traceId", traceId);
        
        try {
            // 所有日志自动带上traceId
            log.info("开始处理消息");
            doProcess(message);
            log.info("消息处理完成");
        } finally {
            MDC.remove("traceId");
        }
    }
}

实战场景

场景一:支付成功但积分没到账

问题现象: 用户支付成功,但积分账户没有增加积分

排查步骤:

java
/**
 * 支付成功但积分没到账 - 排查示例
 */
@Service
public class PointsIssueTroubleshoot {
    
    /**
     * 排查步骤1: 确认消息是否发送
     */
    public void checkSendLog(String orderNo) {
        // 查询发送日志
        List<MqSendLog> sendLogs = mqSendLogMapper.selectByOrderNo(orderNo);
        
        if (sendLogs.isEmpty()) {
            log.error("未找到发送记录: orderNo={}", orderNo);
            // 可能原因: 支付回调没有发送消息
            return;
        }
        
        for (MqSendLog log : sendLogs) {
            log.info("发送记录: messageId={}, status={}, time={}", 
                log.getMessageId(), 
                log.getStatus(), 
                log.getCreateTime());
        }
    }
    
    /**
     * 排查步骤2: 确认Broker是否收到
     */
    public void checkBrokerStatus(String messageId) {
        // RabbitMQ: 通过管理API查询
        // curl -u guest:guest http://localhost:15672/api/queues/...
        
        // Kafka: 查询消费组lag
        // kafka-consumer-groups.sh --describe --group points-group
    }
    
    /**
     * 排查步骤3: 确认消费端是否处理
     */
    public void checkConsumeLog(String messageId) {
        // 查询消费日志
        MqConsumeLog consumeLog = mqConsumeLogMapper.selectByMessageId(messageId);
        
        if (consumeLog == null) {
            log.error("未找到消费记录: messageId={}", messageId);
            // 可能原因: 消费者未启动或队列堵塞
            return;
        }
        
        log.info("消费记录: status={}, error={}, time={}", 
            consumeLog.getStatus(),
            consumeLog.getErrorMessage(),
            consumeLog.getCreateTime());
    }
    
    /**
     * 排查步骤4: 确认是否被幂等拦截
     */
    public void checkIdempotent(String messageId) {
        IdempotentRecord record = idempotentMapper.selectByMessageId(messageId);
        
        if (record != null) {
            log.info("幂等记录: status={}", record.getStatus());
            // 如果status=成功,说明业务已处理
            // 检查积分是否真的没到账,还是查询缓存问题
        }
    }
    
    /**
     * 排查步骤5: 确认业务数据
     */
    public void checkPointsAccount(String userId, String orderNo) {
        // 查询积分流水
        PointsRecord record = pointsMapper.selectByOrderNo(userId, orderNo);
        
        if (record == null) {
            log.error("未找到积分记录: userId={}, orderNo={}", userId, orderNo);
            // 业务处理失败,查看是否有异常记录
        }
    }
    
    /**
     * 排查步骤6: 检查死信队列
     */
    public void checkDeadLetter(String messageId) {
        DeadLetterLog dlxLog = deadLetterMapper.selectByMessageId(messageId);
        
        if (dlxLog != null) {
            log.error("消息在死信队列: reason={}, time={}", 
                dlxLog.getReason(), 
                dlxLog.getCreateTime());
            // 根据死信原因处理
        }
    }
}

常见原因和解决方案:

原因排查方法解决方案
消息未发送查发送日志修复发送逻辑
消息路由失败查Broker退回日志检查路由配置
消费者异常查消费日志修复消费逻辑
幂等拦截查幂等表确认是否重复
业务异常查业务日志修复业务逻辑
死信队列查死信队列人工处理或补偿

场景二:重复发券

问题现象: 用户收到多张优惠券

原因分析:

java
/**
 * 重复发券问题分析
 */
@Service
public class DuplicateCouponIssue {
    
    // 错误示例: 没有幂等控制
    @RabbitListener(queues = "coupon.issue.queue")
    public void issueCouponBad(Message message) {
        CouponMessage couponMsg = JSON.parseObject(
            new String(message.getBody()),
            CouponMessage.class
        );
        
        // 直接发券,没有幂等检查
        Coupon coupon = new Coupon();
        coupon.setUserId(couponMsg.getUserId());
        coupon.setAmount(couponMsg.getAmount());
        couponMapper.insert(coupon); // 可能插入多次
        
        log.info("发券成功: userId={}, amount={}", 
            couponMsg.getUserId(), couponMsg.getAmount());
    }
    
    // 正确示例: 幂等控制
    @RabbitListener(queues = "coupon.issue.queue")
    public void issueCouponGood(Message message) {
        CouponMessage couponMsg = JSON.parseObject(
            new String(message.getBody()),
            CouponMessage.class
        );
        
        String messageId = message.getMessageProperties().getMessageId();
        
        // 1. 幂等检查
        if (idempotentMapper.exists(messageId)) {
            log.warn("优惠券已发放,幂等拦截: messageId={}", messageId);
            return;
        }
        
        // 2. 加锁(防止并发重复)
        String lockKey = "coupon:issue:" + couponMsg.getActivityId() + 
                        ":" + couponMsg.getUserId();
        Boolean locked = redisTemplate.opsForValue()
            .setIfAbsent(lockKey, "1", 10, TimeUnit.SECONDS);
        
        if (!locked) {
            log.warn("发券处理中,稍后重试: userId={}", couponMsg.getUserId());
            throw new RuntimeException("处理中,请重试");
        }
        
        try {
            // 3. 插入幂等记录
            idempotentMapper.insert(messageId, "COUPON");
            
            // 4. 发放优惠券(数据库唯一约束兜底)
            try {
                Coupon coupon = new Coupon();
                coupon.setCouponNo(buildCouponNo(couponMsg.getActivityId(), 
                                                 couponMsg.getUserId()));
                coupon.setUserId(couponMsg.getUserId());
                coupon.setAmount(couponMsg.getAmount());
                couponMapper.insert(coupon);
            } catch (DuplicateKeyException e) {
                log.info("优惠券已发放: userId={}", couponMsg.getUserId());
            }
            
        } finally {
            // 5. 释放锁
            redisTemplate.delete(lockKey);
        }
    }
    
    private String buildCouponNo(Long activityId, Long userId) {
        return activityId + ":" + userId; // 同一活动同一用户只能有一张券
    }
}

场景三:重试风暴

问题现象: 下游服务抖动,导致整个系统雪崩

问题分析:

java
/**
 * 重试风暴问题分析
 */
@Service
public class RetryStormExample {
    
    // 错误示例: 无限重试
    @RabbitListener(queues = "order.queue")
    public void processOrderBad(Message message) {
        Order order = parseOrder(message);
        
        // 调用下游服务
        inventoryService.deductStock(order); // 可能失败
        
        // 没有重试限制,会一直重试
        // 如果下游服务不可用,重试会压垮系统
    }
    
    // 正确示例: 重试控制
    @RabbitListener(queues = "order.queue")
    public void processOrderGood(Message message) {
        String messageId = message.getMessageProperties().getMessageId();
        
        // 1. 检查重试次数
        Integer retryCount = message.getMessageProperties()
            .getHeader("retryCount");
        
        if (retryCount != null && retryCount > 3) {
            log.error("重试次数超限,转死信: messageId={}, retryCount={}", 
                messageId, retryCount);
            throw new AmqpRejectAndDontRequeueException("重试次数超限");
        }
        
        try {
            Order order = parseOrder(message);
            
            // 2. 调用下游服务(带熔断)
            inventoryService.deductStockWithCircuitBreaker(order);
            
        } catch (ServiceUnavailableException e) {
            // 3. 服务不可用,指数退避重试
            log.warn("下游服务不可用,等待重试: {}", e.getMessage());
            
            // 设置重试次数
            message.getMessageProperties()
                .setHeader("retryCount", (retryCount == null ? 0 : retryCount) + 1);
            
            throw e; // 触发重试
        }
    }
}

重试风暴防护策略:

java
/**
 * 重试风暴防护
 */
@Configuration
public class RetryStormProtection {
    
    /**
     * 1. 熔断器保护
     */
    @Bean
    public CircuitBreaker circuitBreaker() {
        CircuitBreakerConfig config = CircuitBreakerConfig.custom()
            .failureRateThreshold(50) // 失败率50%触发熔断
            .waitDurationInOpenState(Duration.ofSeconds(30)) // 熔断30秒
            .ringBufferSizeInHalfOpenState(10) // 半开状态尝试10次
            .ringBufferSizeInClosedState(100) // 关闭状态记录100次
            .build();
        
        return CircuitBreaker.of("inventoryBreaker", config);
    }
    
    /**
     * 2. 限流保护
     */
    @Bean
    public RateLimiter rateLimiter() {
        return RateLimiter.create(100); // 每秒100个请求
    }
    
    /**
     * 3. 超时控制
     */
    @Bean
    public SimpleRabbitListenerContainerFactory rabbitFactory(
            ConnectionFactory connectionFactory) {
        SimpleRabbitListenerContainerFactory factory = 
            new SimpleRabbitListenerContainerFactory();
        factory.setConnectionFactory(connectionFactory);
        
        // 消费超时
        factory.setReceiveTimeout(5000L); // 5秒
        
        // 重试配置
        factory.setAdviceChain(RetryInterceptorBuilder
            .stateless()
            .maxAttempts(3) // 最多3次
            .backOffOptions(1000, 2.0, 10000) // 1s,2s,4s
            .build());
        
        return factory;
    }
}

排查与治理思路

排查顺序

code
问题报告: "消息发了但数据没对"
    ↓
1. 确认问题类型
    ├─ "消息没消费" → 查消费日志
    └─ "消费了但没收敛" → 查业务数据
    ↓
2. 检查幂等记录
    ├─ 有记录 → 业务是否正常
    └─ 无记录 → 消费是否成功
    ↓
3. 检查重试配置
    ├─ 重试次数 → 是否合理
    └─ 重试间隔 → 是否合理
    ↓
4. 检查死信和补偿
    ├─ 死信队列 → 是否有消息
    └─ 补偿记录 → 是否已补偿

治理重点

java
/**
 * MQ治理清单
 */
public class MQGovernanceChecklist {
    
    /**
     * 1. 幂等治理
     */
    public void checkIdempotent() {
        // √ 所有消费者都有幂等控制
        // √ 幂等键设计合理(业务单号或消息ID)
        // √ 幂等记录有过期清理机制
        // √ 幂等失败有明确的日志和告警
    }
    
    /**
     * 2. 重试治理
     */
    public void checkRetry() {
        // √ 重试次数有上限(建议3-5次)
        // √ 重试间隔使用指数退避
        // √ 临时错误和永久错误分类处理
        // √ 重试失败有死信兜底
    }
    
    /**
     * 3. 死信治理
     */
    public void checkDeadLetter() {
        // √ 所有队列都配置死信队列
        // √ 死信队列有消费者处理
        // √ 死信消息有分类和记录
        // √ 死信有补偿和人工处理入口
    }
    
    /**
     * 4. 监控治理
     */
    public void checkMonitoring() {
        // √ 消息积压告警
        // √ 消费失败率告警
        // √ 死信队列数量告警
        // √ 端到端延迟监控
    }
    
    /**
     * 5. 日志治理
     */
    public void checkLogging() {
        // √ 发送端有完整日志
        // √ 消费端有详细日志
        // √ 关键节点记录traceId
        // √ 异常堆栈完整记录
    }
}

常见误区

误区一: 只要有 MQ 就默认最终一致

java
/**
 * 误区: MQ能自动保证最终一致性
 * 
 * 实际: MQ只是传输工具,一致性需要业务保证
 */
@Service
public class ConsistencyMyth {
    
    // 错误理解
    public void wrongUnderstanding() {
        // × "用了MQ就一定能保证数据一致性"
        // × "消息发送成功就完事了"
        // × "消费者一定能处理成功"
    }
    
    // 正确理解
    public void correctUnderstanding() {
        // √ MQ只是消息传输通道
        // √ 需要幂等保证不重复处理
        // √ 需要补偿机制处理失败情况
        // √ 需要监控告警及时发现问题
    }
}

误区二: 只盯 Broker 不看消费者

java
/**
 * 误区: 只关注Broker状态,忽略消费端
 */
@Service
public class BrokerOnlyMyth {
    
    // 错误做法
    public void wrongApproach() {
        // × 只看Broker消息积压
        // × 只看队列消费者数量
        // × 不看消费端日志和异常
        // × 不看幂等记录
    }
    
    // 正确做法
    public void correctApproach() {
        // √ 同时查看Broker和消费端
        // √ 查看消费端处理日志
        // √ 查看幂等拦截记录
        // √ 查看业务数据变化
    }
}

误区三: 重试设计没有边界

java
/**
 * 误区: 无限重试或无重试策略
 */
@Service
public class RetryBoundaryMyth {
    
    // 错误示例: 无限重试
    @RabbitListener(queues = "order.queue")
    public void unlimitedRetry(Message message) {
        // × 没有重试次数限制
        // × 没有重试间隔控制
        // × 所有异常都重试
        
        doProcess(message); // 可能无限重试
    }
    
    // 错误示例: 参数错误也重试
    @RabbitListener(queues = "order.queue")
    public void retryOnValidationError(Message message) {
        try {
            Order order = parseOrder(message);
            
            if (order.getAmount() == null || order.getAmount() <= 0) {
                // × 参数错误也重试,永远不会成功
                throw new BusinessException("金额无效");
            }
            
            doProcess(order);
            
        } catch (BusinessException e) {
            // × 业务异常不应该重试
            throw e;
        }
    }
    
    // 正确示例: 分类处理
    @RabbitListener(queues = "order.queue")
    public void correctRetry(Message message) {
        try {
            Order order = parseOrder(message);
            
            // 参数校验失败: 不重试
            if (order.getAmount() == null || order.getAmount() <= 0) {
                log.error("参数校验失败: messageId={}", 
                    message.getMessageProperties().getMessageId());
                return; // 正常ACK
            }
            
            doProcess(order);
            
        } catch (ServiceUnavailableException e) {
            // 服务不可用: 重试
            log.warn("服务不可用,等待重试: {}", e.getMessage());
            throw e;
            
        } catch (Exception e) {
            // 其他异常: 重试几次后转死信
            log.error("处理异常: {}", e.getMessage(), e);
            throw new AmqpRejectAndDontRequeueException(e);
        }
    }
}

误区四: 没有死信和补偿入口

java
/**
 * 误区: 消息失败了就丢了
 */
@Service
public class NoDeadLetterMyth {
    
    // 错误做法
    @RabbitListener(queues = "order.queue")
    public void noDeadLetter(Message message) {
        try {
            doProcess(message);
        } catch (Exception e) {
            log.error("处理失败: {}", e.getMessage());
            // × 异常被吞掉,消息丢失
            // × 没有死信队列
            // × 没有补偿机制
        }
    }
    
    // 正确做法
    @RabbitListener(queues = "order.queue")
    public void withDeadLetter(Message message) {
        try {
            doProcess(message);
            
        } catch (BusinessException e) {
            // 业务异常: 记录日志,正常ACK
            log.error("业务处理失败: {}", e.getMessage());
            recordFailedMessage(message, e);
            return;
            
        } catch (Exception e) {
            // 其他异常: 转死信队列
            log.error("系统异常,转死信: {}", e.getMessage(), e);
            throw new AmqpRejectAndDontRequeueException(e);
        }
    }
    
    private void recordFailedMessage(Message message, Exception e) {
        // 记录失败消息到补偿表
        FailedMessage failed = new FailedMessage();
        failed.setMessageId(message.getMessageProperties().getMessageId());
        failed.setContent(new String(message.getBody()));
        failed.setErrorType(e.getClass().getName());
        failed.setErrorMessage(e.getMessage());
        failed.setCreateTime(new Date());
        
        failedMessageMapper.insert(failed);
    }
}

最佳实践

1. 幂等设计最佳实践

java
/**
 * 幂等设计最佳实践
 */
@Service
public class IdempotentBestPractices {
    
    /**
     * 完整的幂等控制流程
     */
    @Transactional
    public void processWithIdempotent(Message message) {
        String messageId = message.getMessageProperties().getMessageId();
        String businessKey = extractBusinessKey(message);
        
        // 1. 并发控制(分布式锁)
        String lockKey = "lock:" + businessKey;
        boolean locked = tryLock(lockKey, 10, TimeUnit.SECONDS);
        
        if (!locked) {
            throw new ConcurrentProcessingException("处理中,请稍后");
        }
        
        try {
            // 2. 幂等检查(Redis快速检查)
            if (isProcessedInRedis(messageId)) {
                log.info("消息已处理(Redis): messageId={}", messageId);
                return;
            }
            
            // 3. 幂等检查(数据库精确检查)
            if (isProcessedInDB(messageId)) {
                log.info("消息已处理(DB): messageId={}", messageId);
                return;
            }
            
            // 4. 插入幂等记录
            insertIdempotentRecord(messageId, businessKey);
            
            // 5. 执行业务逻辑
            doBusiness(message);
            
            // 6. 更新幂等状态为成功
            updateIdempotentStatus(messageId, "SUCCESS");
            
            // 7. 记录Redis缓存
            markProcessedInRedis(messageId);
            
        } catch (Exception e) {
            // 8. 更新幂等状态为失败
            updateIdempotentStatus(messageId, "FAILED");
            throw e;
        } finally {
            // 9. 释放锁
            unlock(lockKey);
        }
    }
    
    private boolean isProcessedInRedis(String messageId) {
        return Boolean.TRUE.equals(
            redisTemplate.hasKey("idempotent:" + messageId)
        );
    }
    
    private boolean isProcessedInDB(String messageId) {
        return idempotentMapper.selectByMessageId(messageId) != null;
    }
    
    private void markProcessedInRedis(String messageId) {
        redisTemplate.opsForValue().set(
            "idempotent:" + messageId,
            "1",
            7,
            TimeUnit.DAYS
        );
    }
}

2. 重试设计最佳实践

java
/**
 * 重试设计最佳实践
 */
@Service
public class RetryBestPractices {
    
    /**
     * 分类重试策略
     */
    @Retryable(
        value = {
            TimeoutException.class,
            ServiceUnavailableException.class
        },
        maxAttempts = 3,
        backoff = @Backoff(
            delay = 1000,
            multiplier = 2,
            maxDelay = 10000
        )
    )
    public void processWithRetry(Message message) {
        // 仅对临时性错误重试
        doProcess(message);
    }
    
    /**
     * 重试失败处理
     */
    @Recover
    public void recoverRetry(Exception e, Message message) {
        String messageId = message.getMessageProperties().getMessageId();
        
        log.error("重试失败: messageId={}, error={}", messageId, e.getMessage());
        
        // 记录失败消息
        FailedMessage failed = new FailedMessage();
        failed.setMessageId(messageId);
        failed.setErrorType(e.getClass().getName());
        failed.setErrorMessage(e.getMessage());
        failedMessageMapper.insert(failed);
        
        // 转入补偿队列
        compensationService.addToCompensationQueue(message);
    }
}

3. 监控告警最佳实践

java
/**
 * 监控告警最佳实践
 */
@Configuration
public class MonitoringBestPractices {
    
    /**
     * 消息积压监控
     */
    @Scheduled(fixedRate = 60000) // 每分钟检查
    public void checkQueueBacklog() {
        // RabbitMQ
        Properties props = rabbitAdmin.getQueueProperties("order.queue");
        Integer messageCount = (Integer) props.get("messageCount");
        
        if (messageCount > 10000) {
            // 发送告警
            alertService.sendAlert(
                "消息积压告警",
                String.format("队列order.queue积压%d条消息", messageCount)
            );
        }
        
        // Kafka
        // 使用kafka-consumer-groups命令检查lag
    }
    
    /**
     * 消费失败率监控
     */
    @Scheduled(fixedRate = 60000)
    public void checkConsumeFailRate() {
        // 查询最近1分钟的消费统计
        ConsumeStats stats = consumeLogMapper.selectStatsLastMinute();
        
        double failRate = (double) stats.getFailCount() / stats.getTotalCount();
        
        if (failRate > 0.1) { // 失败率超过10%
            alertService.sendAlert(
                "消费失败率告警",
                String.format("消费失败率%.2f%%,请检查", failRate * 100)
            );
        }
    }
    
    /**
     * 死信队列监控
     */
    @Scheduled(fixedRate = 60000)
    public void checkDeadLetterQueue() {
        Properties props = rabbitAdmin.getQueueProperties("dlx.order.queue");
        Integer messageCount = (Integer) props.get("messageCount");
        
        if (messageCount > 100) {
            alertService.sendAlert(
                "死信队列告警",
                String.format("死信队列有%d条消息待处理", messageCount)
            );
        }
    }
}

面试要点

基础问题

1. 为什么 MQ 场景下幂等是底线?

code
答:
1. 消息重复是必然的,不是偶发的
   - 生产者重试导致重复
   - 消费者ACK失败导致重复
   - Broker重启导致重复投递

2. 如果没有幂等,会导致:
   - 数据重复(重复发券、重复扣款)
   - 资金损失(重复发放奖励)
   - 业务逻辑错误(订单状态混乱)

3. 幂等是MQ可靠性的基础保障
   - 没有幂等,MQ就是不可靠的
   - 幂等是业务方的责任,不是MQ的责任

2. 重试为什么必须分类?

code
答:
1. 不同类型的错误,处理方式不同:
   - 临时错误(网络超时): 应该重试
   - 永久错误(参数错误): 不应该重试
   - 业务错误: 不应该重试

2. 不分类会导致:
   - 永久错误无限重试,浪费资源
   - 业务错误重试,导致错误放大
   - 重试风暴压垮系统

3. 正确的分类策略:
   - 网络超时: 指数退避重试
   - 服务不可用: 限次重试+熔断
   - 参数错误: 直接转死信
   - 业务错误: 记录日志,正常ACK

3. 一致性排障为什么必须看端到端链路?

code
答:
1. MQ只是消息传输的一环,不是全部:
   - 发送端 → Broker → 消费端
   - 任何一环都可能出问题

2. 只看某一环会遗漏问题:
   - 只看Broker: 不知道消费是否成功
   - 只看消费端: 不知道消息是否发送
   - 只看业务: 不知道是MQ问题还是业务问题

3. 端到端排查才能定位根因:
   - 发送日志: 确认消息是否发送
   - Broker状态: 确认是否投递
   - 消费日志: 确认是否处理
   - 幂等记录: 确认是否重复
   - 业务数据: 确认最终结果

进阶问题

4. 如何设计一个通用的幂等框架?

java
/**
 * 通用幂等框架设计
 */
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
public @interface Idempotent {
    // 幂等键表达式(SpEL)
    String key();
    
    // 幂等键前缀
    String prefix() default "idempotent";
    
    // 过期时间(秒)
    int expire() default 86400;
    
    // 幂等失败处理
    IdempotentStrategy strategy() default IdempotentStrategy.RETURN;
}

public enum IdempotentStrategy {
    RETURN,     // 直接返回
    THROW,      // 抛异常
    CUSTOM      // 自定义处理
}

@Aspect
@Component
public class IdempotentAspect {
    
    @Autowired
    private StringRedisTemplate redisTemplate;
    
    @Around("@annotation(idempotent)")
    public Object around(ProceedingJoinPoint point, Idempotent idempotent) 
            throws Throwable {
        // 1. 解析幂等键
        String key = parseKey(point, idempotent.key());
        String fullKey = idempotent.prefix() + ":" + key;
        
        // 2. SETNX检查
        Boolean success = redisTemplate.opsForValue()
            .setIfAbsent(fullKey, "1", idempotent.expire(), TimeUnit.SECONDS);
        
        if (!success) {
            // 3. 幂等失败处理
            return handleIdempotentFail(idempotent.strategy(), key);
        }
        
        try {
            // 4. 执行业务
            return point.proceed();
        } catch (Exception e) {
            // 5. 业务失败,删除幂等键
            redisTemplate.delete(fullKey);
            throw e;
        }
    }
}

// 使用示例
@Service
public class OrderService {
    
    @Idempotent(key = "#orderNo", prefix = "order", expire = 86400)
    public void createOrder(String orderNo, Order order) {
        // 自动幂等保护
        orderMapper.insert(order);
    }
}

5. 如何防止重试风暴?

code
答:
1. 限制重试次数
   - 设置最大重试次数(3-5次)
   - 超过次数转死信队列

2. 指数退避策略
   - 重试间隔递增(1s, 2s, 4s, 8s...)
   - 避免同时大量重试

3. 熔断器保护
   - 失败率达到阈值触发熔断
   - 熔断期间快速失败

4. 限流保护
   - 限制消费速率
   - 防止重试压垮下游

5. 错误分类处理
   - 临时错误才重试
   - 永久错误直接转死信

6. 监控告警
   - 重试次数监控
   - 失败率监控
   - 及时发现问题

6. 如何实现消息的端到端追踪?

java
/**
 * 端到端消息追踪方案
 */
@Configuration
public class EndToEndTracing {
    
    // 1. 发送端注入追踪信息
    @Bean
    public RabbitTemplate rabbitTemplate(ConnectionFactory factory) {
        RabbitTemplate template = new RabbitTemplate(factory);
        
        template.setBeforePublishPostProcessors(message -> {
            // 生成或获取TraceId
            String traceId = MDC.get("traceId");
            if (traceId == null) {
                traceId = UUID.randomUUID().toString();
            }
            
            // 注入追踪信息
            message.getMessageProperties().setHeader("traceId", traceId);
            message.getMessageProperties().setMessageId(traceId);
            message.getMessageProperties().setHeader("sendTime", 
                System.currentTimeMillis());
            
            // 记录发送日志
            logSendTrace(message, traceId);
            
            return message;
        });
        
        return template;
    }
    
    // 2. 消费端提取追踪信息
    @Component
    public class TracingConsumer {
        
        @RabbitListener(queues = "order.queue")
        public void process(Message message) {
            // 提取追踪信息
            String traceId = message.getMessageProperties()
                .getHeader("traceId");
            Long sendTime = message.getMessageProperties()
                .getHeader("sendTime");
            
            MDC.put("traceId", traceId);
            
            try {
                // 记录接收日志
                logReceiveTrace(message, traceId, sendTime);
                
                // 业务处理
                doProcess(message);
                
                // 记录成功日志
                logSuccessTrace(traceId);
                
            } catch (Exception e) {
                // 记录失败日志
                logFailTrace(traceId, e);
                throw e;
            } finally {
                MDC.remove("traceId");
            }
        }
    }
    
    // 3. 日志格式
    // 发送: [TRACE_SEND] traceId=xxx, exchange=xxx, routingKey=xxx, sendTime=xxx
    // 接收: [TRACE_RECEIVE] traceId=xxx, queue=xxx, receiveTime=xxx, latency=xxx
    // 成功: [TRACE_SUCCESS] traceId=xxx, processTime=xxx
    // 失败: [TRACE_FAIL] traceId=xxx, error=xxx
}

实战场景

7. 如何处理支付成功但积分没到账的问题?

code
答:
排查步骤:
1. 确认消息是否发送
   - 查询发送日志表
   - 检查发送状态

2. 确认Broker是否收到
   - RabbitMQ: 管理界面查看队列状态
   - Kafka: 查看消费者组lag

3. 确认消费端是否处理
   - 查询消费日志表
   - 检查处理状态

4. 确认是否被幂等拦截
   - 查询幂等记录表
   - 检查业务数据是否正常

5. 确认业务数据
   - 查询积分流水表
   - 检查积分余额

6. 检查死信队列
   - 查看是否有死信消息
   - 根据死信原因处理

解决方案:
- 消息未发送: 修复发送逻辑,补发消息
- 消息积压: 增加消费者,临时扩容
- 消费失败: 修复bug,手动重试
- 幂等拦截: 确认是否重复,如不是则修复
- 业务异常: 修复业务逻辑,补偿数据
- 死信消息: 分析原因,重新处理

8. 如何设计一个可靠的补偿机制?

java
/**
 * 可靠补偿机制设计
 */
@Service
public class CompensationService {
    
    @Autowired
    private CompensateTaskMapper compensateTaskMapper;
    
    /**
     * 1. 失败消息记录
     */
    public void recordFailedMessage(Message message, Exception e) {
        CompensateTask task = new CompensateTask();
        task.setMessageId(message.getMessageProperties().getMessageId());
        task.setContent(new String(message.getBody()));
        task.setErrorType(e.getClass().getName());
        task.setErrorMessage(e.getMessage());
        task.setStatus("PENDING"); // 待处理
        task.setRetryCount(0);
        task.setMaxRetry(3);
        task.setCreateTime(new Date());
        
        compensateTaskMapper.insert(task);
    }
    
    /**
     * 2. 补偿任务调度
     */
    @Scheduled(fixedRate = 60000) // 每分钟执行
    public void executeCompensateTasks() {
        // 查询待补偿的任务
        List<CompensateTask> tasks = compensateTaskMapper
            .selectPendingTasks(100); // 每次100条
        
        for (CompensateTask task : tasks) {
            try {
                // 执行补偿
                doCompensate(task);
                
                // 更新状态为成功
                task.setStatus("SUCCESS");
                task.setUpdateTime(new Date());
                compensateTaskMapper.updateById(task);
                
            } catch (Exception e) {
                // 更新重试次数
                task.setRetryCount(task.getRetryCount() + 1);
                
                if (task.getRetryCount() >= task.getMaxRetry()) {
                    // 超过最大重试次数,标记为需人工处理
                    task.setStatus("MANUAL");
                }
                
                task.setUpdateTime(new Date());
                compensateTaskMapper.updateById(task);
            }
        }
    }
    
    /**
     * 3. 人工处理入口
     */
    @GetMapping("/api/compensate/tasks")
    public List<CompensateTask> listManualTasks() {
        return compensateTaskMapper.selectByStatus("MANUAL");
    }
    
    @PostMapping("/api/compensate/tasks/{id}/retry")
    public void manualRetry(@PathVariable Long id) {
        CompensateTask task = compensateTaskMapper.selectById(id);
        
        try {
            doCompensate(task);
            task.setStatus("SUCCESS");
        } catch (Exception e) {
            task.setErrorMessage(e.getMessage());
        }
        
        task.setUpdateTime(new Date());
        compensateTaskMapper.updateById(task);
    }
    
    @PostMapping("/api/compensate/tasks/{id}/ignore")
    public void ignoreTask(@PathVariable Long id) {
        CompensateTask task = compensateTaskMapper.selectById(id);
        task.setStatus("IGNORED");
        task.setUpdateTime(new Date());
        compensateTaskMapper.updateById(task);
    }
}

总结:

MQ幂等、重试和一致性排障是消息队列生产环境的核心能力:

  1. 幂等是底线: 所有消费者必须实现幂等控制
  2. 重试要分类: 不同错误类型采用不同重试策略
  3. 排查看全程: 端到端链路追踪才能定位问题
  4. 补偿有机制: 失败消息要有兜底处理方案
  5. 监控要完善: 及时发现和告警问题

只有做好这些,才能保证MQ的可靠性和业务的一致性。

版本差异(旧版 → 当前)

组件旧版(本文编写时)当前
RabbitMQ3.8/3.93.13/4.x(quorum queue 为默认推荐)
Kafka2.x/3.03.7+/4.x(KRaft 模式取代 ZooKeeper)
RocketMQ4.x5.x(gRPC 通信、简化运维)
Java 版本JDK 8JDK 17+(Kafka 3.7+ 客户端要求)

本文讲解的消息可靠性设计(投递确认、死信、幂等、顺序)原理不变;注意各中间件版本升级后的配置差异,如 Kafka 的 KRaft 模式、RabbitMQ 的 quorum queue。