{T}

消费可靠性与消息积压治理

概念与背景

消息队列真正难的地方,往往不是"怎么发消息",而是消息发出去之后能不能稳定落库、能不能被正确消费、积压了之后怎么止血。很多线上事故都不是 Broker 直接挂掉,而是消费速度跟不上、重复消费没兜住、失败重试把下游打崩。

所以 MQ 进阶必须重点看三件事:

  • 消费可靠性如何保证
  • 消息失败后如何重试、回退、隔离
  • 消息积压出现后如何判断和治理

消费可靠性的三大挑战

code
┌─────────────────────────────────────────────┐
│          消费可靠性的三大挑战                │
├─────────────────────────────────────────────┤
│  1. 消息丢失                                 │
│     - 消费者宕机,消息未确认                 │
│     - 自动 ACK 后业务失败                   │
│     - 异常被吞掉                             │
├─────────────────────────────────────────────┤
│  2. 消息重复                                 │
│     - 消费成功但 ACK 失败                   │
│     - 消费者重启导致重新投递                │
│     - 生产者重复发送                        │
├─────────────────────────────────────────────┤
│  3. 消息积压                                 │
│     - 消费速度 < 生产速度                   │
│     - 下游依赖故障                          │
│     - 毒消息阻塞队列                        │
└─────────────────────────────────────────────┘

原理与机制

一条消息的可靠性链路

从生产到消费,一条消息通常要经过:

code
1. 生产者发送到 Broker
      ↓
2. Broker 持久化并复制
      ↓
3. 消费者拉取或接收消息
      ↓
4. 消费者执行业务逻辑
      ↓
5. 消费成功后提交确认(ACK)

真正的风险点主要集中在后三步:

风险点1:消费者收到了消息,但业务还没执行完就宕机

code
消费者收到消息
    ↓
开始执行业务逻辑
    ↓
宕机!(业务未完成)
    ↓
消息状态: 未确认
    ↓
Broker 重新投递
    ↓
消息未丢失 √

前提: 使用手动 ACK,而不是自动 ACK

风险点2:业务执行成功了,但确认没有提交成功

code
消费者收到消息
    ↓
执行业务逻辑 → 成功
    ↓
发送 ACK → 失败(网络问题)
    ↓
Broker 认为消息未消费
    ↓
重新投递
    ↓
重复消费 ×

解决方案: 消费者必须幂等

风险点3:消费失败后无限重试,导致下游雪崩

code
消息处理失败
    ↓
重新入队
    ↓
立即再次消费
    ↓
再次失败
    ↓
重新入队
    ↓
循环... (下游被反复调用)
    ↓
下游雪崩 ×

解决方案: 限制重试次数,失败后转死信队列

消费确认、重试与死信

ACK 机制详解

ACK(确认)的作用:

告诉 Broker 这条消息已经成功处理,可以从队列中删除。

三种确认方式:

确认方式说明可靠性性能适用场景
自动确认(Auto)消息到达消费者立即确认非关键业务
手动确认(Manual)业务成功后再确认关键业务
批量确认(Batch)批量确认多条消息大批量处理

代码示例:

java
// × 自动确认(默认): 不可靠
@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 这条消息处理失败
  • 可以选择是否重新入队
java
// basicNack 参数说明
channel.basicNack(
    deliveryTag,  // 消息标识
    multiple,     // 是否批量确认(false=单条)
    requeue       // 是否重新入队
);

Reject:

  • 拒绝单条消息
  • 比 NACK 简单
java
// basicReject 参数说明
channel.basicReject(
    deliveryTag,  // 消息标识
    requeue       // 是否重新入队
);

NACK vs Reject:

操作批量重新入队适用场景
basicNack支持支持批量失败处理
basicReject不支持支持单条失败处理

重试机制

重试的三种策略:

