一. 消息丢失

1. 生产者发送消息到 broker,消息丢失
2. 消息写入broker 时,消息丢失
    ① 消息写入内存后非正常关机;
    ② 消息写入磁盘后,磁盘坏了;
    ③ master 数据备份到 slave 由于网络原因导致备份失败;
3. 消息者从broker 拉取消息时 ,消息丢失

死信

正常情况下无法被消费的消息 - 死信

死信队列

存放无法被正常消费的消息


1. 保证消息发送到 mq,消息不丢失

* 单向发送

消息被发送,就不再管了 - 消息容易丢失

* 同步发送

消息被发送,等待 broker 响应,设置超时机制|重试机制,超出设定的阈值,转发死信队列处理

* 异步发送

生产者注册一个监听器,通过回调接口处理消息broker响应事件,发送失败|超时可重试

* 事务消息机制

1) 生产者发送半消息给 broker,验证 broker是否正常工作,broker回复生产者一个 ack

2) 执行本地事务,发送一个状态给broker:
        commit:本地事务执行成功,发送一个commit 状态,broker可将消息推送给下游服务
        rollback:本地事务执行失败,发送一个rollback 状态,broker 直接丢弃消息
        unknown:本地事务执行超时,发送一个unknow状态触发事务回查机制,broker 默认重复确认15次,根据回查函数返回的状态,执行提交|回滚操作

2. 保证消息写入不丢失(broker消息的持久化)

1)mq的刷盘策略选择:消息被写入CPU缓存(PageCache)的刷盘策略

同步刷盘:broker.conf 文件配置 flushDiskType=SYNC_FLUSH

异步刷盘:broker.conf 文件配置 flushDisk/XMLSchema=ASYNC_FLUSH

同步刷盘

通过定时任务每10ms 执行一次刷盘操作,而不是一条条写入缓存

异步刷盘

写入数据量达到系统内存阈值,操作系统主动执行刷盘操作

2)避免 broker 单点故障

采用集群方式,避免单节点故障,最大程度保证mq 可用

3. 保证 consumer 端消费消息不丢失

1)同步消费机制

broker 推送消息给 consumer 失败(未收到consumer ack消息),会重新推送给其他消费者,不需要考虑consumer 端消息丢失问题,只需要考虑消息被重复消费问题


2)异步消费机制

broker推送消息给consumer,返回成功,consumer 实际上未接收到消息,从而丢失消息


4)重复消费问题 - 消息的幂等性处理

发送消息时,每条消息指定一个MessageID,基于MessageID判断消息是否被处理过,在数据库中保存被处理过的消息

二. 消息的存储

CommitLog | ConsumerQueue | IndexLog 是存储消息的核心文件,实现高效的消息存储与检索

1. CommitLog - 物理存储文件

rocketmq 所有消息的统一物理存储介质,无论属于哪个主题(topic)| 队列(queue),最终都会按发送顺序写入 CommitLog

1) 内容与组成

内容

每条消息在 CommitLog 中以二进制格式存储,包含完整的消息元数据 payload(消息体)

组成

消息长度、魔数(校验用)、主题(topic)名称、queueId、messageId、消息标志(flag)、消息存储时间戳、消息体(body)、消息属性(延迟级别、事务状态等)、CRC 校验码等

2) 存储形式

文件作为存储单位,每个文件默认大小为 1G(可配置),文件名用消息的起始偏移量命名( 00000000000000000000000000000000001073741824 等,单位字节)。消息按只追加方式写入,避免随机写,保证高写入性能

3) 核心作用

* 保存所有消息的原始数据,是消息可靠性的基础(支持持久化刷盘)

* 提供消息的物理偏移量(offset),用于定位消息在磁盘上的具体位置

2. ConsumeQueue - 消息的逻辑索引文件

ConsumeQueue 是基于 CommitLog 的逻辑索引文件,用于快速定位某个主题(topic)下某个队列(queue)的消息在 CommitLog 中的位置,相当于 “目录” | “索引表”

1)内容与特点

内容

每条记录是一个固定大小(20 字节)的索引条目,包含三部分信息:

CommitLog Offset(8 字节)

消息的物理偏移量(可直接找到消息在 CommitLog 中的存储位置)

Size(4 字节)

消息的长度(结合偏移量可读取完整消息)

Tags HashCode(8 字节)

消息标签|子标题(Tag)的哈希值(用于过滤消息)

2)存储形式

