消息队列实践
异步解耦、流量削峰是消息队列的核心价值。
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. 最佳实践
- 消息持久化:开启消息和队列持久化,防止 Broker 宕机丢失
- 消费者幂等:所有消费者必须实现幂等操作
- 死信队列:配置死信队列处理消费失败的消息
- 监控告警:监控队列积压、消费延迟、失败率
- 分区规划:RocketMQ 按业务分区,避免单队列热点
- 容量评估:预估峰值 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