code
┌─────────────────────────────────────────┐
│  策略1: 立即重试(不推荐)                │
│  - 消息重新入队                         │
│  - 立即被再次消费                       │
│  - 风险: 下游压力持续                   │
├─────────────────────────────────────────┤
│  策略2: 延迟重试(推荐)                  │
│  - 消息转延迟队列                       │
│  - 延迟后再次投递                       │
│  - 优点: 给下游恢复时间                 │
├─────────────────────────────────────────┤
│  策略3: 指数退避重试(推荐)              │
│  - 重试间隔指数增长                     │
│  - 1s → 2s → 4s → 8s → 16s            │
│  - 优点: 避免重试风暴                   │
└─────────────────────────────────────────┘

延迟重试实现(RabbitMQ):

java
@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);
            }
        }
    }
}

指数退避实现:

java
@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)

什么是死信队列?

存储无法被正常消费的消息的队列,用于后续人工处理或补偿。

死信产生的三种情况:

code
1. 消息被拒绝(reject/nack)且不重新入队
      ↓
2. 消息过期(TTL 到期)
      ↓
3. 队列满了(超过最大长度)
      ↓
转入死信队列

死信队列配置:

java
@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();
    }
}

死信队列处理:

java
@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()));
    }
}

一个稳妥的消费流程

code
┌─────────────────────────────────────────────┐
│          稳妥的消费流程(最佳实践)           │
├─────────────────────────────────────────────┤
│  1. 接收消息                                 │
│     - 手动 ACK 模式                         │
├─────────────────────────────────────────────┤
│  2. 幂等性检查                               │
│     - 业务唯一键去重                        │
│     - 已处理则直接 ACK                      │
├─────────────────────────────────────────────┤
│  3. 执行业务逻辑                             │
│     - 区分临时错误和永久错误                │
├─────────────────────────────────────────────┤
│  4. 业务成功                                 │
│     - 手动 ACK                              │
│     - 记录日志                              │
├─────────────────────────────────────────────┤
│  5. 业务失败                                 │
│     - 临时错误: 延迟重试                    │
│     - 永久错误: 转死信队列                  │
├─────────────────────────────────────────────┤
│  6. 监控与告警                               │
│     - 监控死信队列数量                      │
│     - 监控重试次数                          │
└─────────────────────────────────────────────┘

这里最重要的点不是"要不要重试",而是"什么错误值得重试"。

值得重试的错误:

  • 网络抖动
  • 短暂依赖超时
  • 数据库连接池满
  • 第三方服务限流

不值得重试的错误:

  • 参数脏数据
  • 字段缺失
  • 业务状态非法
  • 数据格式错误

幂等为什么是消费可靠性的底线

只要系统容忍重试,就必须容忍重复消费。

典型重复消费场景

场景1:消费者业务成功,但 ACK 失败

code
消费者收到消息
    ↓
执行业务逻辑 → 成功 √
    ↓
发送 ACK → 网络故障 ×
    ↓
Broker 认为消息未消费
    ↓
重新投递
    ↓
再次执行业务逻辑(重复!) ×

场景2:消费者超时重启

code
消费者收到消息
    ↓
执行业务逻辑 → 成功 √
    ↓
还没来得及 ACK
    ↓
消费者超时重启
    ↓
Broker 重新投递
    ↓
再次执行业务逻辑(重复!) ×

场景3:上游应用重复发送

code
生产者发送消息
    ↓
网络超时
    ↓
生产者重试
    ↓
Broker 收到两条相同消息
    ↓
消费者处理两次(重复!) ×

常见幂等手段

1. 业务唯一单号

java
@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. 数据库唯一约束

sql
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;
java
@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 幂等标记

java
@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. 状态机约束

java
@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 故障高并发场景
状态机业务语义清晰状态设计复杂订单、工单

最佳实践:组合方案

java
@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;
        }
    }
}

消息积压的本质

积压本质是:生产速度持续大于消费速度

code
生产速度: 10000 消息/秒
消费速度: 1000 消息/秒
积压速度: 9000 消息/秒

积压的常见原因

1. 消费逻辑太慢

java
// × 慢消费: 每条消息查 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+
}

优化方案:

java
// √ 快速消费: 批量查询 + 缓存
@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. 下游依赖慢