* 按 “topic + queue id” 维度划分,一个队列对应一个独立的 ConsumeQueue 目录( store/consumequeue/TopicA/0/)
* 每个 ConsumeQueue 目录下包含多个索引文件,每个文件默认 30 万条索引条目(固定大小约 5.72MB),文件名用消息在队列中的逻辑偏移量offset)命名

3)核心作用

* 提供高效的消息检索能力

 消费者订阅指定 topic 下 queue 的消息时,可通过 ConsumeQueue 快速定位到消息在 CommitLog 中的存储位置,避免遍历 CommitLog

* 支持按 Tag 过滤消息

通过比较 Tags HashCode 快速过滤掉不符合条件的消息,减少对 CommitLog 文件的 IO 操作

3. CommitLog 与 ConsumerQueue 的关系与协作流程

1)消息写入流程

* 发送消息时,先将消息写入 CommitLog,生成物理偏移量
* 再交由后台线程(ReputMessageService)将消息的索引条目(offset、size、tag hashCode)异步写入 对应 Topic + Queue ConsumeQueue

2)消息读取流程

* 消费者拉取消息时,先根据 “topic + queue id + 消费偏移量(Offset)” 从 ConsumeQueue 中获取索引条目

* 再通过索引条目中的 Offset Size,从 CommitLog 中读取完整的消息内容

3)总结

CommitLog 是 消息的“物理存储”,存储所有消息的完整数据,保证可靠性和写入性能
ConsumeQueue 是消息的 “逻辑索引”,按 topic+queue 分区,存储消息在 CommitLog 中的位置信息,提升消息检索效率。二者协作,既保证了海量消息的高效写入,又支持了按主题、队列的快速查询,是 RocketMQ 高吞吐、高可用的核心设计之一。

4. Index File(索引文件)

按 “消息键(Key)” | “唯一消息 ID” 查询的辅助索引结构,与 CommitLog 和 ConsumeQueue 共同构成消息的3层存储结构(物理存储、逻辑队列索引、键值索引)

1) ConsumeQueue 的使用场景

正常生产与消费消息

所有与 “按队列顺序拉取消息” 的相关操作都必须依赖于ConsumeQueue 实现,具体场景包括:

Consumer 拉取消息(Push | Pull 模式),订阅某个 topic 后,按照 “topic + 队列 id” 维度,通过 ConsumeQueue 定位消息,consumer 记录自身的 “消费偏移量(Offset)”,表示已经消费到对应队列的第几条消息;

消息投递与消费流程

* 每次拉取时,根据 Offset ConsumeQueue 读取对应的索引条目(消息在 CommitLog 中的偏移量和大小);

* 再通过 CommitLog 记录的消息偏移量和大小,完成消息拉取与消费

消息顺序性保证

需要保证 “同一队列内的消息按发送顺序消费” (订单状态变更消息),需要保证 ConsumeQueue 的 “顺序存储” 

* 消息在 ConsumeQueue 中按发送顺序依次存储(逻辑偏移量 Offset 递增)
* 消费者按 Offset 顺序拉取消息,保证消费顺序性


2) Index File 的使用场景

按消息 Key | ID 反向查询( Key | ID 对应 消息在CommitLog文件中的物理偏移量)

通过messageIndexEnable属性配置打开|关闭 Index File 功能

Index File 作为 辅助查询工具,仅用于 “通过消息的(Key | ID)定位消息” 的场景,不参与正常的生产消费流程

具体包括:

业务回溯

通过 Key 查询消息,当业务需要验证 “某个具体消息 Key 对应的消息是否发送成功”、“查询某个具体消息 Key 对应的消息内容” (支付订单号对应的支付结果消息),可通过 IndexFile 索引查询:底层通过 Key 的哈希值从 IndexFile 中找到对应的 CommitLog 偏移量,进而读取消息

例如:

* 调用 DefaultMQAdminExt.viewMessageByKey() 接口;

* 在 RocketMQ 控制台输入 Key 查询;

问题排查

通过 msgId 查询消息,消息发送后,生成唯一的 msgId(客户端|服务端生成),当需要根据 msgId 定位具体消息(排查消息丢失、重复消费问题)时,IndexFile 也会用于创建msgId的索引 (部分版本中 msgId 被当作特殊 Key )

注意

只有当包含了Key的消息被投递到broker Index File 才会为消息建立索引

若消息无 Key,则无法通过 Index File 查询。

Index File 是可选的辅助结构,无 Index File文件,消息的生产和消费也能正常进行(失去了按 Key 查询的能力)

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