这是《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
队列作用:
- 并行写入;
- 并行消费;
- 顺序隔离;
- 扩容并行度;
- 分配给消费组实例。
读写队列数:
| 参数 | 含义 |
|---|---|
| writeQueueNums | 可写入队列数 |
| readQueueNums | 可读取队列数 |
扩容时通常先加写队列再扩读队列,缩容相反,避免消息不可读。
3.3 消费组
集群模式下:
Consumer Group order-service
|-- Consumer A -> queue 0, queue 1
|-- Consumer B -> queue 2, queue 3
特点:
- 同组消费者分摊队列;
- 消费进度按组隔离;
- 重试和死信按组隔离;
- 同组逻辑应一致;
- 组内实例扩容上限受队列数影响。
广播模式下:
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
位点提交原则:
- 集群模式业务成功后提交;
- 失败进入重试;
- 不能为了跳过异常提前提交;
- 重置位点必须审批和备份;
- 位点错误可能导致漏消息或重复消费。
RocketMQ 默认语义是至少一次,因此消费必须幂等。
3.7 重试队列
消费失败后,消息会进入消费组对应的重试队列。
命名通常形如:
%RETRY%order-consumer-group
特点:
- 按消费组隔离;
- 有递增延迟;
- 达到最大重试次数进入死信;
- 重试次数可配置;
- 异常类型应区分可重试与不可重试。
不可重试异常:
- 参数缺失;
- JSON 格式错误;
- 业务规则明确拒绝;
- 版本不兼容。
这类消息应进入死信并人工处理,而不是无限重试。
3.8 死信队列
命名通常形如:
%DLQ%order-consumer-group
治理要求:
- 监控死信数量;
- 查看原始消息和异常;
- 修复后支持重放;
- 无法处理要有业务闭环;
- 定期清理归档;
- 记录处理人、原因和结果。
死信不是垃圾桶,而是待处理异常任务池。
3.9 消息语义
| 语义 | 说明 |
|---|---|
| At most once | 可能丢,不会重复 |
| At least once | 不丢,可能重复 |
| Exactly once | 端到端严格一次 |
RocketMQ 常规业务消费按至少一次设计。要达到业务上的“一次效果”,必须:
- 消息带唯一业务 ID;
- 消费端幂等;
- 数据库唯一约束;
- 状态机限制;
- 对账兜底。
3.10 常见设计错误
| 错误 | 后果 |
|---|---|
| 一个 Topic 混所有业务 | 权限、治理和过滤复杂 |
| Topic 过细 | 管理成本高 |
| 消费组混用不同逻辑 | 位点和重试互相影响 |
| 依赖自动创建 Topic | 资源失控 |
| 用 messageId 做幂等 | 无法表达业务语义 |
| 重置位点不当 | 漏消息或重复风暴 |
| 死信无人处理 | 业务事实丢失 |
| 队列数随意扩缩 | 顺序和分配异常 |
本章小结
RocketMQ 以 Topic 和 MessageQueue 组织并行度,以 ConsumerGroup 隔离消费进度、重试和死信。消息应携带业务 key 和事件 ID,消费按至少一次设计并实现幂等。重试队列和死信队列是业务可靠性的一部分,必须有监控、重放和处理闭环。
思考题
- Topic、Tag 和消费组分别适合什么隔离粒度?
- readQueueNums 和 writeQueueNums 扩缩容要注意什么?
- 消费位点提前提交会造成什么风险?
- 哪些异常不应该重试?
- 如何治理死信队列?