java
// × 下游慢: 调用外部接口超时
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
    orderService.process(event);
    smsService.sendSms(event.getPhone());  // 短信接口超时: 5秒+
    // 消费线程被阻塞
}

优化方案:

java
// √ 异步化: 外部调用异步处理
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
    orderService.process(event);
    
    // 异步发送短信
    CompletableFuture.runAsync(() -> {
        smsService.sendSms(event.getPhone());
    }, asyncExecutor);
}

3. 消费并发不够

java
// × 并发度低: 单线程消费
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
    // 单线程处理
}

优化方案:

java
// √ 提高并发度
@RabbitListener(queues = "order.queue", concurrency = "10-20")
public void handleOrder(OrderEvent event) {
    // 10-20 个线程并发处理
}

// 或增加消费者实例
// Kubernetes: kubectl scale deployment order-consumer --replicas=5

4. 某一类异常消息反复重试

java
// × 毒消息阻塞: 异常消息反复重试
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
    try {
        orderService.process(event);
    } catch (Exception e) {
        log.error("处理失败", e);
        throw new RuntimeException(e);  // 重新入队,再次消费
    }
}

优化方案:

java
// √ 毒消息隔离: 限制重试次数
@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. 顺序消费场景被单条慢消息卡住

java
// × 顺序消费: 单条慢消息阻塞整个队列
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
    if (event.isSlowMessage()) {
        Thread.sleep(30000);  // 慢消息处理 30 秒
    }
    // 后续消息全部被阻塞
}

优化方案:

java
// √ 拆分队列: 快慢消息分离
@RabbitListener(queues = "order.fast.queue")
public void handleFastOrder(OrderEvent event) {
    // 快速消息队列
}

@RabbitListener(queues = "order.slow.queue", concurrency = "5-10")
public void handleSlowOrder(OrderEvent event) {
    // 慢速消息队列,提高并发度
}

积压治理核心原则

所以治理积压不能只盯 Broker 指标,而要把消费耗时、失败率、重试次数、线程池饱和度一起看。

实战场景

场景一:订单支付成功后的事件消费

需求: 支付成功后,需要更新订单状态、扣库存、发优惠券、记积分

挑战:

  • 积分服务短时超时
  • 不能直接丢消息
  • 不能无限同步重试
  • 不能因为一个下游故障把整个消费线程堵死

解决方案:

java
@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

治理顺序:

code
1. 确认积压原因
    ├─ 生产突增? (正常,秒杀流量)
    ├─ 消费变慢? (检查下游依赖)
    └─ 异常阻塞? (检查错误日志)
    ↓
2. 临时扩容
    ├─ 增加消费者实例
    ├─ 提高分区数
    └─ 提升并发度
    ↓
3. 限流降级
    ├─ 非核心链路延后处理
    ├─ 核心链路限流
    └─ 降级非必要功能
    ↓
4. 优化消费逻辑
    ├─ 批量处理
    ├─ 异步化外部调用
    └─ 减少数据库查询
    ↓
5. 恢复后复盘
    ├─ 分析根因
    ├─ 优化架构
    └─ 完善监控

代码实现:

java
// 秒杀消费者(限速处理)
@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);
            }
        }
    }
}

场景三:毒消息拖垮消费组

问题: "毒消息"指无论重试多少次都无法成功处理的消息,例如字段缺失、版本不兼容、脏数据

没有死信隔离时的后果:

code
毒消息到达
    ↓
消费失败
    ↓
重新入队
    ↓
立即再次消费
    ↓
再次失败
    ↓
循环... (消费线程一直被占满)
    ↓
正常消息被堵在后面
    ↓
队列积压越来越严重
    ↓
告警不断,但系统没有真正恢复

解决方案:区分可重试错误和不可重试错误

java
@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);
    }
}

排查与治理思路

先看哪些指标

排查积压时优先关注:

java
@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));
        }
    }
}

指标判断逻辑:

code
如果只有堆积上涨,但消费耗时正常:
    → 生产流量突增(正常)
    → 考虑临时扩容

