第 201 章:MQ 选型对比与应用场景
学习目标
- 掌握 Kafka / RabbitMQ / RocketMQ / Redis Stream 对比
- 根据场景选型
- 理解消息中间件核心指标
- 避免常见踩坑
一、四大消息中间件
| 特性 | Kafka | RabbitMQ | RocketMQ | Redis Stream |
|---|---|---|---|---|
| 语言 | Scala/Java | Erlang | Java | C |
| 吞吐量 | 百万/s | 万/s | 十万/s | 万/s |
| 延迟 | ms 级 | us 级 | ms 级 | ms 级 |
| 持久化 | 磁盘 | 内存+磁盘 | 磁盘 | 内存+磁盘 |
| 消息回溯 | ✅ | ❌ | ✅ | ✅ |
| 事务 | ✅ | ✅ | ✅ | ❌ |
| 死信 | ❌ | ✅ | ✅ | ❌ |
| 延迟消息 | ❌(需插件) | ✅ | ✅ | ❌ |
| 优先级队列 | ❌ | ✅ | ❌ | ❌ |
| 集群 | 强 | 中 | 强 | 弱 |
| 运维难度 | 高 | 中 | 中 | 低 |
二、消息中间件核心指标
2.1 四大维度
2.2 可靠性级别
| 级别 | 保证 | 适用 |
|---|---|---|
| 最多一次 | 可能丢 | 日志 |
| 至少一次 | 可能重复 | 订单 |
| 恰好一次 | 不丢不重 | 金融(难) |
三、典型场景选型
3.1 选型决策树
3.2 业务消息(订单、交易)
推荐:RocketMQ 或 RabbitMQ
原因:
- ✅ 事务消息(订单创建 + 库存扣减)
- ✅ 延迟消息(订单超时)
- ✅ 死信队列(失败重试)
- ✅ 万级吞吐够用
3.3 日志 / 事件流(用户行为、IoT)
推荐:Kafka
原因:
- ✅ 百万级吞吐
- ✅ 消息回溯(可重放)
- ✅ 持久化(磁盘顺序写)
- ✅ 流处理生态(Kafka Streams)
3.4 微服务解耦(简单)
推荐:Redis Stream 或 RabbitMQ
原因:
- ✅ 部署简单
- ✅ 路由灵活(RabbitMQ)
- ✅ 成本低
3.5 金融交易(强事务)
推荐:RocketMQ
原因:
- ✅ 原生事务消息
- ✅ 阿里双 11 验证
- ✅ 高可用架构
四、消息可靠性设计
4.1 生产者不丢
java
// Kafka
props.put("acks", "all");
props.put("enable.idempotence", true);
// RabbitMQ
spring.rabbitmq.publisher-confirm-type=correlated
// RocketMQ
producer.setRetryTimesWhenSendFailed(3);
producer.send(msg, new SendCallback() { ... });4.2 Broker 不丢
properties
# Kafka 多副本
replication.factor=3
min.insync.replicas=2
# RabbitMQ 镜像队列
ha-mode: all
ha-sync-mode: automatic4.3 消费者不丢
java
// 业务成功 + 提交 offset
inventoryService.deduct(event);
ack.acknowledge();4.4 事务消息
java
// RocketMQ 事务消息
TransactionMQProducer producer = new TransactionMQProducer("group");
producer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 本地事务
try {
orderService.create(msg);
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 回查(防止本地事务结果未确定)
return orderService.exists(msg.getKeys())
? LocalTransactionState.COMMIT_MESSAGE
: LocalTransactionState.ROLLBACK_MESSAGE;
}
});五、消息幂等设计
5.1 唯一 ID 方案
java
// 生产
Message msg = new Message(topic, body);
msg.setKeys(orderId); // 唯一键
producer.send(msg);
// 消费
String orderId = msg.getKeys();
if (redisTemplate.opsForValue().setIfAbsent("mq:order:" + orderId, "1", 24, TimeUnit.HOURS)) {
// 首次消费
inventoryService.deduct(orderId);
} else {
log.info("重复消息: {}", orderId);
}5.2 数据库唯一约束
sql
CREATE TABLE order_event (
id BIGINT PRIMARY KEY,
order_id BIGINT,
UNIQUE KEY uk_order (order_id)
);六、消息积压治理
6.1 应急
bash
# 1. 临时扩容消费者
kafka-topics --alter --topic orders --partitions 30
# 2. 直接读底层 log(应急)
kafka-console-consumer --bootstrap-server localhost:9092 \
--topic orders --partition 0 --offset earliest6.2 长期方案
- 异步消费
- 批量消费
- 业务降级
- 限流保护
七、消息中间件监控
| 指标 | 监控方式 |
|---|---|
| 发送成功率 | 应用埋点 |
| 消费 Lag | Kafka Lag / RabbitMQ 队列深度 |
| 消费失败率 | 应用埋点 |
| 延迟 | end-to-end 时间戳 |
八、本章小结
| 场景 | 推荐 |
|---|---|
| 业务消息 | RocketMQ / RabbitMQ |
| 日志事件 | Kafka |
| 微服务解耦 | RabbitMQ / Redis Stream |
| 简单队列 | Redis Stream |
| 金融交易 | RocketMQ |
动手练习
- 对比测试 Kafka / RabbitMQ / Redis Stream 性能
- 用事务消息实现下单 + 库存
- 实现消费幂等(Redis 去重)
- 监控 Lag 并配置告警
推荐阅读
下一章:第 202 章:Elasticsearch 入门 →