尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

RocketMQ 消息查询机制详解:按 MessageId 与 MessageKey 的查询实现原理

RocketMQ 消息查询机制详解:按 MessageId 与 MessageKey 的查询实现原理 消息队列后端微服务流处理【免费下载链接】rocketmqApache RocketMQ is a cloud native messaging and streaming platform, making it simple to build event-driven applications.项目地址https://gitcode.com/gh_mirrors/ro/rocketmq点击查看免费下载Apache RocketMQ 为消息查询提供了两条核心维度——按 MessageId 查询与按 Message Key 查询分别服务于精确定位单条消息与按业务键批量回溯消息两类场景。本文以 Design_Query.md 为骨架结合仓库中store、broker、common、client等模块的源码实现完整剖析两条查询链路的工作原理包括 MessageId 的编码结构、VIEW_MESSAGE_BY_ID请求的处理流程以及 IndexFile 索引文件的二进制布局、写入与查询算法帮助读者在排查消息丢失、消息轨迹回溯等场景时能够精准定位问题、理解查询性能边界。一、RocketMQ 消息查询的两大维度RocketMQ 支持两种消息查询方式查询维度定位依据典型场景按 MessageId 查询MessageId 中编码的 Broker 地址IP 端口与 CommitLog 物理偏移量精确取回某一条消息排查发送、存储结果按 Message Key 查询消息属性中的UNIQ_KEY或KEYS借助 IndexFile 索引按业务键订单号、流水号等回溯一段时间窗口内的消息两者在 RocketMQ 的请求协议中分别对应两个业务请求码QUERY_MESSAGE12与VIEW_MESSAGE_BY_ID33二者都在 RequestCode.java 中定义并由 Broker 端的 QueryMessageProcessor 统一分发处理。二、按 MessageId 查询从 ID 反解出物理位置2.1 MessageId 的 16 字节编码结构RocketMQ 中的 MessageId 总长度为16 字节其内部结构为Broker 地址8 字节前 4 字节为 IP 地址后 4 字节为端口号CommitLog 偏移量8 字节消息在 CommitLog 文件中的物理偏移地址。这个编码逻辑在 MessageDecoder.createMessageId 中实现方法先写入 8 字节的地址addr再写入 8 字节的offset最终通过UtilAll.bytes2string转成十六进制字符串。需要注意的是源码中的msgIDLength addr.limit() 8 ? 16 : 28分支是为了兼容 IPv6 地址28 字节 16 字节地址 8 字节偏移 4 字节端口布局在 IPv4 环境下标准长度就是 16 字节。2.2 查询链路从 ID 到 CommitLog按 MessageId 查询的整体流程如下Client 端解码客户端通过MessageDecoder.decodeMessageId(msgId)将字符串形式的 MessageId 还原为 IP、端口与 CommitLog 偏移量得到 Broker 地址组装 RPC 请求将 Broker 地址与偏移量封装为ViewMessageRequestHeader通过通信层发送给对应 Broker业务请求码为VIEW_MESSAGE_BY_IDBroker 端读取QueryMessageProcessor.viewMessageById 调用messageStore.selectOneMessageByOffset(offset)直接依据 CommitLog 偏移量定位真实消息零拷贝回传命中的消息通过OneMessageTransferNettyFileRegion以页缓存直接写出避免数据从内核态拷贝到用户态再写回的开销。在 Broker 端selectOneMessageByOffset的实现在 DefaultMessageStore.java它首先通过偏移量计算出目标 MappedFile再从 MappedFile 中切片出该消息对应的缓冲区并读取消息大小最终返回SelectMappedBufferResult。此外若请求头中携带了topicBroker 还会对消息做主题匹配校验matchesRequestTopic防止跨主题越权读取——该逻辑对定时消息、事务半消息等真实存储主题RMQ_SYS_SCHEDULE_TOPIC、TIMER_TOPIC、RMQ_SYS_TRANS_HALF_TOPIC等会使用消息属性PROPERTY_REAL_TOPIC还原逻辑主题后比对。由于查询依据是精确的物理偏移量按 MessageId 查询是 O(1) 级别的精确定位不依赖任何索引结构。三、按 Message Key 查询基于 IndexFile 的索引体系按 Message Key 查询的核心是 RocketMQ 的IndexFile索引文件。其逻辑结构与 JDK 中HashMap的实现非常相似通过哈希槽位定位桶再以链表解决哈希冲突。不同的是IndexFile 直接以内存映射文件MappedByteBuffer的形式落盘查询过程无需加载任何额外索引到内存。索引文件默认存放在$HOME/store/index/目录下路径由 StorePathConfigHelper.getStorePathIndex 计算得出storePathRootDir /index文件名以创建时刻的时间戳命名格式为UtilAll.timeMillisToHumanString因此天然按创建时间有序便于按时间范围筛选与过期清理。3.1 文件大小与固定容量IndexFile 的文件大小固定为 420,000,040 字节其计算公式如下40 5,000,000 × 4 20,000,000 × 20 420,000,040 字节即IndexHeader(40B) HashSlotTable(500万 × 4B) IndexLinkedList(2000万 × 20B)。这一容量由 MessageStoreConfig 中的两个配置项决定maxHashSlotNum默认 5,000,000哈希槽数量maxIndexNum默认 20,000,000即5000000 * 4单个索引文件可容纳的最大索引条目数。IndexFile 的物理布局计算与容量校验在 IndexFile.java 的构造函数中完成fileTotalSize IndexHeader.INDEX_HEADER_SIZE hashSlotNum * hashSlotSize indexNum * indexSize其中hashSlotSize 4、indexSize 20。3.2 索引键的构造规则IndexFile 为每个索引条目生成的键key取决于消息属性若消息属性中设置了UNIQ_KEY则使用topic # UNIQ_KEY作为索引键若消息属性中设置了KEYS多个 Key 以空格分隔则对每个 Key 分别使用topic # KEY作为索引键。该构造逻辑位于 IndexService.buildIndex先写入UNIQ_KEY索引再遍历KEYS按MessageConst.KEY_SEPARATOR即空格拆分逐个写入同时还支持对TAGS属性建立topic #T# tags形式的标签索引。对应的索引类型常量定义在 MessageConst.javaINDEX_KEY_TYPE K、INDEX_UNIQUE_TYPE U、INDEX_TAG_TYPE T。说明UNIQ_KEY与KEYS的定义同样位于 MessageConst.java其中PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX UNIQ_KEY、PROPERTY_KEYS KEYS。生产者在发送消息时通过setKeys()、setKeys(tags, keys)等 API 写入KEYS属性客户端会自动生成UNIQ_KEY。四、IndexFile 的二进制结构深度剖析4.1 IndexHeader文件头的 40 字节统计信息IndexFile 的文件头Header包含6 个字段共 40 字节字段字节数偏移含义beginTimestamp80第一条索引对应消息的storeTimestampendTimestamp88最后一条索引对应消息的storeTimestampbeginPhyOffset816第一条索引对应消息的物理偏移量endPhyOffset824最后一条索引对应消息的物理偏移量hashSlotCount432已使用的哈希槽数量indexCount436已写入的索引条目数量beginTimestamp与endTimestamp刻画了该索引文件覆盖的时间窗口beginPhyOffset与endPhyOffset刻画了覆盖的 CommitLog 偏移区间。这些字段的读写由 IndexHeader.java 通过直接操作ByteBuffer完成且使用AtomicLong/AtomicInteger保证并发可见性updateByteBuffer()在 flush 时把内存中的最新统计值写回文件。4.2 索引条目20 字节的四个字段每一条索引数据Index Store Unit包含4 个字段共 20 字节字段字节数说明Key HashkeyHash4索引键的哈希值key.hashCode()取绝对值CommitLog offsetphyOffset8消息在 CommitLog 中的物理偏移量TimestamptimeDiff4时间差值(storeTimestamp - beginTimestamp) / 1000秒NextIndex offsetprevIndex4同一哈希槽链表中的下一条索引位置关于这 4 个字段的实现细节见 IndexFile.java 及putKey方法哈希冲突链当新写入的索引键哈希值keyHash与某槽位已有索引的 keyHash 相同时新条目的NextIndex offset会指向前一条索引数据的位置从而将所有冲突的索引串成一条单向链表。槽位表保存的正是这条单链表的头结点位置指向最新写入的索引序号。时间戳的精巧设计Timestamp字段记录的不是具体时间而是当前消息storeTimestamp与IndexHeader.beginTimestamp的差值单位为秒。这样 4 字节即足以表示一个较大的时间跨度同时查询时只需beginTimestamp timeDiff即可还原真实时间见selectPhyOffset中的timeDiff * 1000L; timeRead beginTimestamp timeDiff。极端值处理写入时若timeDiff超过Integer.MAX_VALUE则截断为最大值小于 0 则置 0IndexFile.java读取时若timeDiff 0则直接终止链表遍历表示该区域尚未写入有效数据。4.3 三层分区结构整个 IndexFile 由三部分连续区域组成Header40 字节存储上述通用统计信息Slot Table4 × 500 万字节 20MB不保存真实索引数据每个槽位只保存其对应单链表的头结点4 字节即该槽位最新一条索引的序号Index Linked List20 × 2000 万字节 400MB真实的索引数据区一个 IndexFile 最多可容纳 2000 万条索引。这种哈希槽 链表的设计决定了 IndexFile 的查询复杂度理想情况下为 O(1) 定位最坏情况下退化为 O(n)同一 keyHash 大量冲突时。五、索引的写入与查询算法5.1 写入putKey 与索引构建索引写入的入口是 IndexService.buildIndex它作为CommitLogDispatchStore被DefaultMessageStore在消息落盘后异步回调属于存储层的分发dispatch组件。写入前还会做两类过滤事务过滤TRANSACTION_ROLLBACK_TYPE事务回滚的消息不建立索引PREPARED与COMMIT类型则照常建立偏移回退保护若msg.getCommitLogOffset() endPhyOffset说明该消息属于已处理过的历史数据如异常恢复场景直接跳过避免重复建索引。putKeyIndexFile.java的写入步骤为计算keyHash indexKeyHashMethod(key)槽位slotPos keyHash % hashSlotNum读取槽位当前值作为新索引条目的prevIndex即链表头在索引区追加写入keyHash、phyOffset、timeDiff、prevIndex四个字段回写槽位使其指向新索引序号并递增indexCount若indexCount 1首条索引初始化beginPhyOffset/beginTimestamp无论是否首条都更新endPhyOffset/endTimestamp。5.2 写入满与滚动换文件当一个 IndexFile 写满indexCount indexNum时isWriteFull()返回 true后续写入会滚动创建新的 IndexFile。在 IndexService.getAndCreateLastIndexFile 中可以看到新文件会以上一个文件的endPhyOffset与endTimestamp作为自己的beginPhyOffset/beginTimestamp保证跨文件的时间与偏移区间连续创建完成后还会异步启动FlushIndexFileThread将旧文件刷盘并把其endTimestamp写入StoreCheckpoint的indexMsgTimestamp用于异常退出后的索引文件恢复校验见load方法中lastExitOK false时对超期文件的销毁逻辑。5.3 查询selectPhyOffset 的链表遍历按 Key 查询的存储层入口是 IndexService.queryOffset它从最新的索引文件倒序向前遍历for (int i indexFileList.size(); i 0; i--)只对isTimeMatched(begin, end)命中的文件执行selectPhyOffset并在f.getBeginTimestamp() begin或结果数量达到上限时提前终止——这是查询按时间窗口收敛的关键优化。selectPhyOffsetIndexFile.java的查询算法为计算 keyHash 与槽位读取链表头序号slotValue沿prevIndex遍历链表依次读取keyHashRead、phyOffsetRead、timeDiff、prevIndexRead还原timeRead beginTimestamp timeDiff仅当keyHash keyHashRead且timeRead落在[begin, end]时间窗口内时将phyOffsetRead加入结果集当结果数达到maxNum、链表越界或timeRead begin时终止遍历。最终queryOffset返回QueryOffsetResult见 QueryOffsetResult.java其中包含物理偏移量列表phyOffsets以及索引文件最新的更新时间和偏移量indexLastUpdateTimestamp/indexLastUpdatePhyoffset供客户端判断索引是否落后于 CommitLog。5.4 从偏移量到消息DefaultMessageStore.queryMessage拿到物理偏移量列表后DefaultMessageStore.queryMessage 负责把偏移量转换为真实消息通过indexService.queryOffset(...)或indexRocksDBStore获得偏移量列表排序后逐个处理用commitLog.getData(offset, false)从 CommitLog 读取该偏移处的SelectMappedBufferResult以首 4 字节为消息大小截取消息将多个消息封装进QueryMessageResult返回给QueryMessageProcessor由处理器通过QueryMessageTransferFileRegion 零拷贝写回客户端若结果为空还会将lastQueryMsgTime更新为最新一条消息的storeTimestamp并最多重试 3 轮以应对索引尚未同步完的场景全部未命中则返回ResponseCode.QUERY_NOT_FOUND提示maybe time range not correct。5.5 相关配置项一览按 Key 查询行为可通过以下 Broker 配置控制均在 MessageStoreConfig.java 中定义可在broker.conf中覆盖配置项默认值说明messageIndexEnabletrue是否启用消息索引构建关闭后按 Key 查询将无索引可用messageIndexSafefalse索引写入是否要求同步刷盘flush 到磁盘后才认为可靠maxHashSlotNum5000000单个 IndexFile 的哈希槽数量maxIndexNum20000000单个 IndexFile 的最大索引条目数indexFileReadEnabletrue是否启用 IndexFile 文件读取索引indexRocksDBEnablefalse是否启用 RocksDB 作为索引存储与 IndexFile 二选一defaultQueryMaxNum32客户端未指定时的默认最大查询条数maxMsgsNumBatch64单次查询结果条数的上限queryOffset中会对 maxNum 做Math.min收敛从源码结构看RocketMQ 的索引体系目前支持两条并行的实现路径传统的IndexFile内存映射文件与IndexRocksDBStoreRocksDB 存储由indexFileReadEnable/indexRocksDBEnable两个开关决定实际生效的查询后端文档与本文重点剖析的是默认启用的 IndexFile 实现。六、两条查询链路总结对比维度按 MessageId 查询按 Message Key 查询请求码VIEW_MESSAGE_BY_ID33QUERY_MESSAGE12定位依据MessageId 中编码的 IP、端口与 CommitLog 偏移IndexFile 中 keyHash 槽位链表依赖结构无直接按偏移读 CommitLogIndexFileHeader Slot Table Index Linked List复杂度O(1) 精确定位理想 O(1)哈希冲突时退化为 O(n)典型用途精确取单条消息如消息轨迹、排查丢失按业务 Key 在时间窗口内回溯批量消息时间约束无查询必须指定begin/end时间窗口并受maxNum限制七、延伸阅读消息查询相关的架构总览可参考 Design_Store.md存储层整体设计与 Design_Remoting.md通信层协议设计Broker 端请求处理上下文见 QueryMessageProcessor.java其测试用例 QueryMessageProcessorTest.java 覆盖了QUERY_MESSAGE与VIEW_MESSAGE_BY_ID两条路径索引文件核心实现 IndexFile.java、IndexHeader.java、IndexService.java存储配置项全集见 MessageStoreConfig.javaBroker 配置示例可参考 distribution/conf/broker.conf。赞分享消息队列后端微服务流处理【免费下载链接】rocketmqApache RocketMQ is a cloud native messaging and streaming platform, making it simple to build event-driven applications.项目地址https://gitcode.com/gh_mirrors/ro/rocketmq点击查看免费下载相关推荐Apache RocketMQ 消息查询机制详解Message Id 直查与 Message Key 索引查询Apache RocketMQ 消息查询机制详解Message Id 直查与 Message Key 索引查询 导读 本文围绕 RocketMQ 的消息查询能消息队列流处理后端Apache RocketMQ消息查询机制深度解析Apache RocketMQ消息查询机制深度解析 一、消息查询概述 在分布式消息系统中消息查询是一个非常重要的功能。Apache RocketMQ作为一款高消息队列后端微服务流处理【快速上手】Apache RocketMQ消息查询机制深度解析Apache RocketMQ消息查询机制深度解析 一、消息查询概述 Apache RocketMQ作为一款分布式消息中间件提供了强大的消息查询能力这对于消消息队列流处理后端上一篇使用Azure认知服务语音SDK在iOS上实现Swift文本转语音快速入门下一篇MM2-0/Kvaesitso 插件开发入门指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表