如果堆积上涨,同时失败率飙升:
    → 下游依赖或消息内容出问题
    → 检查错误日志
    → 检查下游依赖健康状态

如果堆积上涨,同时消费耗时上升:
    → 消费逻辑变慢
    → 检查是否有慢 SQL
    → 检查外部接口响应时间

如果堆积上涨,同时重试次数异常:
    → 有毒消息阻塞
    → 检查死信队列
    → 隔离毒消息

治理动作怎么排优先级

优先级排序:

code
┌─────────────────────────────────────────────┐
│  优先级1: 先止血                             │
│  - 限流(保护下游)                           │
│  - 降级(暂停非核心消费者)                   │
│  - 熔断(避免雪崩)                           │
├─────────────────────────────────────────────┤
│  优先级2: 再恢复吞吐                         │
│  - 扩容消费者实例                           │
│  - 提升分区并发                             │
│  - 拆分慢逻辑                               │
├─────────────────────────────────────────────┤
│  优先级3: 再隔离异常                         │
│  - 把毒消息转死信队列                       │
│  - 避免阻塞主消费链路                       │
│  - 人工介入处理                             │
├─────────────────────────────────────────────┤
│  优先级4: 最后回到根因                       │
│  - 补幂等                                   │
│  - 补监控                                   │
│  - 补重试分类和补偿机制                     │
└─────────────────────────────────────────────┘

代码实现:

java
@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);
            }
        }
    }
}

顺序消费要特别谨慎

顺序消费的特点:

优点:

  • 保证消息按发送顺序消费
  • 适合订单状态流转、账户流水处理

代价:

  • 单分区吞吐受限
  • 一条慢消息可能阻塞整个分区
  • 重试策略设计不当时,积压会迅速放大

什么时候需要顺序消费?

code
√ 需要顺序消费的场景:
- 同一订单的状态流转: CREATED → PAID → SHIPPED
- 同一账户的余额变动: 充值 → 消费 → 提现
- 同一用户的操作日志

× 不需要顺序消费的场景:
- 不同订单的处理
- 不同用户的通知
- 独立的业务事件

顺序消费的实现:

RabbitMQ: 单队列单消费者

java
// 单队列单消费者,保证顺序
@RabbitListener(queues = "order.sequence.queue", concurrency = "1-1")
public void handleOrderSequence(OrderEvent event) {
    // 单线程处理,保证顺序
    orderService.process(event);
}

Kafka: 单分区单消费者

java
// 单分区单消费者,保证顺序
@KafkaListener(
    topics = "order.topic",
    groupId = "order-group",
    concurrency = 1  // 单线程
)
public void handleOrderSequence(OrderEvent event) {
    orderService.process(event);
}

RocketMQ: 顺序消息

java
// 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);
    }
}

顺序消费的优化:

java
// × 问题: 单条慢消息阻塞整个队列
@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);
}

示例代码

完整的消费可靠性实现

java
@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,数据就最终一致了"
  • "消息投递成功,就等于业务成功"

实际情况:

  • 消息到了,不代表业务一定正确处理了
  • 需要补充:补偿、幂等、告警
  • 需要监控:消费延迟、失败率、死信数量

× 误区二:所有异常都一律重试

错误示例:

java
// × 所有异常都重试
@RabbitListener(queues = "order.queue")
public void handleOrder(OrderEvent event) {
    try {
        orderService.process(event);
    } catch (Exception e) {
        log.error("处理失败", e);
        throw new RuntimeException(e);  // 所有异常都重新入队
    }
}

正确做法:

java
// √ 区分异常类型
@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 活着,不代表消费正常
  • 消费者可能卡住
  • 积压可能越来越严重

正确做法:

java
// 监控消费延迟
@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("消费异常告警", ...);
    }
}

× 误区四:看到积压就一味扩容消费者

问题:

  • 可能下游依赖才是瓶颈
  • 可能消息处理逻辑本身慢
  • 扩容可能加剧下游压力

正确做法:

