)
RocketMQ ConsumeQueue 逻辑消费队列源码详解基于主题的消费索引设计与实现4.9.3【免费下载链接】source-code-hunter 从源码层面剖析挖掘互联网行业主流技术的底层实现原理为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶Mybatis、Netty、Dubbo 框架及 Redis、Tomcat 中间件等项目地址: https://gitcode.com/doocs/source-code-hunter本文基于 doocs/source-code-hunter 仓库中 RocketMQ 系列文档 展开聚焦 RocketMQ 存储层中 ConsumeQueue消费队列的设计动机、存储结构、消息定位与文件生命周期管理。阅读完本文你将掌握ConsumeQueue 为何能替代直接遍历 CommitLog 实现高效消费、20 字节索引单元CQ_STORE_UNIT_SIZE的构成、按存储时间二分定位偏移量getOffsetInQueueByTime、跨文件偏移滚动rollNextFile、写入与删除的完整源码路径并能结合仓库内 CommitLog、MappedFile 内存映射、消息发送存储 与 消息消费流程 等姊妹文档形成完整知识链路。一、为什么需要 ConsumeQueueCommitLog 的消费目录文件RocketMQ 基于主题订阅模式实现消息消费消费者关注的是某一个主题Topic下的所有消息。但在存储层RocketMQ 采用单一日志CommitLog的方式顺序存储所有消息同一主题下的消息在 CommitLog 文件中是不连续的——它们会按到达顺序与其它主题的消息交错排列。如果消费者每次消费都直接从 CommitLog 中遍历查找某个主题下的消息就需要在巨大的物理文件中做全量扫描效率极低。为此RocketMQ 设计了 ConsumeQueue 文件可以把它看作CommitLog 的消费目录文件索引文件每条消息在 CommitLog 中的物理位置都会按主题 队列的组织方式登记到 ConsumeQueue 中消费时先查索引、再按物理偏移量精确定位消息从而把遍历大文件变成顺序读小索引 随机定位。ConsumeQueue 采用两级目录结构第一级目录消息主题名称Topic第二级目录该主题下的队列 idQueueId。也就是说Broker 落盘后会出现形如store/consumequeue/{topic}/{queueId}/的目录树其下的每个文件对应这个队列的一段逻辑偏移区间文件以起始逻辑偏移量命名与 CommitLog 文件命名规则一致。更完整的存储布局可对照 CommitLog 详解commitlog目录顺序存储消息本体而consumequeue目录为每个主题队列维护消费索引。二、ConsumeQueue 存储单元20 字节的关键信息为了加速查询并节省磁盘空间ConsumeQueue 并不会存储消息的全量信息每条消息只登记三个关键字段合计CQ_STORE_UNIT_SIZE 20字节字段字节数含义CommitLog 物理偏移量8 字节long消息在 CommitLog 文件中的起始物理偏移量消息大小4 字节int消息在 CommitLog 中占用的字节数Tag 哈希码8 字节longTag 的哈希码tagsCode用于快速过滤其中CQ_STORE_UNIT_SIZE 8偏移量 4大小 8tag 哈希码 20 字节正因为索引单元是定长的ConsumeQueue 文件内的定位才可以转化为纯粹的算术运算某个逻辑偏移量offset对应的单元在文件内的位置就是offset * CQ_STORE_UNIT_SIZE。这也为下一节按时间二分查找提供了前提——所有单元等长二分查找天然成立。三、根据消息存储时间查找物理偏移量getOffsetInQueueByTime在按时间回溯消费、重置消费位点等场景中需要根据消息存储时间找到对应的 ConsumeQueue 偏移量其入口是org.apache.rocketmq.store.ConsumeQueue#getOffsetInQueueByTime。整个过程分为两步按时间戳定位物理文件再在文件内利用二分查找法加速检索。3.1 第一步根据时间戳定位物理文件public MappedFile getMappedFileByTime(final long timestamp) { Object[] mfs this.copyMappedFiles(0); if (null mfs) return null; for (int i 0; i mfs.length; i) { MappedFile mappedFile (MappedFile) mfs[i]; if (mappedFile.getLastModifiedTimestamp() timestamp) { return mappedFile; } } return (MappedFile) mfs[mfs.length - 1]; }算法思想很朴素从第一个文件开始找到第一个最后修改时间大于等于目标时间戳的文件。由于 ConsumeQueue 文件按写入时间顺序创建最后一个文件的修改时间必然是最新的如果所有文件的修改时间都早于目标时间戳则退化返回最后一个文件mfs[mfs.length - 1]。这里用到的MappedFile是 RocketMQ 的内存映射文件抽象其初始化RandomAccessFilefileChannel.map(MapMode.READ_WRITE, 0, fileSize)与提交/刷盘逻辑可参考仓库文档 MappedFile 内存映射文件详解ConsumeQueue 文件同样基于这一套 MappedFile/MappedFileQueue 机制管理。3.2 第二步文件内二分查找加速检索定位到物理文件后下一步是在文件内部用二分法逼近目标时间戳对应的单元。先计算最低查找偏移量int low minLogicOffset mappedFile.getFileFromOffset() ? (int) (minLogicOffset - mappedFile.getFileFromOffset()) : 0;如果消息队列的最小逻辑偏移量minLogicOffset大于文件起始偏移量则low取两者差值否则取 0——这保证二分查找不会回溯到已被删除的过期索引。再计算中间偏移量利用CQ_STORE_UNIT_SIZE把偏移量对齐到单元边界midOffset (low high) / (2 * CQ_STORE_UNIT_SIZE) * CQ_STORE_UNIT_SIZE;即先按单元数取中值再乘回单元大小保证midOffset永远是 20 的整数倍落在某个完整单元上。随后读取该单元的物理偏移量与消息大小与当前最小物理偏移量minPhysicOffset比较byteBuffer.position(midOffset); long phyOffset byteBuffer.getLong(); int size byteBuffer.getInt(); if (phyOffset minPhysicOffset) { low midOffset CQ_STORE_UNIT_SIZE; leftOffset midOffset; continue; }若phyOffset minPhysicOffset说明该索引指向的物理消息已被删除物理偏移量太旧属于无效区域待查找消息的物理偏移量必然更大因此将low上移一个单元继续二分。若phyOffset不小于最小物理偏移量说明该单元是有效信息则根据消息物理偏移量和消息长度从 CommitLog 中反查该消息的存储时间戳long storeTime this.defaultMessageStore.getCommitLog().pickupStoreTimestamp(phyOffset, size);拿到storeTime后按以下四分支收敛if (storeTime 0) { return 0; } else if (storeTime timestamp) { targetOffset midOffset; break; } else if (storeTime timestamp) { high midOffset - CQ_STORE_UNIT_SIZE; rightOffset midOffset; rightIndexValue storeTime; } else { low midOffset CQ_STORE_UNIT_SIZE; leftOffset midOffset; leftIndexValue storeTime; }storeTime 0消息无效直接返回 0storeTime timestamp恰好命中设置targetOffset并跳出循环storeTime timestamp说明待查找消息的物理偏移量小于midOffset将high收缩到midOffset - CQ_STORE_UNIT_SIZE并记录右侧边界值storeTime timestamp说明待查找消息的物理偏移量大于midOffset将low上移到midOffset CQ_STORE_UNIT_SIZE并记录左侧边界值。循环结束后依据收敛结果决定返回哪个偏移量if (targetOffset ! -1) { offset targetOffset; } else { if (leftIndexValue -1) { offset rightOffset; } else if (rightIndexValue -1) { offset leftOffset; } else { offset Math.abs(timestamp - leftIndexValue) Math.abs(timestamp - rightIndexValue) ? rightOffset : leftOffset; } }targetOffset ! -1找到了存储时间戳恰好等于待查找时间戳的消息leftIndexValue -1没有比目标时间更早的消息返回大于且最接近目标时间戳的偏移量rightOffsetrightIndexValue -1没有比目标时间更晚的消息返回小于且最接近目标时间戳的偏移量leftOffset两侧都存在时取时间距离更近的一侧Math.abs比较。这段逻辑在按时间回放消息从某时间点开始消费等 RocketMQ 管理/查询功能中是核心路径。四、根据当前偏移量获取下一个文件的偏移量rollNextFile消费逻辑偏移量是跨文件连续递增的每个 ConsumeQueue 文件默认固定大小文件以起始逻辑偏移量命名。当当前文件写满或消费越过文件边界时需要滚动到下一个文件对应方法org.apache.rocketmq.store.ConsumeQueue#rollNextFilepublic long rollNextFile(final long index) { int mappedFileSize this.mappedFileSize; int totalUnitsInFile mappedFileSize / CQ_STORE_UNIT_SIZE; return index totalUnitsInFile - index % totalUnitsInFile; }推导逻辑totalUnitsInFile mappedFileSize / CQ_STORE_UNIT_SIZE一个文件最多容纳多少个 20 字节索引单元index % totalUnitsInFile当前偏移量在本文件内的单元序号index totalUnitsInFile - index % totalUnitsInFile直接跳到下一个文件的起始偏移量。该方法的语义与 CommitLog 侧的rollNextFile见 CommitLog 详解遥相呼应一个按物理偏移量滚动文件一个按逻辑偏移量滚动文件共同保证偏移量 → 文件的换算只依赖文件大小无需遍历。五、ConsumeQueue 添加消息putMessagePositionInfoBroker 写入 CommitLog 成功后会通过转发Dispatch机制把消息的索引信息追加到对应的 ConsumeQueue。核心方法是org.apache.rocketmq.store.ConsumeQueue#putMessagePositionInfoprivate boolean putMessagePositionInfo(final long offset, final int size, final long tagsCode, final long cqOffset) { if (offset size this.maxPhysicOffset) { log.warn(Maybe try to build consume queue repeatedly maxPhysicOffset{} phyOffset{}, maxPhysicOffset, offset); return true; } this.byteBufferIndex.flip(); this.byteBufferIndex.limit(CQ_STORE_UNIT_SIZE); this.byteBufferIndex.putLong(offset); this.byteBufferIndex.putInt(size); this.byteBufferIndex.putLong(tagsCode); final long expectLogicOffset cqOffset * CQ_STORE_UNIT_SIZE; MappedFile mappedFile this.mappedFileQueue.getLastMappedFile(expectLogicOffset); if (mappedFile ! null) { if (mappedFile.isFirstCreateInQueue() cqOffset ! 0 mappedFile.getWrotePosition() 0) { this.minLogicOffset expectLogicOffset; this.mappedFileQueue.setFlushedWhere(expectLogicOffset); this.mappedFileQueue.setCommittedWhere(expectLogicOffset); this.fillPreBlank(mappedFile, expectLogicOffset); log.info(fill pre blank space mappedFile.getFileName() expectLogicOffset mappedFile.getWrotePosition()); } if (cqOffset ! 0) { long currentLogicOffset mappedFile.getWrotePosition() mappedFile.getFileFromOffset(); if (expectLogicOffset currentLogicOffset) { log.warn(Build consume queue repeatedly, expectLogicOffset: {} currentLogicOffset: {} Topic: {} QID: {} Diff: {}, expectLogicOffset, currentLogicOffset, this.topic, this.queueId, expectLogicOffset - currentLogicOffset); return true; } if (expectLogicOffset ! currentLogicOffset) { LOG_ERROR.warn([BUG]logic queue order maybe wrong, expectLogicOffset: {} currentLogicOffset: {} Topic: {} QID: {} Diff: {}, expectLogicOffset, currentLogicOffset, this.topic, this.queueId, expectLogicOffset - currentLogicOffset ); } } this.maxPhysicOffset offset size; return mappedFile.appendMessage(this.byteBufferIndex.array()); } return false; }关键步骤拆解幂等保护若offset size maxPhysicOffset说明这是一条已被索引过的旧消息重复转发直接返回 true避免重复构建组装 20 字节索引单元ByteBuffer依次写入 8 字节offset、4 字节size、8 字节tagsCode正好填满CQ_STORE_UNIT_SIZE计算期望逻辑偏移expectLogicOffset cqOffset * CQ_STORE_UNIT_SIZE即逻辑序号 × 单元大小预填充空白当新建文件且cqOffset ! 0时说明该队列之前已存在更早的数据新文件起始位置之前可能有空洞例如跳跃式写入此时调用fillPreBlank填充空白区域同时推进minLogicOffset、flushedWhere、committedWhere保证文件内偏移连续、可被正确换算顺序性校验cqOffset ! 0时检查expectLogicOffset与当前写位置是否一致——小于说明重复构建告警返回不一致则打[BUG]日志仅告警不阻断用于暴露潜在的逻辑队列乱序问题写盘更新maxPhysicOffset offset size后把索引单元追加到 MappedFileappendMessage。值得说明的是putMessagePositionInfo的调用时机与事务消息有关在 消息发送存储流程 中可以看到CommitLog 写入时会根据事务状态处理queueOffset——TRANSACTION_PREPARED_TYPE预提交与TRANSACTION_ROLLBACK_TYPE回滚的消息queueOffset 0不会进入消费队列只有TRANSACTION_NOT_TYPE/TRANSACTION_COMMIT_TYPE普通消息与事务提交才会更新topicQueueTable并继续更新 ConsumeQueue 信息从源头保证了未提交事务消息不可见的语义。六、ConsumeQueue 文件删除destroy当 ConsumeQueue 数据过期被清理时调用org.apache.rocketmq.store.ConsumeQueue#destroypublic void destroy() { this.maxPhysicOffset -1; this.minLogicOffset 0; this.mappedFileQueue.destroy(); if (isExtReadEnable()) { this.consumeQueueExt.destroy(); } }流程要点重置偏移状态maxPhysicOffset -1、minLogicOffset 0让队列回到空的初始语义删除文件调用MappedFileQueue#destroy()将 ConsumeQueue 目录下的文件全部删除扩展属性清理若开启了 ConsumeQueue 扩展存储isExtReadEnable()用于存放 tag 哈希之外的过滤数据同步销毁consumeQueueExt。MappedFileQueue#destroy最终会落到单个 MappedFile 的销毁上先shutdown标记不可用并释放引用等引用计数归零且清理完成后关闭fileChannel并删除磁盘文件——这一完整过程在 MappedFile 内存映射文件详解 中有逐行说明与 CommitLog 文件的销毁逻辑完全一致。七、ConsumeQueue 在存储与消费链路中的位置把 ConsumeQueue 放回 RocketMQ 的整体存储链路中可以更清晰地理解它的价值写入侧消息先顺序追加到 CommitLog物理存储再异步/同步转发索引到对应主题、队列的 ConsumeQueue逻辑索引。CommitLog 是事实数据ConsumeQueue 是可重建的派生物。读取侧Broker 处理拉取请求时见 broker 处理拉取消息请求流程PullMessageProcessor会调用messageStore.getMessage(...)通过消费组、主题、队列、逻辑偏移量在 ConsumeQueue 中定位索引单元拿到物理偏移量与消息大小后再经 CommitLog#getMessage 按偏移量精确定位并返回消息内容。整个过程中ConsumeQueue 承担了逻辑偏移 → 物理偏移的翻译职责。恢复侧Broker 启动恢复时CommitLog 的 recoverNormally 会比对maxPhyOffsetOfConsumeQueue与processOffset若逻辑索引超前于物理数据则调用truncateDirtyLogicFiles(processOffset)截断 ConsumeQueue 中多余的脏索引保证逻辑索引与物理数据的一致性——这也解释了为什么 ConsumeQueue 是可重建、可截断的派生数据。值得一提的是由于 ConsumeQueue 只存 20 字节定长索引天然适合定长单元 二分查找 文件滚动这套高效算法而索引可以随时按 CommitLog 重建的特性也使得 ConsumeQueue 损坏后可以通过mqnrt等机制恢复而不会丢失消息本体。八、小结ConsumeQueue 是 RocketMQ 以空间换时间、索引加速消费思想的核心落地能力关键方法核心机制按时间定位偏移量getOffsetInQueueByTime按文件修改时间定位文件 文件内二分查找跨文件滚动rollNextFile基于CQ_STORE_UNIT_SIZE的纯算术换算索引写入putMessagePositionInfo20 字节定长单元、幂等保护、预填充空白文件删除destroy重置偏移 MappedFileQueue 销毁结合仓库中 CommitLog、MappedFile、消息发送存储、消息拉取流程 与 消息消费流程 等文档继续阅读即可贯通发送 → CommitLog 落盘 → ConsumeQueue 建索引 → 拉取 → 消费的完整源码链路从索引层真正理解 RocketMQ 高性能消息存储的底层原理。【免费下载链接】source-code-hunter 从源码层面剖析挖掘互联网行业主流技术的底层实现原理为广大开发者 “提升技术深度” 提供便利。目前开放 Spring 全家桶Mybatis、Netty、Dubbo 框架及 Redis、Tomcat 中间件等项目地址: https://gitcode.com/doocs/source-code-hunter创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考