CC 咖啡猫的工作空间 Coding Space

消息队列实践

异步解耦、流量削峰是消息队列的核心价值。


1. 主流消息队列对比

特性 RabbitMQ RocketMQ Kafka ActiveMQ
吞吐量 万级 十万级 百万级 万级
延迟 微秒级 毫秒级 毫秒级 毫秒级
可靠性 支持事务 支持事务 异步复制 支持事务
生态 成熟稳定 阿里开源 大数据场景 逐渐淘汰
适合场景 小型系统 电商交易 日志/流处理 老系统迁移

2. RabbitMQ 快速入门

2.1 Spring Boot 集成

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
spring:
  rabbitmq:
    host: ${RABBITMQ_HOST:localhost}
    port: ${RABBITMQ_PORT:5672}
    username: ${RABBITMQ_USER:guest}
    password: ${RABBITMQ_PASS:guest}
    virtual-host: /prod

2.2 基本配置

@Configuration
public class RabbitConfig {

    @Bean
    public Queue orderQueue() {
        return QueueBuilder.durable("order.queue")
                .build();
    }

    @Bean
    public DirectExchange orderExchange() {
        return new DirectExchange("order.exchange");
    }

    @Bean
    public Binding orderBinding(Queue orderQueue, DirectExchange orderExchange) {
        return BindingBuilder.bind(orderQueue)
                .to(orderExchange)
                .with("order.routing.key");
    }
}

2.3 生产者

@Service
@RequiredArgsConstructor
public class OrderProducer {

    private final RabbitTemplate rabbitTemplate;

    public void sendOrder(Order order) {
        rabbitTemplate.convertAndSend(
            "order.exchange",
            "order.routing.key",
            order,
            message -> {
                message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
                message.getMessageProperties().setContentEncoding("UTF-8");
                return message;
            }
        );
    }
}

2.4 消费者

@Component
@Slf4j
public class OrderConsumer {

    @RabbitListener(queues = "order.queue", concurrency = "3-10")
    public void handleOrder(Order order, Message message) {
        log.info("收到订单消息: {}", order.getOrderId());
        try {
            // 处理订单业务
            processOrder(order);
        } catch (Exception e) {
            log.error("处理订单失败: {}", order.getOrderId(), e);
            // 消息重试或进入死信队列
            throw new AmqpRejectAndDontRequeueException(e);
        }
    }

    private void processOrder(Order order) {
        // 实际业务处理
    }
}

3. RocketMQ 实践

3.1 Spring Boot 集成

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
</dependency>
rocketmq:
  name-server: ${ROCKETMQ_NAMESRV:localhost:9876}
  producer:
    group: order-producer-group
    send-message-timeout: 3000

3.2 事务消息

@Service
@RequiredArgsConstructor
public class OrderService {

    private final RocketMQTemplate rocketMQTemplate;

    public void createOrder(Order order) {
        // 本地事务与消息发送绑定
        Transaction transaction = rocketMQTemplate.beginTransaction();

        try {
            // 1. 保存订单到数据库
            orderMapper.insert(order);

            // 2. 发送半消息
            rocketMQTemplate.asyncSend("order:create", order, transaction);

            // 3. 提交事务
            transaction.commit();
        } catch (Exception e) {
            transaction.rollback();
            throw e;
        }
    }
}

@Component
@Slf4j
public class OrderTransactionListener implements RocketMQTransactionListener {

    @Override
    public TransactionState executeLocalTransaction(Message msg, Object arg) {
        try {
            // 本地事务执行(订单已保存)
            return TransactionState.COMMIT;
        } catch (Exception e) {
            log.error("本地事务执行失败", e);
            return TransactionState.ROLLBACK;
        }
    }

    @Override
    public TransactionState checkLocalTransaction(MessageExt msg) {
        // 事务回查:检查订单是否已保存
        String orderId = msg.getKeys();
        if (orderMapper.selectById(orderId) != null) {
            return TransactionState.COMMIT;
        }
        return TransactionState.UNKNOWN;
    }
}

4. 常见问题处理

4.1 消息丢失

生产者 → Broker → 消费者

消息丢失场景:
1. 生产者发送时网络抖动
2. Broker 宕机未持久化
3. 消费者处理时异常未确认

解决方案:

场景 方案
生产者到 Broker 开启 publisher confirms
Broker 持久化 同步刷盘 + 集群副本
消费者未处理 手动 ACK + 补偿机制
// RabbitMQ 开启确认
spring:
  rabbitmq:
    publisher-confirm-type: correlated
    publisher-returns: true

// RocketMQ 同步刷盘
broker:
    flushDiskType: SYNC_FLUSH

4.2 重复消费

幂等处理方案:

// 方案1:数据库唯一索引
INSERT INTO order_log (order_id, status) VALUES (#{orderId}, 'PROCESSED')
ON DUPLICATE KEY UPDATE status = status;

// 方案2:Redis 幂等
String key = "order:processed:" + orderId;
if (redis.setIfAbsent(key, "1", 24 hours)) {
    // 首次处理
    processOrder(order);
} else {
    // 重复消息,直接返回
    log.warn("订单已处理: {}", orderId);
}

// 方案3:业务状态机
if (order.getStatus() != OrderStatus.PENDING) {
    log.warn("订单状态非待处理,无需处理: {}", orderId);
    return;
}

4.3 消息顺序性

问题:生产者 → 消息队列 → 消费者,乱序消费

场景:订单创建 → 支付 → 履约,必须按顺序处理

解决方案:

方案 适用场景
消息分区 RocketMQ 分区键,同一订单进同一分区
消费者串行 单线程消费,按序处理
版本号机制 消息带版本号,乱序时缓存等待
// RocketMQ 分区键保证顺序
rocketMQTemplate.asyncSend("order:create", order, new SendCallback() {
    @Override
    public void onSuccess(SendResult result) {
        // 同一 orderId 的消息会进入同一队列
    }
});

5. 消息队列使用场景

场景 推荐队列 说明
异步解耦 RabbitMQ 订单完成后通知下游系统
流量削峰 RocketMQ/Kafka 秒杀场景,抵挡瞬时流量
日志采集 Kafka 高吞吐日志收集
事务消息 RocketMQ 订单与库存一致性
延迟队列 RabbitMQ 订单超时取消

6. 最佳实践

  1. 消息持久化:开启消息和队列持久化,防止 Broker 宕机丢失
  2. 消费者幂等:所有消费者必须实现幂等操作
  3. 死信队列:配置死信队列处理消费失败的消息
  4. 监控告警:监控队列积压、消费延迟、失败率
  5. 分区规划:RocketMQ 按业务分区,避免单队列热点
  6. 容量评估:预估峰值 QPS,预留 2-3 倍 buffer
# RabbitMQ 死信队列配置
spring:
  rabbitmq:
    listener:
      simple:
        default-requeue-rejected: false
        retry:
          enabled: true
          initial-interval: 1000
          max-attempts: 3

# x-dead-letter-exchange 和 x-dead-letter-routing-key 在队列声明时配置

最后更新:2026/05/13