code
1. 先确认根因:
   - 消费耗时是否正常?
   - 下游依赖是否健康?
   - 是否有毒消息?

2. 再决定方案:
   - 下游瓶颈 → 优化下游或限流
   - 逻辑慢 → 优化代码
   - 毒消息 → 隔离处理
   - 真缺消费者 → 才扩容

× 误区五:顺序消费、批量消费、重试机制一起乱配

问题配置:

java
// × 顺序消费 + 批量消费 + 自动重试,冲突了
@RabbitListener(
    queues = "order.queue",
    concurrency = "5-10",      // 多线程并发
    containerFactory = "batch" // 批量消费
)
public void handleOrders(List<OrderEvent> events) {
    // 无法保证顺序
    // 批量失败时,整批重试,影响效率
}

正确配置:

java
// √ 顺序消费: 单线程
@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. 积压排查第一步看什么?

答案: 先区分是生产突增、消费变慢,还是异常重试导致线程被占满。

排查步骤:

  1. 看队列积压量
  2. 看消费耗时
  3. 看失败率和重试次数
  4. 看消费线程池状态
  5. 看下游依赖健康状态

判断逻辑:

  • 积压 + 耗时正常 → 生产突增
  • 积压 + 耗时上升 → 消费逻辑慢
  • 积压 + 失败率高 → 异常重试阻塞

4. 死信队列的价值是什么?

答案: 把无法自动恢复的消息从主链路隔离出来,避免拖垮正常消费。

价值:

  • 隔离毒消息,避免阻塞主队列
  • 提供排障入口,便于定位问题
  • 支持人工处理或自动补偿
  • 保护系统稳定性

5. 顺序消费为什么不能乱用?

答案: 它会显著压缩吞吐,并放大单条慢消息的阻塞效应。

代价:

  • 单分区吞吐受限
  • 一条慢消息阻塞整个分区
  • 重试策略设计复杂
  • 扩容困难

建议:只有确实要求顺序的业务才应使用顺序消费,例如同一订单状态流转、同一账户流水处理,而不是默认所有消息都要有序。

实战理解题

题目一:设计订单支付成功的消费可靠性方案

需求:

  • 支付成功后,需要更新订单状态、扣库存、发优惠券、记积分
  • 要保证不丢消息、不重复处理
  • 要处理下游服务超时的情况

参考方案:

code
支付成功消息
    ↓
消费者接收
    ↓
幂等性检查(Redis SETNX)
    ↓
参数校验
    ↓
┌─────────┴─────────┐
│ 核心业务(必须成功) │
│ - 更新订单状态     │
│ - 扣库存           │
└─────────┬─────────┘
    ↓
┌─────────┴─────────┐
│ 非核心业务(可补偿) │
│ - 发优惠券         │
│ - 记积分           │
└─────────┬─────────┘
    ↓
成功 ACK
    ↓
失败处理:
- 临时异常 → 延迟重试(指数退避)
- 永久异常 → 转死信队列
- 下游超时 → 异步化 + 补偿

题目二:设计消息积压的应急处理方案

需求:

  • 秒杀活动导致队列积压 100万消息
  • 数据库只能承受 1000 TPS
  • 要在不影响用户体验的情况下处理

参考方案:

code
1. 确认积压原因(1分钟)
   - 查看监控: 队列深度、消费耗时、失败率
   - 判断: 生产突增(正常) or 消费变慢(异常)

2. 临时止血(5分钟)
   - 启用限流: 保护数据库
   - 非核心消费者暂停
   - 启用熔断器

3. 扩容恢复(10分钟)
   - 增加消费者实例: 5 → 20
   - 提升分区数: 10 → 50
   - 提高并发度: 1-1 → 5-10

4. 持续监控(1小时)
   - 监控队列深度变化
   - 监控消费速率
   - 监控下游依赖压力

5. 恢复后复盘
   - 分析根因
   - 优化架构
   - 完善应急预案

参考资料:

最后更新: 2026-03-30

版本差异(旧版 → 当前)

组件旧版(本文编写时)当前
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。