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 -> 16s3. 带限流的重试
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();
}
}一致性排障
一致性排障看什么
当业务反馈"消息已经发了,但状态没对上",排查通常要同时看:
- 生产端发送记录
- Broker 投递情况
- 消费端幂等记录
- 死信和补偿记录
排查工具和方法
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. 正确的分类策略:
- 网络超时: 指数退避重试
- 服务不可用: 限次重试+熔断
- 参数错误: 直接转死信
- 业务错误: 记录日志,正常ACK3. 一致性排障为什么必须看端到端链路?
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幂等、重试和一致性排障是消息队列生产环境的核心能力:
- 幂等是底线: 所有消费者必须实现幂等控制
- 重试要分类: 不同错误类型采用不同重试策略
- 排查看全程: 端到端链路追踪才能定位问题
- 补偿有机制: 失败消息要有兜底处理方案
- 监控要完善: 及时发现和告警问题
只有做好这些,才能保证MQ的可靠性和业务的一致性。
版本差异(旧版 → 当前)
| 组件 | 旧版(本文编写时) | 当前 |
|---|---|---|
| RabbitMQ | 3.8/3.9 | 3.13/4.x(quorum queue 为默认推荐) |
| Kafka | 2.x/3.0 | 3.7+/4.x(KRaft 模式取代 ZooKeeper) |
| RocketMQ | 4.x | 5.x(gRPC 通信、简化运维) |
| Java 版本 | JDK 8 | JDK 17+(Kafka 3.7+ 客户端要求) |
本文讲解的消息可靠性设计(投递确认、死信、幂等、顺序)原理不变;注意各中间件版本升级后的配置差异,如 Kafka 的 KRaft 模式、RabbitMQ 的 quorum queue。