rocketmq
一. 消息丢失
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(可配置),文件名用消息的起始偏移量命名( 00000000000000000000、000000000000001073741824 等,单位字节)。消息按只追加方式写入,避免随机写,保证高写入性能
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 查询的能力)
更多推荐





所有评论(0)