消费可靠性与消息积压治理
概念与背景
消息队列真正难的地方,往往不是"怎么发消息",而是消息发出去之后能不能稳定落库、能不能被正确消费、积压了之后怎么止血。很多线上事故都不是 Broker 直接挂掉,而是消费速度跟不上、重复消费没兜住、失败重试把下游打崩。
所以 MQ 进阶必须重点看三件事:
- 消费可靠性如何保证
- 消息失败后如何重试、回退、隔离
- 消息积压出现后如何判断和治理
消费可靠性的三大挑战
┌─────────────────────────────────────────────┐
│ 消费可靠性的三大挑战 │
├─────────────────────────────────────────────┤
│ 1. 消息丢失 │
│ - 消费者宕机,消息未确认 │
│ - 自动 ACK 后业务失败 │
│ - 异常被吞掉 │
├─────────────────────────────────────────────┤
│ 2. 消息重复 │
│ - 消费成功但 ACK 失败 │
│ - 消费者重启导致重新投递 │
│ - 生产者重复发送 │
├─────────────────────────────────────────────┤
│ 3. 消息积压 │
│ - 消费速度 < 生产速度 │
│ - 下游依赖故障 │
│ - 毒消息阻塞队列 │
└─────────────────────────────────────────────┘原理与机制
一条消息的可靠性链路
从生产到消费,一条消息通常要经过:
1. 生产者发送到 Broker
↓
2. Broker 持久化并复制
↓
3. 消费者拉取或接收消息
↓
4. 消费者执行业务逻辑
↓
5. 消费成功后提交确认(ACK)真正的风险点主要集中在后三步:
风险点1:消费者收到了消息,但业务还没执行完就宕机
消费者收到消息
↓
开始执行业务逻辑
↓
宕机!(业务未完成)
↓
消息状态: 未确认
↓
Broker 重新投递
↓
消息未丢失 √前提: 使用手动 ACK,而不是自动 ACK
风险点2:业务执行成功了,但确认没有提交成功
消费者收到消息
↓
执行业务逻辑 → 成功
↓
发送 ACK → 失败(网络问题)
↓
Broker 认为消息未消费
↓
重新投递
↓
重复消费 ×解决方案: 消费者必须幂等
风险点3:消费失败后无限重试,导致下游雪崩
消息处理失败
↓
重新入队
↓
立即再次消费
↓
再次失败
↓
重新入队
↓
循环... (下游被反复调用)
↓
下游雪崩 ×解决方案: 限制重试次数,失败后转死信队列
消费确认、重试与死信
ACK 机制详解
ACK(确认)的作用:
告诉 Broker 这条消息已经成功处理,可以从队列中删除。
三种确认方式:
| 确认方式 | 说明 | 可靠性 | 性能 | 适用场景 |
|---|---|---|---|---|
| 自动确认(Auto) | 消息到达消费者立即确认 | 低 | 高 | 非关键业务 |
| 手动确认(Manual) | 业务成功后再确认 | 高 | 中 | 关键业务 |
| 批量确认(Batch) | 批量确认多条消息 | 中 | 高 | 大批量处理 |
代码示例:
// × 自动确认(默认): 不可靠
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
// 消息到达立即 ACK
// 如果下面业务失败,消息已丢失
orderService.process(event);
}
// √ 手动确认: 可靠
@RabbitListener(queues = "order.queue", ackMode = "MANUAL")
public void handleOrder(OrderEvent event, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) {
try {
orderService.process(event);
// 业务成功,手动 ACK
channel.basicAck(tag, false);
} catch (Exception e) {
// 业务失败,NACK
channel.basicNack(tag, false, false);
}
}NACK 与 Reject
NACK(Negative Acknowledgment):
- 告诉 Broker 这条消息处理失败
- 可以选择是否重新入队
// basicNack 参数说明
channel.basicNack(
deliveryTag, // 消息标识
multiple, // 是否批量确认(false=单条)
requeue // 是否重新入队
);Reject:
- 拒绝单条消息
- 比 NACK 简单
// basicReject 参数说明
channel.basicReject(
deliveryTag, // 消息标识
requeue // 是否重新入队
);NACK vs Reject:
| 操作 | 批量 | 重新入队 | 适用场景 |
|---|---|---|---|
| basicNack | 支持 | 支持 | 批量失败处理 |
| basicReject | 不支持 | 支持 | 单条失败处理 |
重试机制
重试的三种策略:
┌─────────────────────────────────────────┐
│ 策略1: 立即重试(不推荐) │
│ - 消息重新入队 │
│ - 立即被再次消费 │
│ - 风险: 下游压力持续 │
├─────────────────────────────────────────┤
│ 策略2: 延迟重试(推荐) │
│ - 消息转延迟队列 │
│ - 延迟后再次投递 │
│ - 优点: 给下游恢复时间 │
├─────────────────────────────────────────┤
│ 策略3: 指数退避重试(推荐) │
│ - 重试间隔指数增长 │
│ - 1s → 2s → 4s → 8s → 16s │
│ - 优点: 避免重试风暴 │
└─────────────────────────────────────────┘延迟重试实现(RabbitMQ):
@Configuration
public class RetryConfig {
// 延迟队列
@Bean
public Queue retryQueue() {
return QueueBuilder.durable("order.retry.queue")
.withArgument("x-dead-letter-exchange", "order.exchange")
.withArgument("x-dead-letter-routing-key", "order.queue")
.withArgument("x-message-ttl", 10000) // 10秒后转主队列
.build();
}
}
@Service
public class OrderConsumer {
@RabbitListener(queues = "order.queue", ackMode = "MANUAL")
public void handleOrder(OrderEvent event, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag,
@Header("retry_count", required = false) Integer retryCount) {
retryCount = retryCount == null ? 0 : retryCount;
try {
orderService.process(event);
channel.basicAck(tag, false);
} catch (Exception e) {
log.error("处理失败,准备重试: orderId={}, retryCount={}",
event.getOrderId(), retryCount, e);
if (retryCount < 3) {
// 转延迟队列重试
rabbitTemplate.convertAndSend(
"order.exchange",
"order.retry",
event,
message -> {
message.getMessageProperties().setHeader("retry_count", retryCount + 1);
return message;
}
);
channel.basicAck(tag, false);
} else {
// 超过重试次数,转死信队列
channel.basicNack(tag, false, false);
}
}
}
}指数退避实现:
@Service
public class OrderConsumer {
// 计算延迟时间(指数退避)
private long calculateDelay(int retryCount) {
return (long) Math.pow(2, retryCount) * 1000; // 1s, 2s, 4s, 8s...
}
@RabbitListener(queues = "order.queue", ackMode = "MANUAL")
public void handleOrder(OrderEvent event, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag,
@Header("retry_count", required = false) Integer retryCount) {
retryCount = retryCount == null ? 0 : retryCount;
try {
orderService.process(event);
channel.basicAck(tag, false);
} catch (Exception e) {
if (retryCount < 5) {
// 指数退避延迟
long delay = calculateDelay(retryCount);
log.warn("处理失败,{}毫秒后重试: orderId={}", delay, event.getOrderId());
// 发送到延迟队列
rabbitTemplate.convertAndSend(
"order.delay.exchange",
"order.delay",
event,
message -> {
message.getMessageProperties().setDelay((int) delay);
message.getMessageProperties().setHeader("retry_count", retryCount + 1);
return message;
}
);
channel.basicAck(tag, false);
} else {
// 超过最大重试次数,转死信
log.error("超过最大重试次数,转死信: orderId={}", event.getOrderId());
channel.basicNack(tag, false, false);
}
}
}
}死信队列(Dead Letter Queue)
什么是死信队列?
存储无法被正常消费的消息的队列,用于后续人工处理或补偿。
死信产生的三种情况:
1. 消息被拒绝(reject/nack)且不重新入队
↓
2. 消息过期(TTL 到期)
↓
3. 队列满了(超过最大长度)
↓
转入死信队列死信队列配置:
@Configuration
public class DeadLetterConfig {
// 死信交换机
@Bean
public DirectExchange deadLetterExchange() {
return new DirectExchange("dlx.exchange");
}
// 死信队列
@Bean
public Queue deadLetterQueue() {
return new Queue("dlx.queue", true);
}
// 死信绑定
@Bean
public Binding deadLetterBinding() {
return BindingBuilder
.bind(deadLetterQueue())
.to(deadLetterExchange())
.with("dlx.routing.key");
}
// 业务队列(配置死信队列)
@Bean
public Queue orderQueue() {
return QueueBuilder.durable("order.queue")
.withArgument("x-dead-letter-exchange", "dlx.exchange")
.withArgument("x-dead-letter-routing-key", "dlx.routing.key")
.build();
}
}死信队列处理:
@Service
public class DeadLetterConsumer {
@Autowired
private AlertService alertService;
@Autowired
private FailedMessageService failedMessageService;
@RabbitListener(queues = "dlx.queue")
public void handleDeadLetter(Message message) {
log.error("收到死信消息: {}", new String(message.getBody()));
// 1. 记录到数据库
FailedMessage failedMessage = new FailedMessage();
failedMessage.setPayload(new String(message.getBody()));
failedMessage.setReason("处理失败,转入死信");
failedMessage.setCreateTime(LocalDateTime.now());
failedMessageService.save(failedMessage);
// 2. 发送告警
alertService.sendAlert("死信告警",
String.format("消息ID: %s, 队列: %s",
message.getMessageProperties().getMessageId(),
message.getMessageProperties().getConsumerQueue()));
}
}一个稳妥的消费流程
┌─────────────────────────────────────────────┐
│ 稳妥的消费流程(最佳实践) │
├─────────────────────────────────────────────┤
│ 1. 接收消息 │
│ - 手动 ACK 模式 │
├─────────────────────────────────────────────┤
│ 2. 幂等性检查 │
│ - 业务唯一键去重 │
│ - 已处理则直接 ACK │
├─────────────────────────────────────────────┤
│ 3. 执行业务逻辑 │
│ - 区分临时错误和永久错误 │
├─────────────────────────────────────────────┤
│ 4. 业务成功 │
│ - 手动 ACK │
│ - 记录日志 │
├─────────────────────────────────────────────┤
│ 5. 业务失败 │
│ - 临时错误: 延迟重试 │
│ - 永久错误: 转死信队列 │
├─────────────────────────────────────────────┤
│ 6. 监控与告警 │
│ - 监控死信队列数量 │
│ - 监控重试次数 │
└─────────────────────────────────────────────┘这里最重要的点不是"要不要重试",而是"什么错误值得重试"。
值得重试的错误:
- 网络抖动
- 短暂依赖超时
- 数据库连接池满
- 第三方服务限流
不值得重试的错误:
- 参数脏数据
- 字段缺失
- 业务状态非法
- 数据格式错误
幂等为什么是消费可靠性的底线
只要系统容忍重试,就必须容忍重复消费。
典型重复消费场景
场景1:消费者业务成功,但 ACK 失败
消费者收到消息
↓
执行业务逻辑 → 成功 √
↓
发送 ACK → 网络故障 ×
↓
Broker 认为消息未消费
↓
重新投递
↓
再次执行业务逻辑(重复!) ×场景2:消费者超时重启
消费者收到消息
↓
执行业务逻辑 → 成功 √
↓
还没来得及 ACK
↓
消费者超时重启
↓
Broker 重新投递
↓
再次执行业务逻辑(重复!) ×场景3:上游应用重复发送
生产者发送消息
↓
网络超时
↓
生产者重试
↓
Broker 收到两条相同消息
↓
消费者处理两次(重复!) ×常见幂等手段
1. 业务唯一单号
@Service
public class OrderConsumer {
@Autowired
private StringRedisTemplate redisTemplate;
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
String dedupKey = "order:dedup:" + event.getOrderNo();
// SETNX: 如果 key 不存在则设置
Boolean isFirst = redisTemplate.opsForValue()
.setIfAbsent(dedupKey, "1", Duration.ofDays(7));
if (Boolean.FALSE.equals(isFirst)) {
log.info("重复消息,跳过: orderNo={}", event.getOrderNo());
return;
}
orderService.process(event);
}
}2. 数据库唯一约束
CREATE TABLE `t_order` (
`id` bigint NOT NULL AUTO_INCREMENT,
`order_no` varchar(64) NOT NULL,
PRIMARY KEY (`id`),
UNIQUE KEY `uk_order_no` (`order_no`) -- 唯一约束
) ENGINE=InnoDB;@Service
public class OrderService {
@Transactional
public void createOrder(OrderEvent event) {
Order order = new Order();
order.setOrderNo(event.getOrderNo());
try {
orderMapper.insert(order);
} catch (DuplicateKeyException e) {
log.warn("订单已存在,跳过: orderNo={}", event.getOrderNo());
// 唯一约束冲突,说明订单已存在,幂等返回
}
}
}3. Redis 幂等标记
@Service
public class PaymentConsumer {
@Autowired
private StringRedisTemplate redisTemplate;
public void handlePayment(PaymentEvent event) {
String key = "payment:idempotent:" + event.getPaymentId();
Boolean success = redisTemplate.opsForValue()
.setIfAbsent(key, UUID.randomUUID().toString(), Duration.ofDays(7));
if (Boolean.FALSE.equals(success)) {
log.warn("重复支付: paymentId={}", event.getPaymentId());
return;
}
try {
paymentService.process(event);
} catch (Exception e) {
// 失败时删除标记,允许重试
redisTemplate.delete(key);
throw e;
}
}
}4. 状态机约束
@Service
public class OrderService {
// 订单状态流转: CREATED → PAID → SHIPPED → COMPLETED
public boolean markPaid(Long orderId) {
Order order = orderMapper.selectById(orderId);
// 只有 CREATED 状态才能变为 PAID
if (order.getStatus() != OrderStatus.CREATED) {
log.warn("订单状态不允许支付: orderId={}, status={}",
orderId, order.getStatus());
return false; // 幂等返回
}
// 乐观锁更新
int rows = orderMapper.updateStatus(
orderId, OrderStatus.PAID, OrderStatus.CREATED);
return rows > 0;
}
}幂等方案对比
| 方案 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| 业务唯一键 | 简单,性能好 | 需要业务有唯一键 | 订单号、流水号 |
| 数据库唯一约束 | 强一致性 | 影响性能 | 核心业务 |
| Redis 标记 | 高性能 | 需要考虑 Redis 故障 | 高并发场景 |
| 状态机 | 业务语义清晰 | 状态设计复杂 | 订单、工单 |
最佳实践:组合方案
@Service
public class OrderConsumer {
@Autowired
private StringRedisTemplate redisTemplate;
@Autowired
private OrderService orderService;
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
String dedupKey = "order:dedup:" + event.getOrderNo();
// 1. Redis 快速去重(第一道防线)
Boolean isFirst = redisTemplate.opsForValue()
.setIfAbsent(dedupKey, "1", Duration.ofDays(7));
if (Boolean.FALSE.equals(isFirst)) {
log.info("Redis 去重命中: orderNo={}", event.getOrderNo());
return;
}
try {
// 2. 数据库唯一约束兜底(第二道防线)
orderService.createOrder(event);
} catch (DuplicateKeyException e) {
log.warn("数据库去重命中: orderNo={}", event.getOrderNo());
} catch (Exception e) {
// 其他异常,删除标记,允许重试
redisTemplate.delete(dedupKey);
throw e;
}
}
}消息积压的本质
积压本质是:生产速度持续大于消费速度
生产速度: 10000 消息/秒
消费速度: 1000 消息/秒
积压速度: 9000 消息/秒积压的常见原因
1. 消费逻辑太慢
// × 慢消费: 每条消息查 3 次数据库
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
Order order = orderMapper.selectById(event.getOrderId()); // 查 DB: 10ms
User user = userMapper.selectById(order.getUserId()); // 查 DB: 10ms
Product product = productMapper.selectById(order.getProductId()); // 查 DB: 10ms
// 单条消息耗时: 30ms+
}优化方案:
// √ 快速消费: 批量查询 + 缓存
@RabbitListener(queues = "order.queue", containerFactory = "batchFactory")
public void handleOrders(List<OrderEvent> events) {
// 批量查询
List<Long> orderIds = events.stream()
.map(OrderEvent::getOrderId)
.collect(Collectors.toList());
Map<Long, Order> orderMap = orderMapper.selectBatchIds(orderIds)
.stream()
.collect(Collectors.toMap(Order::getId, Function.identity()));
// 批量处理
events.forEach(event -> {
Order order = orderMap.get(event.getOrderId());
// 处理逻辑
});
// 单条消息耗时: 3ms
}2. 下游依赖慢
// × 下游慢: 调用外部接口超时
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
orderService.process(event);
smsService.sendSms(event.getPhone()); // 短信接口超时: 5秒+
// 消费线程被阻塞
}优化方案:
// √ 异步化: 外部调用异步处理
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
orderService.process(event);
// 异步发送短信
CompletableFuture.runAsync(() -> {
smsService.sendSms(event.getPhone());
}, asyncExecutor);
}3. 消费并发不够
// × 并发度低: 单线程消费
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
// 单线程处理
}优化方案:
// √ 提高并发度
@RabbitListener(queues = "order.queue", concurrency = "10-20")
public void handleOrder(OrderEvent event) {
// 10-20 个线程并发处理
}
// 或增加消费者实例
// Kubernetes: kubectl scale deployment order-consumer --replicas=54. 某一类异常消息反复重试
// × 毒消息阻塞: 异常消息反复重试
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
try {
orderService.process(event);
} catch (Exception e) {
log.error("处理失败", e);
throw new RuntimeException(e); // 重新入队,再次消费
}
}优化方案:
// √ 毒消息隔离: 限制重试次数
@RabbitListener(queues = "order.queue", ackMode = "MANUAL")
public void handleOrder(OrderEvent event, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag,
@Header("retry_count", required = false) Integer retryCount) {
retryCount = retryCount == null ? 0 : retryCount;
try {
orderService.process(event);
channel.basicAck(tag, false);
} catch (Exception e) {
if (retryCount < 3) {
// 重新入队
channel.basicNack(tag, false, true);
} else {
// 超过重试次数,转死信队列
channel.basicNack(tag, false, false);
}
}
}5. 顺序消费场景被单条慢消息卡住
// × 顺序消费: 单条慢消息阻塞整个队列
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
if (event.isSlowMessage()) {
Thread.sleep(30000); // 慢消息处理 30 秒
}
// 后续消息全部被阻塞
}优化方案:
// √ 拆分队列: 快慢消息分离
@RabbitListener(queues = "order.fast.queue")
public void handleFastOrder(OrderEvent event) {
// 快速消息队列
}
@RabbitListener(queues = "order.slow.queue", concurrency = "5-10")
public void handleSlowOrder(OrderEvent event) {
// 慢速消息队列,提高并发度
}积压治理核心原则
所以治理积压不能只盯 Broker 指标,而要把消费耗时、失败率、重试次数、线程池饱和度一起看。
实战场景
场景一:订单支付成功后的事件消费
需求: 支付成功后,需要更新订单状态、扣库存、发优惠券、记积分
挑战:
- 积分服务短时超时
- 不能直接丢消息
- 不能无限同步重试
- 不能因为一个下游故障把整个消费线程堵死
解决方案:
@Service
@Slf4j
public class PaymentConsumer {
@Autowired
private OrderService orderService;
@Autowired
private InventoryService inventoryService;
@Autowired
private CouponService couponService;
@Autowired
private PointsService pointsService;
@Autowired
private StringRedisTemplate redisTemplate;
@Autowired
private DeadLetterService deadLetterService;
@RabbitListener(queues = "payment.queue", ackMode = "MANUAL")
public void handlePaymentSuccess(PaymentEvent event, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) {
String dedupKey = "payment:dedup:" + event.getPaymentId();
try {
// 1. 幂等性检查
Boolean isFirst = redisTemplate.opsForValue()
.setIfAbsent(dedupKey, "1", Duration.ofDays(7));
if (Boolean.FALSE.equals(isFirst)) {
log.info("重复消息,跳过: paymentId={}", event.getPaymentId());
channel.basicAck(tag, false);
return;
}
// 2. 更新订单状态(核心逻辑,必须成功)
orderService.markPaid(event.getOrderId());
// 3. 扣库存(核心逻辑,必须成功)
inventoryService.deduct(event.getOrderId());
// 4. 发优惠券(非核心,失败可补偿)
try {
couponService.grant(event.getUserId());
} catch (Exception e) {
log.error("优惠券发放失败,记录补偿: userId={}", event.getUserId(), e);
compensationService.record(new CompensationTask(
"coupon", event.getUserId(), "grant"
));
}
// 5. 记积分(非核心,失败可补偿)
try {
pointsService.grant(event.getUserId(), event.getAmount());
} catch (Exception e) {
log.error("积分发放失败,记录补偿: userId={}", event.getUserId(), e);
compensationService.record(new CompensationTask(
"points", event.getUserId(), "grant"
));
}
// 6. 业务成功,ACK
channel.basicAck(tag, false);
log.info("支付成功处理完成: paymentId={}", event.getPaymentId());
} catch (Exception e) {
log.error("支付成功处理失败: paymentId={}", event.getPaymentId(), e);
// 删除幂等标记,允许重试
redisTemplate.delete(dedupKey);
// NACK,重新入队
try {
channel.basicNack(tag, false, true);
} catch (IOException ex) {
log.error("NACK 失败", ex);
}
}
}
}场景二:秒杀流量导致消息积压
需求: 秒杀活动,瞬间 10万请求,数据库只能承受 1000 TPS
治理顺序:
1. 确认积压原因
├─ 生产突增? (正常,秒杀流量)
├─ 消费变慢? (检查下游依赖)
└─ 异常阻塞? (检查错误日志)
↓
2. 临时扩容
├─ 增加消费者实例
├─ 提高分区数
└─ 提升并发度
↓
3. 限流降级
├─ 非核心链路延后处理
├─ 核心链路限流
└─ 降级非必要功能
↓
4. 优化消费逻辑
├─ 批量处理
├─ 异步化外部调用
└─ 减少数据库查询
↓
5. 恢复后复盘
├─ 分析根因
├─ 优化架构
└─ 完善监控代码实现:
// 秒杀消费者(限速处理)
@Service
@Slf4j
public class SeckillConsumer {
@Autowired
private SeckillService seckillService;
@Autowired
private RedisTemplate redisTemplate;
@Autowired
private RateLimiter rateLimiter;
@RabbitListener(
queues = "seckill.queue",
concurrency = "1-1", // 并发度为 1,限速
ackMode = "MANUAL"
)
public void handleSeckill(SeckillRequest request, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) {
try {
// 限流: 每秒最多 100 条
if (!rateLimiter.tryAcquire()) {
log.warn("触发限流,延迟处理: userId={}", request.getUserId());
Thread.sleep(100); // 延迟 100ms
channel.basicNack(tag, false, true); // 重新入队
return;
}
// 幂等性检查
String dedupKey = "seckill:dedup:" + request.getUserId() + ":" + request.getProductId();
Boolean isFirst = redisTemplate.opsForValue()
.setIfAbsent(dedupKey, "1", Duration.ofDays(1));
if (Boolean.FALSE.equals(isFirst)) {
log.info("重复秒杀请求: userId={}, productId={}",
request.getUserId(), request.getProductId());
channel.basicAck(tag, false);
return;
}
// 创建订单
seckillService.createOrder(request.getUserId(), request.getProductId());
channel.basicAck(tag, false);
log.info("秒杀成功: userId={}, productId={}",
request.getUserId(), request.getProductId());
} catch (Exception e) {
log.error("秒杀失败: userId={}, productId={}",
request.getUserId(), request.getProductId(), e);
// 恢复库存
redisTemplate.opsForValue().increment("seckill:stock:" + request.getProductId());
try {
channel.basicAck(tag, false); // 失败也 ACK,避免反复重试
} catch (IOException ex) {
log.error("ACK 失败", ex);
}
}
}
}场景三:毒消息拖垮消费组
问题: "毒消息"指无论重试多少次都无法成功处理的消息,例如字段缺失、版本不兼容、脏数据
没有死信隔离时的后果:
毒消息到达
↓
消费失败
↓
重新入队
↓
立即再次消费
↓
再次失败
↓
循环... (消费线程一直被占满)
↓
正常消息被堵在后面
↓
队列积压越来越严重
↓
告警不断,但系统没有真正恢复解决方案:区分可重试错误和不可重试错误
@Service
@Slf4j
public class OrderConsumer {
@Autowired
private OrderService orderService;
@Autowired
private DeadLetterService deadLetterService;
private static final int MAX_RETRY = 3;
@RabbitListener(queues = "order.queue", ackMode = "MANUAL")
public void handleOrder(OrderEvent event, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag,
@Header("retry_count", required = false) Integer retryCount) {
retryCount = retryCount == null ? 0 : retryCount;
try {
// 1. 参数校验
validateEvent(event);
// 2. 业务处理
orderService.process(event);
// 3. 成功 ACK
channel.basicAck(tag, false);
log.info("订单处理成功: orderId={}", event.getOrderId());
} catch (ValidationException e) {
// 参数校验失败,不应该重试,转死信
log.error("参数校验失败,转死信: orderId={}", event.getOrderId(), e);
deadLetterService.sendToDeadLetter(event, e.getMessage());
channel.basicAck(tag, false);
} catch (BusinessException e) {
// 业务异常,不应该重试,转死信
log.error("业务异常,转死信: orderId={}", event.getOrderId(), e);
deadLetterService.sendToDeadLetter(event, e.getMessage());
channel.basicAck(tag, false);
} catch (TransientException e) {
// 临时异常(网络超时、数据库连接池满),可以重试
if (retryCount < MAX_RETRY) {
log.warn("临时异常,准备重试: orderId={}, retryCount={}",
event.getOrderId(), retryCount);
try {
channel.basicNack(tag, false, true);
} catch (IOException ex) {
log.error("NACK 失败", ex);
}
} else {
// 超过最大重试次数,转死信
log.error("超过最大重试次数,转死信: orderId={}", event.getOrderId());
deadLetterService.sendToDeadLetter(event, "重试次数超限");
channel.basicAck(tag, false);
}
} catch (Exception e) {
// 未知异常,记录日志,转死信
log.error("未知异常,转死信: orderId={}", event.getOrderId(), e);
deadLetterService.sendToDeadLetter(event, e.getMessage());
channel.basicAck(tag, false);
}
}
private void validateEvent(OrderEvent event) throws ValidationException {
if (event.getOrderId() == null) {
throw new ValidationException("订单ID不能为空");
}
if (event.getUserId() == null) {
throw new ValidationException("用户ID不能为空");
}
if (event.getAmount() == null || event.getAmount().compareTo(BigDecimal.ZERO) <= 0) {
throw new ValidationException("订单金额必须大于0");
}
}
}
// 自定义异常
public class ValidationException extends Exception {
public ValidationException(String message) {
super(message);
}
}
public class BusinessException extends Exception {
public BusinessException(String message) {
super(message);
}
}
public class TransientException extends Exception {
public TransientException(String message) {
super(message);
}
}排查与治理思路
先看哪些指标
排查积压时优先关注:
@Component
@Slf4j
public class MqMonitor {
@Autowired
private RabbitAdmin rabbitAdmin;
@Autowired
private MeterRegistry meterRegistry;
@Scheduled(fixedRate = 60000)
public void monitor() {
// 1. 堆积消息数
Properties props = rabbitAdmin.getQueueProperties("order.queue");
Integer messageCount = (Integer) props.get("messageCount");
// 2. 消费延迟
Long lagTime = calculateLagTime();
// 3. 单条消费耗时
Double consumeTime = meterRegistry.get("rabbitmq.consume.time")
.timer()
.mean(TimeUnit.MILLISECONDS);
// 4. 重试次数和失败率
Long retryCount = meterRegistry.get("rabbitmq.retry.count")
.counter()
.count();
Double failureRate = meterRegistry.get("rabbitmq.failure.rate")
.gauge()
.value();
// 5. 消费线程池活跃数
Integer activeThreads = getActiveConsumerThreads();
// 6. 下游依赖 RT 与错误率
Double dbRt = meterRegistry.get("database.rt")
.gauge()
.value();
Double dbErrorRate = meterRegistry.get("database.error.rate")
.gauge()
.value();
log.info("队列监控: messageCount={}, lagTime={}ms, consumeTime={}ms, " +
"retryCount={}, failureRate={}%, activeThreads={}, dbRt={}ms, dbErrorRate={}%",
messageCount, lagTime, consumeTime, retryCount, failureRate * 100,
activeThreads, dbRt, dbErrorRate * 100);
// 告警
if (messageCount > 10000) {
alertService.sendAlert("消息积压告警",
String.format("队列: order.queue, 积压: %d", messageCount));
}
if (failureRate > 0.05) {
alertService.sendAlert("消费失败率告警",
String.format("失败率: %.2f%%", failureRate * 100));
}
}
}指标判断逻辑:
如果只有堆积上涨,但消费耗时正常:
→ 生产流量突增(正常)
→ 考虑临时扩容
如果堆积上涨,同时失败率飙升:
→ 下游依赖或消息内容出问题
→ 检查错误日志
→ 检查下游依赖健康状态
如果堆积上涨,同时消费耗时上升:
→ 消费逻辑变慢
→ 检查是否有慢 SQL
→ 检查外部接口响应时间
如果堆积上涨,同时重试次数异常:
→ 有毒消息阻塞
→ 检查死信队列
→ 隔离毒消息治理动作怎么排优先级
优先级排序:
┌─────────────────────────────────────────────┐
│ 优先级1: 先止血 │
│ - 限流(保护下游) │
│ - 降级(暂停非核心消费者) │
│ - 熔断(避免雪崩) │
├─────────────────────────────────────────────┤
│ 优先级2: 再恢复吞吐 │
│ - 扩容消费者实例 │
│ - 提升分区并发 │
│ - 拆分慢逻辑 │
├─────────────────────────────────────────────┤
│ 优先级3: 再隔离异常 │
│ - 把毒消息转死信队列 │
│ - 避免阻塞主消费链路 │
│ - 人工介入处理 │
├─────────────────────────────────────────────┤
│ 优先级4: 最后回到根因 │
│ - 补幂等 │
│ - 补监控 │
│ - 补重试分类和补偿机制 │
└─────────────────────────────────────────────┘代码实现:
@Service
@Slf4j
public class OrderConsumer {
@Autowired
private OrderService orderService;
@Autowired
private RateLimiter rateLimiter;
@Autowired
private CircuitBreaker circuitBreaker;
@RabbitListener(queues = "order.queue", ackMode = "MANUAL")
public void handleOrder(OrderEvent event, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) {
// 1. 限流(优先级1)
if (!rateLimiter.tryAcquire()) {
log.warn("触发限流: orderId={}", event.getOrderId());
try {
channel.basicNack(tag, false, true); // 重新入队,延后处理
} catch (IOException e) {
log.error("NACK 失败", e);
}
return;
}
// 2. 熔断(优先级1)
if (circuitBreaker.isOpen()) {
log.warn("熔断器打开,暂停消费: orderId={}", event.getOrderId());
try {
channel.basicNack(tag, false, true);
} catch (IOException e) {
log.error("NACK 失败", e);
}
return;
}
try {
// 3. 业务处理
orderService.process(event);
circuitBreaker.recordSuccess();
channel.basicAck(tag, false);
} catch (Exception e) {
circuitBreaker.recordFailure();
log.error("处理失败: orderId={}", event.getOrderId(), e);
// 4. 降级(优先级1)
if (event.isNonCritical()) {
log.info("非核心业务,降级处理: orderId={}", event.getOrderId());
channel.basicAck(tag, false);
return;
}
try {
channel.basicNack(tag, false, true);
} catch (IOException ex) {
log.error("NACK 失败", ex);
}
}
}
}顺序消费要特别谨慎
顺序消费的特点:
优点:
- 保证消息按发送顺序消费
- 适合订单状态流转、账户流水处理
代价:
- 单分区吞吐受限
- 一条慢消息可能阻塞整个分区
- 重试策略设计不当时,积压会迅速放大
什么时候需要顺序消费?
√ 需要顺序消费的场景:
- 同一订单的状态流转: CREATED → PAID → SHIPPED
- 同一账户的余额变动: 充值 → 消费 → 提现
- 同一用户的操作日志
× 不需要顺序消费的场景:
- 不同订单的处理
- 不同用户的通知
- 独立的业务事件顺序消费的实现:
RabbitMQ: 单队列单消费者
// 单队列单消费者,保证顺序
@RabbitListener(queues = "order.sequence.queue", concurrency = "1-1")
public void handleOrderSequence(OrderEvent event) {
// 单线程处理,保证顺序
orderService.process(event);
}Kafka: 单分区单消费者
// 单分区单消费者,保证顺序
@KafkaListener(
topics = "order.topic",
groupId = "order-group",
concurrency = 1 // 单线程
)
public void handleOrderSequence(OrderEvent event) {
orderService.process(event);
}RocketMQ: 顺序消息
// RocketMQ 顺序消息
@Service
public class OrderProducer {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void sendOrderMessage(OrderEvent event) {
// 相同订单号发送到同一队列
rocketMQTemplate.syncSendOrderly(
"order.topic",
event,
event.getOrderId().toString() // 队列选择键
);
}
}
@Service
@RocketMQMessageListener(
topic = "order.topic",
consumerGroup = "order-group",
consumeMode = ConsumeMode.ORDERLY // 顺序消费
)
public class OrderConsumer implements RocketMQListener<OrderEvent> {
@Override
public void onMessage(OrderEvent event) {
orderService.process(event);
}
}顺序消费的优化:
// × 问题: 单条慢消息阻塞整个队列
@RabbitListener(queues = "order.sequence.queue", concurrency = "1-1")
public void handleOrderSequence(OrderEvent event) {
if (event.isSlowMessage()) {
Thread.sleep(30000); // 慢消息处理 30 秒
}
// 后续消息全部被阻塞
}
// √ 优化: 快慢消息分离
@Service
public class OrderProducer {
@Autowired
private RabbitTemplate rabbitTemplate;
public void sendOrderMessage(OrderEvent event) {
if (event.isSlowMessage()) {
// 慢消息发到慢队列
rabbitTemplate.convertAndSend("order.slow.exchange", "order.slow", event);
} else {
// 快消息发到快队列
rabbitTemplate.convertAndSend("order.fast.exchange", "order.fast", event);
}
}
}
// 快队列消费者
@RabbitListener(queues = "order.fast.queue", concurrency = "1-1")
public void handleFastOrder(OrderEvent event) {
orderService.process(event);
}
// 慢队列消费者(提高并发度)
@RabbitListener(queues = "order.slow.queue", concurrency = "5-10")
public void handleSlowOrder(OrderEvent event) {
orderService.process(event);
}示例代码
完整的消费可靠性实现
@Service
@Slf4j
public class OrderConsumer {
@Autowired
private StringRedisTemplate redisTemplate;
@Autowired
private OrderService orderService;
@Autowired
private DeadLetterService deadLetterService;
@Autowired
private MeterRegistry meterRegistry;
private static final int MAX_RETRY = 3;
@RabbitListener(queues = "order.queue", ackMode = "MANUAL")
public void onMessage(OrderEvent event, Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag,
@Header("retry_count", required = false) Integer retryCount) {
long start = System.currentTimeMillis();
String dedupKey = "order:dedup:" + event.getOrderId();
retryCount = retryCount == null ? 0 : retryCount;
try {
// 1. 幂等性检查
Boolean isFirst = redisTemplate.opsForValue()
.setIfAbsent(dedupKey, "1", Duration.ofHours(6));
if (Boolean.FALSE.equals(isFirst)) {
log.info("重复消息,跳过: orderId={}", event.getOrderId());
channel.basicAck(tag, false);
return;
}
// 2. 参数校验
validateEvent(event);
// 3. 业务处理
orderService.processOrder(event);
// 4. 成功 ACK
channel.basicAck(tag, false);
// 5. 记录指标
long cost = System.currentTimeMillis() - start;
meterRegistry.timer("rabbitmq.consume.time").record(cost, TimeUnit.MILLISECONDS);
meterRegistry.counter("rabbitmq.consume.success").increment();
log.info("订单处理成功: orderId={}, cost={}ms", event.getOrderId(), cost);
} catch (ValidationException e) {
// 参数校验失败,不应该重试
handlePermanentFailure(event, channel, tag, dedupKey, "参数校验失败: " + e.getMessage());
} catch (BusinessException e) {
// 业务异常,不应该重试
handlePermanentFailure(event, channel, tag, dedupKey, "业务异常: " + e.getMessage());
} catch (TransientException e) {
// 临时异常,可以重试
handleTransientFailure(event, channel, tag, dedupKey, retryCount, e);
} catch (Exception e) {
// 未知异常
handlePermanentFailure(event, channel, tag, dedupKey, "未知异常: " + e.getMessage());
meterRegistry.counter("rabbitmq.consume.unknown.error").increment();
}
}
// 处理永久性失败
private void handlePermanentFailure(OrderEvent event, Channel channel, long tag,
String dedupKey, String reason) {
log.error("永久性失败,转死信: orderId={}, reason={}", event.getOrderId(), reason);
// 删除幂等标记
redisTemplate.delete(dedupKey);
// 转死信队列
deadLetterService.sendToDeadLetter(event, reason);
try {
channel.basicAck(tag, false);
} catch (IOException e) {
log.error("ACK 失败", e);
}
meterRegistry.counter("rabbitmq.consume.permanent.failure").increment();
}
// 处理临时性失败
private void handleTransientFailure(OrderEvent event, Channel channel, long tag,
String dedupKey, Integer retryCount, Exception e) {
// 删除幂等标记,允许重试
redisTemplate.delete(dedupKey);
if (retryCount < MAX_RETRY) {
log.warn("临时异常,准备重试: orderId={}, retryCount={}",
event.getOrderId(), retryCount, e);
try {
// NACK,重新入队
channel.basicNack(tag, false, true);
} catch (IOException ex) {
log.error("NACK 失败", ex);
}
meterRegistry.counter("rabbitmq.consume.retry").increment();
} else {
// 超过最大重试次数,转死信
log.error("超过最大重试次数,转死信: orderId={}", event.getOrderId());
deadLetterService.sendToDeadLetter(event, "重试次数超限");
try {
channel.basicAck(tag, false);
} catch (IOException ex) {
log.error("ACK 失败", ex);
}
meterRegistry.counter("rabbitmq.consume.retry.exhausted").increment();
}
}
private void validateEvent(OrderEvent event) throws ValidationException {
if (event.getOrderId() == null) {
throw new ValidationException("订单ID不能为空");
}
if (event.getUserId() == null) {
throw new ValidationException("用户ID不能为空");
}
if (event.getAmount() == null || event.getAmount().compareTo(BigDecimal.ZERO) <= 0) {
throw new ValidationException("订单金额必须大于0");
}
}
}
// 死信服务
@Service
public class DeadLetterService {
@Autowired
private RabbitTemplate rabbitTemplate;
@Autowired
private FailedMessageMapper failedMessageMapper;
@Autowired
private AlertService alertService;
public void sendToDeadLetter(OrderEvent event, String reason) {
// 1. 发送到死信队列
DeadLetterMessage deadLetter = new DeadLetterMessage();
deadLetter.setPayload(JSON.toJSONString(event));
deadLetter.setReason(reason);
deadLetter.setCreateTime(LocalDateTime.now());
rabbitTemplate.convertAndSend("dlx.exchange", "dlx.order", deadLetter);
// 2. 记录到数据库
FailedMessage failedMessage = new FailedMessage();
failedMessage.setQueueName("order.queue");
failedMessage.setPayload(JSON.toJSONString(event));
failedMessage.setReason(reason);
failedMessage.setCreateTime(LocalDateTime.now());
failedMessageMapper.insert(failedMessage);
// 3. 发送告警
alertService.sendAlert("死信告警",
String.format("队列: order.queue, 原因: %s, 订单ID: %s",
reason, event.getOrderId()));
}
}
// 死信消费者
@Service
@Slf4j
public class DeadLetterConsumer {
@Autowired
private FailedMessageService failedMessageService;
@RabbitListener(queues = "dlx.order.queue")
public void handleDeadLetter(DeadLetterMessage message) {
log.error("收到死信消息: {}", message);
// 记录到数据库
failedMessageService.save(message);
// 可以在这里实现人工处理或自动补偿逻辑
}
}常见误区
× 误区一:以为 MQ 引进来之后就自动具备最终一致性
错误认知:
- "用了 MQ,数据就最终一致了"
- "消息投递成功,就等于业务成功"
实际情况:
- 消息到了,不代表业务一定正确处理了
- 需要补充:补偿、幂等、告警
- 需要监控:消费延迟、失败率、死信数量
× 误区二:所有异常都一律重试
错误示例:
// × 所有异常都重试
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
try {
orderService.process(event);
} catch (Exception e) {
log.error("处理失败", e);
throw new RuntimeException(e); // 所有异常都重新入队
}
}正确做法:
// √ 区分异常类型
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
try {
orderService.process(event);
} catch (BusinessException e) {
// 业务异常,不应该重试,转死信
deadLetterService.sendToDeadLetter(event, e.getMessage());
} catch (TransientException e) {
// 临时异常,可以重试
throw new RuntimeException(e);
}
}× 误区三:只监控 Broker 存活,不监控消费延迟
问题:
- Broker 活着,不代表消费正常
- 消费者可能卡住
- 积压可能越来越严重
正确做法:
// 监控消费延迟
@Scheduled(fixedRate = 60000)
public void monitor() {
// 队列深度
Integer messageCount = getQueueMessageCount();
// 消费延迟
Long lagTime = calculateLagTime();
// 消费速率
Double consumeRate = getConsumeRate();
// 失败率
Double failureRate = getFailureRate();
// 死信数量
Integer deadLetterCount = getDeadLetterCount();
// 告警
if (messageCount > 10000 || lagTime > 60000 || failureRate > 0.05) {
alertService.sendAlert("消费异常告警", ...);
}
}× 误区四:看到积压就一味扩容消费者
问题:
- 可能下游依赖才是瓶颈
- 可能消息处理逻辑本身慢
- 扩容可能加剧下游压力
正确做法:
1. 先确认根因:
- 消费耗时是否正常?
- 下游依赖是否健康?
- 是否有毒消息?
2. 再决定方案:
- 下游瓶颈 → 优化下游或限流
- 逻辑慢 → 优化代码
- 毒消息 → 隔离处理
- 真缺消费者 → 才扩容× 误区五:顺序消费、批量消费、重试机制一起乱配
问题配置:
// × 顺序消费 + 批量消费 + 自动重试,冲突了
@RabbitListener(
queues = "order.queue",
concurrency = "5-10", // 多线程并发
containerFactory = "batch" // 批量消费
)
public void handleOrders(List<OrderEvent> events) {
// 无法保证顺序
// 批量失败时,整批重试,影响效率
}正确配置:
// √ 顺序消费: 单线程
@RabbitListener(queues = "order.sequence.queue", concurrency = "1-1")
public void handleOrderSequence(OrderEvent event) {
orderService.process(event);
}
// √ 批量消费: 多线程
@RabbitListener(queues = "order.batch.queue", concurrency = "10-20")
public void handleOrderBatch(List<OrderEvent> events) {
orderService.processBatch(events);
}面试补充
1. 为什么 MQ 更现实的语义通常是"至少一次"?
答案: 因为网络、宕机、确认失败都会导致重复投递。
- 生产者发送成功,但确认失败 → 生产者重发
- 消费者处理成功,但 ACK 失败 → Broker 重新投递
- 消费者宕机,未确认的消息重新投递
所以"恰好一次"在分布式系统中很难实现,通常接受"至少一次" + 消费者幂等。
2. 消费者为什么必须幂等?
答案: 因为消息重复在真实系统里是常态,不是偶发异常。
重复原因:
- 网络问题导致 ACK 丢失
- 消费者重启导致重新投递
- 生产者重试导致重复发送
如果消费逻辑没有幂等保护,会导致重复扣款、重复发货等严重问题。
3. 积压排查第一步看什么?
答案: 先区分是生产突增、消费变慢,还是异常重试导致线程被占满。
排查步骤:
- 看队列积压量
- 看消费耗时
- 看失败率和重试次数
- 看消费线程池状态
- 看下游依赖健康状态
判断逻辑:
- 积压 + 耗时正常 → 生产突增
- 积压 + 耗时上升 → 消费逻辑慢
- 积压 + 失败率高 → 异常重试阻塞
4. 死信队列的价值是什么?
答案: 把无法自动恢复的消息从主链路隔离出来,避免拖垮正常消费。
价值:
- 隔离毒消息,避免阻塞主队列
- 提供排障入口,便于定位问题
- 支持人工处理或自动补偿
- 保护系统稳定性
5. 顺序消费为什么不能乱用?
答案: 它会显著压缩吞吐,并放大单条慢消息的阻塞效应。
代价:
- 单分区吞吐受限
- 一条慢消息阻塞整个分区
- 重试策略设计复杂
- 扩容困难
建议:只有确实要求顺序的业务才应使用顺序消费,例如同一订单状态流转、同一账户流水处理,而不是默认所有消息都要有序。
实战理解题
题目一:设计订单支付成功的消费可靠性方案
需求:
- 支付成功后,需要更新订单状态、扣库存、发优惠券、记积分
- 要保证不丢消息、不重复处理
- 要处理下游服务超时的情况
参考方案:
支付成功消息
↓
消费者接收
↓
幂等性检查(Redis SETNX)
↓
参数校验
↓
┌─────────┴─────────┐
│ 核心业务(必须成功) │
│ - 更新订单状态 │
│ - 扣库存 │
└─────────┬─────────┘
↓
┌─────────┴─────────┐
│ 非核心业务(可补偿) │
│ - 发优惠券 │
│ - 记积分 │
└─────────┬─────────┘
↓
成功 ACK
↓
失败处理:
- 临时异常 → 延迟重试(指数退避)
- 永久异常 → 转死信队列
- 下游超时 → 异步化 + 补偿题目二:设计消息积压的应急处理方案
需求:
- 秒杀活动导致队列积压 100万消息
- 数据库只能承受 1000 TPS
- 要在不影响用户体验的情况下处理
参考方案:
1. 确认积压原因(1分钟)
- 查看监控: 队列深度、消费耗时、失败率
- 判断: 生产突增(正常) or 消费变慢(异常)
2. 临时止血(5分钟)
- 启用限流: 保护数据库
- 非核心消费者暂停
- 启用熔断器
3. 扩容恢复(10分钟)
- 增加消费者实例: 5 → 20
- 提升分区数: 10 → 50
- 提高并发度: 1-1 → 5-10
4. 持续监控(1小时)
- 监控队列深度变化
- 监控消费速率
- 监控下游依赖压力
5. 恢复后复盘
- 分析根因
- 优化架构
- 完善应急预案参考资料:
- RabbitMQ Reliability Guide
- Kafka Consumer Configuration
- RocketMQ Message Reliability
- 《深入理解 Kafka:核心设计与实践原理》
- 《RabbitMQ 实战指南》
最后更新: 2026-03-30
版本差异(旧版 → 当前)
| 组件 | 旧版(本文编写时) | 当前 |
|---|---|---|
| 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。