Lesson 65 · 数据库与中间件实战
消息队列核心原理:RocketMQ vs Kafka 架构对比
消息队列核心原理:RocketMQ vs Kafka 架构对比
消息可靠投递、顺序消息、重复消费幂等性。RocketMQ 的事务消息、Kafka 的分区与消费者组——两大 MQ 的架构差异。
为什么需要消息队列?
消息队列(MQ)解决的核心问题可以用三个词概括:异步、解耦、削峰。
引入 MQ 有什么代价?
MQ 引入了额外的复杂性:① 系统可用性依赖 MQ(MQ 挂了怎么办?);② 一致性问题(消息发了但消费者处理失败怎么办?);③ 增加了运维成本(MQ 集群的部署、监控、调优)。所以不是所有场景都要用 MQ——同步调用能满足的,不要为了用而用。
RocketMQ 架构
RocketMQ 是阿里开源的分布式消息中间件,专为金融级消息设计,支持事务消息、延迟消息、消息过滤等高级特性。
| 组件 | 角色 | 关键特性 |
|---|---|---|
| NameServer | 轻量注册中心 | 无状态,各节点独立(不用 ZK),存 Broker 路由表 |
| Broker | 消息存储和转发 | Master-Slave 高可用,Topic 分为多个 MessageQueue(默认 4 个) |
| Producer | 消息生产者 | 支持同步/异步/单向发送,消息路由到 MessageQueue |
| Consumer | 消息消费者 | Push(长轮询)/ Pull 模式,支持集群/广播消费 |
Kafka 架构——分区与消费者组
Kafka 最初是 LinkedIn 开发的分布式日志系统,现在广泛用于大数据流处理。它的核心概念是 Topic → Partition → Consumer Group。
| 概念 | RocketMQ 对应 | Kafka 对应 | 说明 |
|---|---|---|---|
| 消息主题 | Topic | Topic | 逻辑分类 |
| 物理分片 | MessageQueue | Partition | 并行度的基本单位 |
| 消费者组 | ConsumerGroup | Consumer Group | 组内负载均衡消费 |
| 注册中心 | NameServer | ZooKeeper / KRaft | 元数据管理 |
| 存储节点 | Broker | Broker | 消息存储和转发 |
Kafka 的 Partition 是并行度的上限——3 个 Partition 最多 3 个 Consumer 并行消费。增加吞吐量 = 增加 Partition 数。RocketMQ 的 MessageQueue 同理。
消息可靠投递——三端保障
消息可靠投递需要从三个环节保障:生产端 → Broker 端 → 消费端。
Kafka Producer 的 acks 参数: acks=0 // 不等确认 → 性能最高,可能丢消息 acks=1 // Leader 写入就确认(默认)→ Leader 挂了可能丢 acks=all // 所有 ISR 副本都写入才确认 → 最安全 配合 min.insync.replicas=2: // ISR 中至少 2 个副本存活才接受写入 // 如果只剩 1 个副本 → 拒绝写入(保护数据一致性)
顺序消息
很多业务场景需要保证消息的顺序性(如订单状态流转:创建→支付→发货)。两个 MQ 的顺序消息实现思路类似:同一业务 ID 的消息路由到同一个队列,同一队列由一个消费者串行消费。
// 生产者:指定 MessageQueueSelector producer.send(msg, (mqs, msg, arg) -> { int index = Math.abs(orderId.hashCode()) % mqs.size(); return mqs.get(index); }, orderId); // 消费者:MessageListenerOrderly(非 MessageListenerConcurrently) consumer.registerMessageListener(new MessageListenerOrderly() { public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ...) { // 顺序消费逻辑 return ConsumeOrderlyStatus.SUCCESS; } });
顺序消息的代价是什么?
顺序消息会大幅降低吞吐量——因为同一队列只能单线程消费,无法并行。如果一条消息消费失败重试,会阻塞该队列后续所有消息。所以只在确实需要顺序的场景使用(如订单状态流转),不要在所有消息上都用顺序模式。
事务消息(RocketMQ 独有)
RocketMQ 的事务消息解决了"本地事务 + 消息发送"的原子性问题。Kafka 也有事务消息,但语义不同(Kafka 事务是多 Partition 原子写入)。
TransactionMQProducer producer = new TransactionMQProducer("tx-group"); // 设置事务监听器 producer.setTransactionListener(new TransactionListener() { // 执行本地事务(半消息发送成功后回调) @Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { try { orderService.createOrder((OrderDTO) arg); // 本地事务 return LocalTransactionState.COMMIT_MESSAGE; } catch (Exception e) { return LocalTransactionState.ROLLBACK_MESSAGE; } } // 回查接口(Broker 主动查询事务状态) @Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { // 查数据库确认本地事务是否执行成功 boolean success = orderService.checkOrderExists(msg.getKeys()); return success ? COMMIT_MESSAGE : ROLLBACK_MESSAGE; } });
RocketMQ 事务消息 = 半消息(暂存)+ 本地事务 + commit/rollback + 回查机制。保证"本地事务"和"消息投递"的最终一致性。适用于跨服务的分布式事务场景。
幂等消费——重复消费怎么办?
MQ 的 at-least-once 语义意味着消息可能被重复投递(网络闪断、消费者重启、手动 ACK 超时等)。消费者必须保证幂等性——同一消息处理多次和处理一次效果相同。
-- 去重表 CREATE TABLE msg_dedup ( msg_id VARCHAR(64) PRIMARY KEY, created_at DATETIME ); -- 消费逻辑(伪代码) @Transactional public void consume(Message msg) { try { // 1. 去重:INSERT IGNORE(唯一键冲突会抛异常) dedupDao.insert(msg.getId()); } catch (DuplicateKeyException e) { log.info("重复消息,跳过: {}", msg.getId()); return; // 幂等:直接返回成功 } // 2. 业务逻辑 + 状态机防重 int rows = orderDao.updateStatus( orderId, "PAID", "UNPAID" // WHERE status='UNPAID' ); if (rows == 0) { log.info("订单状态已变更,跳过"); } }
"MQ 消息可能重复投递,消费者必须保证幂等。常见做法是用消息唯一 ID 做去重表(INSERT IGNORE),结合业务状态机(UPDATE WHERE status='xxx')双重保障。高并发场景可以用 Redis SETNX 做前置快速过滤。"
全篇回顾
- MQ 价值:异步(减少响应时间)、解耦(上下游独立)、削峰(保护 DB)
- RocketMQ 架构:NameServer(路由)+ Broker(存储)+ Producer + Consumer
- Kafka 架构:Topic → Partition → Consumer Group,Partition 数 = 并行度上限
- 消息可靠:生产端(同步 + 重试)+ Broker 端(刷盘 + 复制)+ 消费端(手动 ACK)
- 顺序消息:相同业务 ID 路由到同一队列,单线程消费(牺牲吞吐换顺序)
- 事务消息:RocketMQ 半消息 + 本地事务 + 回查,保证最终一致性
- 幂等消费:去重表 + 唯一键 + 状态机 + Redis SETNX