Lesson 65 · 数据库与中间件实战

消息队列核心原理:RocketMQ vs Kafka 架构对比

高级·🔥 极高·#MQ·#RocketMQ·#Kafka

高级 · 面试极高

消息队列核心原理:RocketMQ vs Kafka 架构对比

消息可靠投递、顺序消息、重复消费幂等性。RocketMQ 的事务消息、Kafka 的分区与消费者组——两大 MQ 的架构差异。

异步解耦削峰填谷RocketMQKafka事务消息幂等消费
第 1 站

为什么需要消息队列?

消息队列(MQ)解决的核心问题可以用三个词概括:异步、解耦、削峰

MQ 三大核心价值 异步 同步调用: A→B→C→D 响应时间 = 200+300+500 = 1000ms 异步: A 发 MQ → 返回 200ms B/C/D 异步消费 解耦 订单系统不直接调用 库存/积分/通知系统 发一条消息 → 谁关心谁消费 下游宕机不影响上游 削峰填谷 秒杀瞬时 10 万 QPS DB 只能扛 2000 QPS MQ 缓冲 → 消费者匀速消费 2000 QPS 处理完所有请求 典型架构 生产者 消息队列(Broker) 消息持久化 + 高可用集群 消费者 A 消费者 B 消费者 C 生产者只管发,消费者只管处理,MQ 负责可靠传递
图 1-1 MQ 三大核心价值:异步、解耦、削峰填谷

引入 MQ 有什么代价?

MQ 引入了额外的复杂性: 系统可用性依赖 MQ(MQ 挂了怎么办?); 一致性问题(消息发了但消费者处理失败怎么办?); 增加了运维成本(MQ 集群的部署、监控、调优)。所以不是所有场景都要用 MQ——同步调用能满足的,不要为了用而用。

第 2 站

RocketMQ 架构

RocketMQ 是阿里开源的分布式消息中间件,专为金融级消息设计,支持事务消息、延迟消息、消息过滤等高级特性。

RocketMQ 四大组件 NameServer 集群 轻量注册中心,存 Broker 路由信息 Broker Master 存储消息、转发消息 Topic → MessageQueue(4个) Broker Slave 异步复制 Master 数据 高可用备份 注册 Producer 发送消息到 Broker Consumer 从 Broker 拉取消息 查路由 查路由 Producer/Consumer 从 NameServer 获取 Broker 路由 → 直连 Broker 收发消息
图 2-1 RocketMQ 四大组件:NameServer、Broker、Producer、Consumer
组件角色关键特性
NameServer轻量注册中心无状态,各节点独立(不用 ZK),存 Broker 路由表
Broker消息存储和转发Master-Slave 高可用,Topic 分为多个 MessageQueue(默认 4 个)
Producer消息生产者支持同步/异步/单向发送,消息路由到 MessageQueue
Consumer消息消费者Push(长轮询)/ Pull 模式,支持集群/广播消费
第 3 站

Kafka 架构——分区与消费者组

Kafka 最初是 LinkedIn 开发的分布式日志系统,现在广泛用于大数据流处理。它的核心概念是 Topic → Partition → Consumer Group

Kafka 架构:Topic → Partition → Consumer Group Topic: order-events (3 Partitions) Partition-0 (Broker-1) offset: 0,1,2,3,4... Partition-1 (Broker-2) offset: 0,1,2,3,4... Partition-2 (Broker-3) offset: 0,1,2,3,4... Consumer Group: order-service (3 consumers) Consumer-A 消费 Partition-0 Consumer-B 消费 Partition-1 Consumer-C 消费 Partition-2 核心规则 ① 一个 Partition 只能被同一个 Consumer Group 中的一个 Consumer 消费 ② 一个 Consumer 可以消费多个 Partition | ③ Consumer 数 <= Partition 数(否则有闲置消费者)
图 3-1 Kafka 分区与消费者组的对应关系
概念RocketMQ 对应Kafka 对应说明
消息主题TopicTopic逻辑分类
物理分片MessageQueuePartition并行度的基本单位
消费者组ConsumerGroupConsumer Group组内负载均衡消费
注册中心NameServerZooKeeper / KRaft元数据管理
存储节点BrokerBroker消息存储和转发
关键结论

Kafka 的 Partition 是并行度的上限——3 个 Partition 最多 3 个 Consumer 并行消费。增加吞吐量 = 增加 Partition 数。RocketMQ 的 MessageQueue 同理。

第 4 站

消息可靠投递——三端保障

消息可靠投递需要从三个环节保障:生产端 → Broker 端 → 消费端

