RocketMQNotes

第 03 章:核心概念

zjc 于 2026-01-03 发布

这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 理解 RocketMQ 的核心概念,是正确设计 Topic、消费组、位点、重试和死信的前提。本章把对象模型和生产语义集中梳理。

3.1 对象模型

NameServer
  -> Broker
     -> Topic
        -> MessageQueue
           -> CommitLog offset
           -> ConsumeQueue index

Consumer Group
  -> Consumer instance
  -> MessageQueue assignment
  -> ConsumeOffset

核心对象:

对象 说明
NameServer 路由注册与发现
Broker 消息存储与服务节点
Topic 逻辑消息分类
MessageQueue Topic 内并行单元
ProducerGroup 生产者逻辑分组
ConsumerGroup 消费进度和重试隔离单位
Tag Topic 内过滤标签
Key 消息业务键,用于查询
Offset 队列内消息位点

3.2 Topic 与 MessageQueue

一个 Topic 有多个队列:

OrderTopic
  |-- queue 0
  |-- queue 1
  |-- queue 2
  |-- queue 3

队列作用:

  1. 并行写入;
  2. 并行消费;
  3. 顺序隔离;
  4. 扩容并行度;
  5. 分配给消费组实例。

读写队列数:

参数 含义
writeQueueNums 可写入队列数
readQueueNums 可读取队列数

扩容时通常先加写队列再扩读队列,缩容相反,避免消息不可读。

3.3 消费组

集群模式下:

Consumer Group order-service
  |-- Consumer A -> queue 0, queue 1
  |-- Consumer B -> queue 2, queue 3

特点:

  1. 同组消费者分摊队列;
  2. 消费进度按组隔离;
  3. 重试和死信按组隔离;
  4. 同组逻辑应一致;
  5. 组内实例扩容上限受队列数影响。

广播模式下:

Consumer A 消费 queue 0-3 全量
Consumer B 消费 queue 0-3 全量

广播进度通常保存在消费者本地,适合本地缓存刷新等场景。

3.4 消息标识

每条消息常见标识:

字段 用途
messageId Broker 生成的消息 ID
uniqKey 客户端生成的唯一键
keys 业务 key,可多个,用于查询
tag Topic 内过滤
shardingKey 顺序消息选择队列
transactionId 事务消息 ID

业务建议:

keys = orderId / orderNo
tag = OrderCreated / OrderPaid
property: eventId / traceId

不要依赖 messageId 做业务幂等,应使用业务唯一键。

3.5 消息过滤

Tag 过滤:

consumer.subscribe("OrderTopic", "OrderCreated || OrderPaid");

SQL92 过滤:

consumer.subscribe("OrderTopic", MessageSelector.bySql("region = 'east'"));
方式 特点
Tag 简单,Broker 端先过滤
SQL92 表达式更强,需 Broker 开启
客户端过滤 灵活但浪费网络

过滤不能替代领域模型。如果多个事件语义差异巨大,应评估拆 Topic。

3.6 消费位点

位点记录消费组在队列中的进度:

queue maxOffset = 1000
consumer offset = 800
lag = 200

位点提交原则:

  1. 集群模式业务成功后提交;
  2. 失败进入重试;
  3. 不能为了跳过异常提前提交;
  4. 重置位点必须审批和备份;
  5. 位点错误可能导致漏消息或重复消费。

RocketMQ 默认语义是至少一次,因此消费必须幂等。

3.7 重试队列

消费失败后,消息会进入消费组对应的重试队列。

命名通常形如:

%RETRY%order-consumer-group

特点:

  1. 按消费组隔离;
  2. 有递增延迟;
  3. 达到最大重试次数进入死信;
  4. 重试次数可配置;
  5. 异常类型应区分可重试与不可重试。

不可重试异常:

  1. 参数缺失;
  2. JSON 格式错误;
  3. 业务规则明确拒绝;
  4. 版本不兼容。

这类消息应进入死信并人工处理,而不是无限重试。

3.8 死信队列

命名通常形如:

%DLQ%order-consumer-group

治理要求:

  1. 监控死信数量;
  2. 查看原始消息和异常;
  3. 修复后支持重放;
  4. 无法处理要有业务闭环;
  5. 定期清理归档;
  6. 记录处理人、原因和结果。

死信不是垃圾桶,而是待处理异常任务池。

3.9 消息语义

语义 说明
At most once 可能丢,不会重复
At least once 不丢,可能重复
Exactly once 端到端严格一次

RocketMQ 常规业务消费按至少一次设计。要达到业务上的“一次效果”,必须:

  1. 消息带唯一业务 ID;
  2. 消费端幂等;
  3. 数据库唯一约束;
  4. 状态机限制;
  5. 对账兜底。

3.10 常见设计错误

错误 后果
一个 Topic 混所有业务 权限、治理和过滤复杂
Topic 过细 管理成本高
消费组混用不同逻辑 位点和重试互相影响
依赖自动创建 Topic 资源失控
用 messageId 做幂等 无法表达业务语义
重置位点不当 漏消息或重复风暴
死信无人处理 业务事实丢失
队列数随意扩缩 顺序和分配异常

本章小结

RocketMQ 以 Topic 和 MessageQueue 组织并行度,以 ConsumerGroup 隔离消费进度、重试和死信。消息应携带业务 key 和事件 ID,消费按至少一次设计并实现幂等。重试队列和死信队列是业务可靠性的一部分,必须有监控、重放和处理闭环。

思考题

  1. Topic、Tag 和消费组分别适合什么隔离粒度?
  2. readQueueNums 和 writeQueueNums 扩缩容要注意什么?
  3. 消费位点提前提交会造成什么风险?
  4. 哪些异常不应该重试?
  5. 如何治理死信队列?