消息可靠投递三道防线 ① 生产端 RocketMQ: 同步发送 + 重试 Kafka: acks=all + retries=3 确保 Broker 确认收到 失败则重试或记录到本地表 ② Broker 端 RocketMQ: 同步刷盘 + Master-Slave 同步复制 Kafka: replication.factor >= 3 + min.insync.replicas >= 2 ③ 消费端 手动提交 offset (处理完再 ACK) RocketMQ: 返回 CONSUME_SUCCESS 才推进 offset 可靠性等级 最高: 同步发送 + 同步刷盘 + 同步复制 + 手动ACK → 零丢失(但性能最低) 平衡: 异步发送 + 异步刷盘 + 异步复制 + 手动ACK → 允许极端情况丢失(高吞吐)
图 4-1 消息可靠投递的三道防线
Kafka acks 机制
Kafka Producer 的 acks 参数:

acks=0   // 不等确认 → 性能最高,可能丢消息
acks=1   // Leader 写入就确认(默认)→ Leader 挂了可能丢
acks=all // 所有 ISR 副本都写入才确认 → 最安全

配合 min.insync.replicas=2:
// ISR 中至少 2 个副本存活才接受写入
// 如果只剩 1 个副本 → 拒绝写入(保护数据一致性)
第 5 站

顺序消息

很多业务场景需要保证消息的顺序性(如订单状态流转:创建→支付→发货)。两个 MQ 的顺序消息实现思路类似:同一业务 ID 的消息路由到同一个队列,同一队列由一个消费者串行消费

顺序消息保证机制 订单消息 orderId % 4 Queue-0: 订单1001 Queue-1: 订单1002 Queue-2: 订单1003 Queue-3: 订单1004 Thread-0 Thread-1 Thread-2 Thread-3 同一订单的 创建/支付/发货 都在同一队列 → 顺序消费
图 5-1 顺序消息:相同业务 ID 路由到同一队列,单线程消费
RocketMQ 顺序消息
// 生产者:指定 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;
    }
});

顺序消息的代价是什么?

顺序消息会大幅降低吞吐量——因为同一队列只能单线程消费,无法并行。如果一条消息消费失败重试,会阻塞该队列后续所有消息。所以只在确实需要顺序的场景使用(如订单状态流转),不要在所有消息上都用顺序模式。

第 6 站

事务消息(RocketMQ 独有)

RocketMQ 的事务消息解决了"本地事务 + 消息发送"的原子性问题。Kafka 也有事务消息,但语义不同(Kafka 事务是多 Partition 原子写入)。

RocketMQ 事务消息流程 Producer Broker Consumer 发送半消息(half message)→ Consumer 不可见 Broker 返回 half 发送成功 执行本地事务(如扣库存、创建订单) ④a 本地事务成功 → Commit → Consumer 可见 ④b 本地事务失败 → Rollback → Broker 删除消息 Broker 回查:长时间未收到 commit/rollback → 回调 Producer 查询事务状态 核心保证:本地事务和消息发送要么都成功,要么都失败(最终一致)
图 6-1 RocketMQ 事务消息:半消息 + 本地事务 + 回查机制
事务消息代码示例
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 + 回查机制。保证"本地事务"和"消息投递"的最终一致性。适用于跨服务的分布式事务场景。

第 7 站

幂等消费——重复消费怎么办?

MQ 的 at-least-once 语义意味着消息可能被重复投递(网络闪断、消费者重启、手动 ACK 超时等)。消费者必须保证幂等性——同一消息处理多次和处理一次效果相同。

四种幂等消费方案 ① 唯一消息 ID + 去重表 INSERT IGNORE INTO dedup(msg_id) 插入成功 → 处理;重复 → 跳过 适合:通用场景 ② 数据库唯一键 订单号做 UNIQUE KEY INSERT 失败(唯一冲突)→ 跳过 适合:创建类操作 ③ 乐观锁 / 状态机 UPDATE SET status='PAID' WHERE order_id=? AND status='UNPAID' 适合:状态流转类操作 ④ Redis SETNX 防重 SETNX dedup:{msgId} 1 EX 3600 返回 1 → 处理;返回 0 → 已处理 适合:高并发场景,快速判断
图 7-1 四种幂等消费方案
最佳实践:去重表 + 状态机
-- 去重表
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 做前置快速过滤。"

全篇回顾

  1. MQ 价值:异步(减少响应时间)、解耦(上下游独立)、削峰(保护 DB)
  2. RocketMQ 架构:NameServer(路由)+ Broker(存储)+ Producer + Consumer
  3. Kafka 架构:Topic → Partition → Consumer Group,Partition 数 = 并行度上限
  4. 消息可靠:生产端(同步 + 重试)+ Broker 端(刷盘 + 复制)+ 消费端(手动 ACK)
  5. 顺序消息:相同业务 ID 路由到同一队列,单线程消费(牺牲吞吐换顺序)
  6. 事务消息:RocketMQ 半消息 + 本地事务 + 回查,保证最终一致性
  7. 幂等消费:去重表 + 唯一键 + 状态机 + Redis SETNX