资讯动态

Apache RocketMQ 设计原理深度解析:消息存储、通信机制、消息过滤、负载均衡、事务消息与消息查询

发布时间:2026/9/20 19:03:40 来源:尧图企业网站定制
Apache RocketMQ 设计原理深度解析消息存储、通信机制、消息过滤、负载均衡、事务消息与消息查询【免费下载链接】rocketmqApache RocketMQ is a cloud native messaging and streaming platform, making it simple to build event-driven applications.项目地址: https://gitcode.com/gh_mirrors/ro/rocketmqApache RocketMQ 作为云原生消息与流处理平台其核心设计文档系统性地阐述了六大关键技术领域消息存储架构、Remoting 通信机制、消息过滤、负载均衡、事务消息与消息查询。本篇文章以仓库内 docs/cn/design.md 为骨架结合store、remoting、broker、client等模块的真实源码实现逐层拆解 RocketMQ 的底层工作原理帮助读者建立从「文件落盘」到「客户端负载均衡」再到「事务最终一致」的完整认知并掌握每项设计背后的工程权衡与可配置参数。1 消息存储消息存储是 RocketMQ 中最为复杂和最为重要的一部分。本节分别从消息存储整体架构、PageCache 与 Mmap 内存映射以及两种不同的刷盘方式三个方面展开叙述。1.1 消息存储整体架构RocketMQ 采用混合型存储结构Broker 单个实例下所有队列共用一个日志数据文件CommitLog来存储消息实体。这种设计针对 Producer 和 Consumer 分别采用了「数据与索引分离」的存储结构整体架构主要由以下三个文件构成。(1) CommitLog消息主体与元数据的存储主体CommitLog 存储 Producer 端写入的消息主体内容消息内容不定长。单个文件大小默认 1G文件名为 20 位、左边补零剩余位为起始偏移量第一个文件00000000000000000000起始偏移量为 0文件大小1G 1073741824字节当第一个文件写满第二个文件为00000000001073741824起始偏移量为1073741824以此类推消息主要是顺序写入日志文件文件满了则写入下一个文件。在源码中CommitLog 的存储单位是 MappedFile多个 MappedFile 由MappedFileQueue统一管理见 store/src/main/java/org/apache/rocketmq/store/MappedFileQueue.java 与 store/src/main/java/org/apache/rocketmq/store/CommitLog.java。消息写入的物理落盘行为由PutMessageLock默认PutMessageReentrantLock也提供PutMessageSpinLock自旋锁实现保护二者都位于store/src/main/java/org/apache/rocketmq/store/下用于在单机多线程并发写入时保证顺序追加。(2) ConsumeQueue消息消费索引ConsumeQueue 的引入是为了提高消息消费性能。由于 RocketMQ 是基于 Topic 的订阅模式消息消费针对主题进行若直接遍历 CommitLog 按 Topic 检索消息会非常低效。Consumer 可根据 ConsumeQueue 查找待消费的消息ConsumeQueue 保存指定 Topic 下队列消息在 CommitLog 中的起始物理偏移量 offset8 字节、消息大小 size4 字节、消息 Tag 的 HashCode8 字节每个条目共 20 字节采用topic/queue/file三层组织结构存储路径为$HOME/store/consumequeue/{topic}/{queueId}/{fileName}单个文件由 30W 个条目组成每个文件约 5.72M可以像数组一样随机访问每一个条目。源码中 ConsumeQueue 由 store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java 实现其中的CQ_STORE_UNIT_SIZE20 字节等常量定义了条目定长结构与文档描述完全一致。(3) IndexFileKey 与时间区间查询索引IndexFile 提供通过 Key 或时间区间查询消息的能力存储位置为$HOME/store/index/{fileName}文件名以创建时的时间戳命名单个 IndexFile 大小固定约 400M可保存 2000W 个索引底层在文件系统中实现 HashMap 结构因此 RocketMQ 的索引文件本质是hash 索引。IndexFile 的实现位于 store/src/main/java/org/apache/rocketmq/store/index/IndexFile.java其结构与配套测试可参考 store/src/test/java/org/apache/rocketmq/store/index/IndexFileTest.java。ReputMessageService 异步构建索引RocketMQ 的整体写入链路为Producer 发送消息至 BrokerBroker 使用同步或异步方式刷盘持久化到 CommitLog。只要消息被刷盘持久化Producer 发送的消息就不会丢失Consumer 也就肯定有机会消费。RocketMQ 使用 Broker 端的后台服务线程ReputMessageService不停地分发请求异步构建 ConsumeQueue逻辑消费队列和 IndexFile 数据。从源码看ReputMessageService定义在 store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java 中它继承ServiceThread持续从 CommitLog 拉取新写入的数据并解析为DispatchRequest随后调用doDispatch(dispatchRequest)将请求分发给 ConsumeQueue 存储与 IndexServicedoDispatch与putMessagePositionInfo均位于同一文件的 L2064-L2076。也就是说消息主数据与索引数据是异步解耦的先落 CommitLog再由后台线程回放构建二级索引。当消费端无法拉取到消息时可以等待下一次拉取同时服务端支持长轮询模式如果消息拉取请求未拉到消息Broker 允许等待 30s只要这段时间内有新消息到达就直接返回给消费端。1.2 页缓存与内存映射页缓存PageCache是 OS 对文件的缓存用于加速文件读写。程序对文件进行顺序读写的速度几乎接近内存读写速度主要就是 OS 使用 PageCache 机制做了性能优化将一部分内存用作 PageCache数据写入时OS 先写入 Cache随后由pdflush内核线程异步刷盘至物理磁盘数据读取时若一次读取未命中 PageCacheOS 从物理磁盘访问文件的同时会顺序对其他相邻块数据文件进行预读取。在 RocketMQ 中ConsumeQueue 逻辑消费队列存储数据较少且顺序读取在 PageCache 预读取作用下其读性能几乎接近读内存即使在消息堆积情况下也不影响性能CommitLog 消息存储日志数据文件读取消息内容时会产生较多随机访问严重影响性能。如果选择合适的系统 IO 调度算法例如块存储采用 SSD 时设置为Deadline调度算法随机读性能也会有所提升。RocketMQ 主要通过MappedByteBuffer对文件进行读写利用 NIO 的 FileChannel 模型将磁盘物理文件直接映射到用户态内存地址中Mmap 方式减少了传统 IO 在内核态缓冲区与用户态缓冲区之间来回拷贝的开销将对文件的操作转化为直接对内存地址进行操作极大提高文件读写效率。正因为需要使用内存映射机制RocketMQ 的文件存储都采用定长结构方便一次将整个文件映射至内存。1.3 消息刷盘RocketMQ 提供两种刷盘方式(1) 同步刷盘只有在消息真正持久化至磁盘后Broker 端才会向 Producer 端返回成功 ACK。同步刷盘对消息可靠性是很好的保障但性能上有较大影响一般适用于金融业务应用。(2) 异步刷盘异步刷盘充分利用 OS 的 PageCache 优势只要消息写入 PageCache 即可将成功 ACK 返回给 Producer 端。消息刷盘由后台异步线程提交降低了读写延迟提高了 MQ 的性能和吞吐量。刷盘策略的开关与参数集中在 Broker 的配置体系BrokerConfig中通过flushDiskType等相关配置项控制实际执行逻辑位于 store/src/main/java/org/apache/rocketmq/store/CommitLog.java 的 GroupCommitService / FlushRealTimeService 等刷盘服务线程中。2 通信机制RocketMQ 消息队列集群主要包括 NameServer、BrokerMaster/Slave、Producer、Consumer 四个角色基本通信流程如下Broker 启动后完成一次将自己注册至 NameServer 的操作随后每隔 30s 定时向 NameServer 上报 Topic 路由信息消息生产者 Producer发送消息时根据消息 Topic 从本地缓存的TopicPublishInfoTable获取路由信息若没有则更新路由信息并从 NameServer 重新拉取Producer 默认每隔 30s 向 NameServer 拉取一次路由信息Producer根据路由信息选择一个队列MessageQueue进行消息发送Broker 作为消息接收者接收消息并落盘存储消息消费者 Consumer根据路由信息在完成客户端负载均衡后选择其中一个或某几个消息队列拉取消息并消费。从上述流程可以看出在消息生产者、Broker 和 NameServer 之间都会发生通信因此网络通信模块的设计直接决定 RocketMQ 集群整体的消息传输能力与最终性能。rocketmq-remoting模块是 RocketMQ 负责网络通信的模块几乎被所有需要网络通信的模块如 rocketmq-client、rocketmq-broker、rocketmq-namesrv所依赖和引用。RocketMQ 自定义了通信协议并在 Netty 基础之上扩展了通信模块。2.1 Remoting 通信类结构Remoting 模块围绕客户端与服务器两端组织类结构客户端侧以NettyRemotingClient为核心remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingClient.java服务器侧以NettyRemotingServer为核心remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyRemotingServer.java编解码器为NettyEncoder/NettyDecoder其中NettyDecoder继承了 Netty 的LengthFieldBasedFrameDecoder以 4 字节长度字段进行拆包见 remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyDecoder.java请求处理器接口为NettyRequestProcessor通过processRequest处理业务请求。2.2 协议设计与编解码Client 与 Server 之间完成一次消息发送需要对消息进行协议约定因此 RocketMQ 自定义了消息协议。在 RocketMQ 中RemotingCommand类封装了消息传输过程中的所有数据内容不但包含所有数据结构还包含编码解码操作remoting/src/main/java/org/apache/rocketmq/remoting/protocol/RemotingCommand.java。RemotingCommand 的 Header 字段如下Header 字段类型Request 说明Response 说明codeint请求操作码应答方根据不同的请求码进行不同的业务处理应答响应码。0 表示成功非 0 则表示各种错误languageLanguageCode请求方实现的语言应答方实现的语言versionint请求方程序的版本应答方程序的版本opaqueint相当于 requestId在同一个连接上的不同请求标识码与响应消息中的相对应应答不做修改直接返回flagint区分是普通 RPC 还是 oneway RPC 的标志区分是普通 RPC 还是 oneway RPC 的标志remarkString传输自定义文本信息传输自定义文本信息extFieldsHashMapString, String请求自定义扩展信息响应自定义扩展信息一次完整的网络传输内容主要分为以下 4 部分消息长度总长度四个字节存储占用一个 int 类型序列化类型 消息头长度同样占用一个 int 类型第一个字节表示序列化类型后面三个字节表示消息头长度消息头数据经过序列化后的消息头数据消息主体数据消息主体的二进制字节数据内容。源码印证RemotingCommand.encodeHeader(int bodyLength)会通过markProtocolType(headerLength, serializeTypeCurrentRPC)将序列化类型标记进第二个 int 字段见 RemotingCommand.java L476-L483而解码端decode方法 L201-L222则先读 4 字节总长度再读 4 字节头部标记随后按 JSON 或 ROCKETMQ 序列化方式解析 Header。当前仓库默认的序列化类型为 JSONSerializeTypeConfigInThisServer SerializeType.JSON并支持通过rocketmq.remoting.serialize.type等协议配置切换。2.3 消息的通信方式和流程RocketMQ 支持同步sync、异步async、单向oneway三种通信方式同步sync发送方阻塞等待应答可靠性最高异步async通过回调处理应答适合高吞吐场景单向oneway只发不收相对简单一般用在发送心跳包场景无需关注 Response。在服务端Netty 的 pipeline 中依次编排了编解码、空闲检测、连接管理等 HandlerNettyRemotingServer收到完整的RemotingCommand后根据业务请求码code从processorTable本地缓存中查找对应的NettyRequestProcessor并执行。2.4 Reactor 多线程设计RocketMQ 的 RPC 通信采用 Netty 组件作为底层通信库遵循 Reactor 多线程模型并在此基础上做了扩展和优化。NettyRemotingServer的 Reactor 多线程模型大致如下一个Reactor 主线程eventLoopGroupBoss负责监听 TCP 网络连接请求建立连接、创建 SocketChannel 并注册到 selector 上RocketMQ 源码会自动根据 OS 类型选择 NIO 与 Epoll也可通过参数配置然后监听真正的网络数据拿到网络数据后丢给Worker 线程池eventLoopGroupSelector源码中默认设置为 3在真正执行业务逻辑之前需要进行 SSL 验证、编解码、空闲检查、网络连接管理这些工作交给defaultEventExecutorGroup源码中默认设置为 8完成处理业务操作放在业务线程池中执行根据 RemotingCommand 的业务请求码code到processorTable本地缓存变量中找到对应 processor封装成 task 后提交给对应的业务 processor 处理线程池执行如发送消息的sendMessageExecutor。从入口到业务逻辑的几步中线程池一直在增加这与每一步逻辑复杂性相关越复杂需要的并发通道越宽。线程组成如下表线程数线程名线程具体说明1NettyBoss_%dReactor 主线程NNettyServerEPOLLSelector_%d_%dReactor 线程池M1NettyServerCodecThread_%dWorker 线程池M2RemotingExecutorThread_%d业务 processor 处理线程池源码中的默认值可以从 remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyServerConfig.java 确认serverSelectorThreads 3Reactor 线程数 N、serverWorkerThreads 8Worker 线程数 M1对应NettyServerCodecThread、serverCallbackExecutorCount等均可配置。eventLoopGroupBoss与eventLoopGroupSelector的构建逻辑位于 NettyRemotingServer 的buildEventLoopGroupBoss()/buildEventLoopGroupSelector()见 NettyRemotingServer.java L145-L146。3 消息过滤RocketMQ 的消息过滤方式有别于其它 MQ 中间件是在Consumer 端订阅消息时再做过滤的。这源于其「Producer 写入消息与 Consumer 订阅消息分离存储」的机制Consumer 端订阅消息需通过 ConsumeQueue 逻辑队列拿到索引再从 CommitLog 读取真正的消息实体内容。ConsumeQueue 存储结构中有 8 字节存储 Message Tag 的哈希值基于 Tag 的消息过滤正是基于这个字段值。主要支持两种过滤方式(1) Tag 过滤方式Consumer 端在订阅消息时除了指定 Topic 还可以指定 TAG如果一个消息有多个 TAG可以用||分隔。流程如下Consumer 端将订阅请求构建成一个SubscriptionData发送一个 Pull 消息请求给 Broker 端Broker 端从文件存储层 Store 读取数据之前会用这些数据先构建一个MessageFilter传给 StoreStore 从 ConsumeQueue 读取一条记录后用记录中的消息 tag hash 值做过滤。由于服务端只是根据 hashcode 判断无法精确对 tag 原始字符串过滤因此在消息消费端拉取到消息后还需要对消息的原始 tag 字符串进行比对如果不同则丢弃该消息不进行消费。(2) SQL92 的过滤方式这种方式的大致做法与 Tag 过滤一样只是在 Store 层的具体过滤过程不同真正的 SQL expression 的构建和执行由rocketmq-filter 模块负责。每次过滤都去执行 SQL 表达式会影响效率所以 RocketMQ 使用BloomFilter避免每次都去执行。SQL92 的表达式上下文为消息的属性。与过滤相关的实现可以参考 store/src/main/java/org/apache/rocketmq/store/MessageFilter.java、store/src/main/java/org/apache/rocketmq/store/DefaultMessageFilter.javaTag 过滤的默认实现以及filter模块下的 filter/src/main/java/org/apache/rocketmq/filter/SqlFilter.java 与 filter/src/main/java/org/apache/rocketmq/filter/FilterFactory.javaSQL92 表达式构建与执行。4 负载均衡RocketMQ 中的负载均衡都在Client 端完成主要分为 Producer 端发送消息的负载均衡和 Consumer 端订阅消息的负载均衡。4.1 Producer 的负载均衡Producer 端在发送消息时会先根据 Topic 找到指定的TopicPublishInfo。获取路由信息后客户端在默认方式下通过selectOneMessageQueue()方法从TopicPublishInfo中的messageQueueList选择一个队列MessageQueue发送消息。具体的容错策略均在MQFaultStrategy类中定义client/src/main/java/org/apache/rocketmq/client/latency/MQFaultStrategy.java。这里有一个sendLatencyFaultEnable开关变量开启时在随机递增取模的基础上再过滤掉 not available 的 Broker 代理。所谓latencyFaultTolerance是指对之前失败的 Broker 按一定的时间做退避。例如如果上次请求的 latency 超过 550ms就退避 30000ms超过 1000ms就退避 60000ms关闭时采用随机递增取模的方式选择一个队列MessageQueue发送消息。latencyFaultTolerance机制是实现消息发送高可用的核心关键所在。4.2 Consumer 的负载均衡Consumer 端的两种消费模式Push/Pull都是基于拉模式获取消息的。Push 模式只是对 Pull 模式的一种封装其本质实现为消息拉取线程从服务器拉取到一批消息后提交到消息消费线程池然后继续向服务器再次尝试拉取消息如果未拉取到消息则延迟一下继续拉取。在两种基于拉模式的消费方式中均需要 Consumer 端知道从 Broker 端的哪一个消息队列中获取消息。因此有必要在 Consumer 端做负载均衡——即 Broker 端多个 MessageQueue 分配给同一个 ConsumerGroup 中的哪些 Consumer 消费。1. Consumer 端的心跳包发送Consumer 启动后会通过定时任务不断向 RocketMQ 集群中所有 Broker 实例发送心跳包包含消息消费分组名称、订阅关系集合、消息通信模式和客户端 id 等信息。Broker 端收到心跳消息后将其维护在ConsumerManager的本地缓存变量consumerTable中同时将封装后的客户端网络通道信息保存在本地缓存变量channelInfoTable中为之后的负载均衡提供元数据依据。2. Consumer 端负载均衡核心类 —— RebalanceImpl在 Consumer 实例的启动流程中启动 MQClientInstance 实例部分会完成负载均衡服务线程RebalanceService的启动每隔 20s 执行一次。从源码看RebalanceService线程的 run() 方法最终调用的是RebalanceImpl类的rebalanceByTopic()方法client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalanceImpl.java该方法是实现 Consumer 端负载均衡的核心。rebalanceByTopic()会根据消费者通信类型为「广播模式」还是「集群模式」做不同逻辑处理集群模式下的主要处理流程如下从rebalanceImpl实例的本地缓存变量topicSubscribeInfoTable中获取该 Topic 主题下的消息消费队列集合mqSet以 topic 和 consumerGroup 为参数调用mQClientFactory.findConsumerIdList()方法向 Broker 端发送获取该消费组下消费者 Id 列表的 RPC 请求Broker 端基于前面 Consumer 上报的心跳包数据构建的consumerTable做出响应返回业务请求码GET_CONSUMER_LIST_BY_GROUP先对 Topic 下的消息消费队列、消费者 Id 排序然后用消息队列分配策略算法默认为消息队列的平均分配算法计算出待拉取的消息队列。这里的平均分配算法类似于分页算法将所有 MessageQueue 排好序类似于记录将所有 Consumer 排好序类似页数求出每一页需要包含的平均 size 和每页记录的范围 range最后遍历整个 range 计算出当前 Consumer 端应该分配到的记录即 MessageQueue调用updateProcessQueueTableInRebalance()方法将分配到的消息队列集合mqSet与processQueueTable做过滤比对红色部分与分配到的 mqSet 互不包含的队列将这些队列设置 Dropped 属性为 true然后查看这些队列是否可以移除出processQueueTable缓存变量具体执行removeUnnecessaryMessageQueue()方法即每隔 1s 查看是否可以获取当前消费处理队列的锁拿到则返回 true等待 1s 后仍拿不到锁则返回 false。返回 true 时从processQueueTable缓存变量中移除对应 Entry绿色部分与分配到的 mqSet 的交集判断该 ProcessQueue 是否已经过期。Pull 模式不用管如果是 Push 模式设置 Dropped 属性为 true并调用removeUnnecessaryMessageQueue()方法像上面一样尝试移除 Entry最后为过滤后的消息队列集合mqSet中每个 MessageQueue 创建一个ProcessQueue对象并存入RebalanceImpl的processQueueTable队列中其中调用computePullFromWhere(MessageQueue mq)方法获取该 MessageQueue 的下一个进度消费值 offset填充至接下来要创建的 pullRequest 对象属性中并创建拉取请求对象 pullRequest 添加到拉取列表pullRequestList中最后执行dispatchPullRequest()方法将 Pull 消息请求对象依次放入PullMessageService服务线程的阻塞队列pullRequestQueue中client/src/main/java/org/apache/rocketmq/client/impl/consumer/PullMessageService.java待该服务线程取出后向 Broker 端发起 Pull 消息请求。消息消费队列在同一消费组不同消费者之间的负载均衡其核心设计理念是一个消息消费队列在同一时间只允许被同一消费组内的一个消费者消费一个消息消费者可以同时消费多个消息队列。5 事务消息Apache RocketMQ 在 4.3.0 版本中开始支持分布式事务消息采用2PC两阶段提交思想实现提交事务消息并增加一个补偿逻辑处理二阶段超时或失败的消息。5.1 RocketMQ 事务消息流程概要事务消息整体分为两个流程正常事务消息的发送及提交、事务消息的补偿流程。1. 事务消息发送及提交发送消息half 消息服务端响应消息写入结果根据发送结果执行本地事务如果写入失败此时 half 消息对业务不可见本地逻辑不执行根据本地事务状态执行 Commit 或 RollbackCommit 操作生成消息索引消息对消费者可见。2. 补偿流程对没有 Commit/Rollback 的事务消息pending 状态的消息从服务端发起一次「回查」Producer 收到回查消息检查回查消息对应的本地事务状态根据本地事务状态重新 Commit 或 Rollback。其中补偿阶段用于解决消息 Commit 或 Rollback 发生超时或失败的情况。5.2 RocketMQ 事务消息设计1. 事务消息在一阶段对用户不可见事务消息相对普通消息最大的特点是一阶段发送的消息对用户不可见。RocketMQ 的做法是如果消息是 half 消息备份原消息的主题与消息消费队列然后改变主题为RMQ_SYS_TRANS_HALF_TOPIC。由于消费组未订阅该主题消费端无法消费 half 类型的消息。随后 RocketMQ 开启一个定时任务从 Topic 为RMQ_SYS_TRANS_HALF_TOPIC中拉取消息进行消费根据生产者组获取服务提供者发送回查事务状态请求根据事务状态决定提交或回滚消息。具体实现策略是写入的如果是事务消息对消息的 Topic 和 Queue 等属性进行替换同时将原来的 Topic 和 Queue 信息存储到消息的属性中。因为消息主题被替换消息不会转发到原主题的消息消费队列消费者无法感知消息的存在。改变消息主题是 RocketMQ 的常用「套路」——延时消息的实现机制与此类似。在源码中事务相关核心类位于 broker/src/main/java/org/apache/rocketmq/broker/transaction/ 目录包括TransactionalMessageService接口定义 commitMessage / rollbackMessage 等操作见 TransactionalMessageService.javaTransactionalMessageServiceImpl默认实现见 TransactionalMessageServiceImpl.javaTransactionalMessageBridge负责 half 消息写入、Op 消息写入与普通消息恢复见 TransactionalMessageBridge.javaTransactionalMessageUtil提供buildOpTopic()等内部 Topic 构建方法见 TransactionalMessageUtil.java。2. Commit 和 Rollback 操作以及 Op 消息的引入完成一阶段写入一条对用户不可见的消息后二阶段如果是Commit需要让消息对用户可见如果是Rollback则需要撤销一阶段的消息。对于 Rollback一阶段的消息本身对用户不可见其实不需要真正撤销消息实际上 RocketMQ 也无法真正删除一条消息因为消息是顺序写文件的。但为了区别于「消息没有确定状态」Pending 状态事务悬而未决需要一个操作来标识这条消息的最终状态。RocketMQ 引入了Op 消息的概念用 Op 消息标识事务消息已经确定的状态Commit 或 Rollback。如果一条事务消息没有对应的 Op 消息说明这个事务的状态还无法确定可能是二阶段失败了。引入 Op 消息后事务消息无论 Commit 还是 Rollback 都会记录一个 Op 操作Commit 相对于 Rollback只是在写入 Op 消息前创建 Half 消息的索引。3. Op 消息的存储和对应关系RocketMQ 将 Op 消息写入全局一个特定的 Topic 中通过源码中的方法TransactionalMessageUtil.buildOpTopic()构建。这个 Topic 是内部 Topic像 Half 消息的 Topic 一样不会被用户消费。Op 消息的内容为对应的 Half 消息的存储 Offset这样通过 Op 消息能索引到 Half 消息进行后续的回查操作。4. Half 消息的索引构建在执行二阶段 Commit 操作时需要构建 Half 消息的索引。一阶段的 Half 消息写入特殊 Topic因此二阶段构建索引时需要读取出 Half 消息将 Topic 和 Queue 替换成真正的目标 Topic 和 Queue之后通过一次普通消息的写入操作生成一条对用户可见的消息。也就是说事务消息二阶段其实是利用一阶段存储的消息内容在二阶段恢复出一条完整的普通消息然后走一遍消息写入流程。5. 如何处理二阶段失败的消息如果在二阶段过程中失败例如做 Commit 操作时出现网络问题导致 Commit 失败需要通过一定策略使这条消息最终被 Commit。RocketMQ 采用一种补偿机制称为「回查」Broker 端对未确定状态的消息发起回查将消息发送到对应的 Producer 端同一个 Group 的 ProducerProducer 根据消息检查本地事务状态进而执行 Commit 或 RollbackBroker 端通过对比 Half 消息和 Op 消息进行事务消息的回查并且推进 CheckPoint记录那些事务消息的状态是确定的。从源码看回查任务由TransactionalMessageCheckService调度执行其check次数上限读取自 Broker 配置的transactionCheckMax见 broker/src/main/java/org/apache/rocketmq/broker/transaction/TransactionalMessageCheckService.java回查请求由AbstractTransactionalMessageCheckListener通过broker2Client.checkProducerTransactionState(...)发送给 Producer见 AbstractTransactionalMessageCheckListener.java。值得注意的是RocketMQ 并不会无休止地对消息事务状态回查默认回查 15 次如果 15 次回查仍无法得知事务状态RocketMQ 默认回滚该消息。6 消息查询RocketMQ 支持按照两种维度进行消息查询按照 Message Id 查询消息、按照 Message Key 查询消息。6.1 按照 MessageId 查询消息RocketMQ 中的 MessageId 长度总共有 16 字节其中包含消息存储主机地址IP 地址和端口和消息 Commit Log offset。「按照 MessageId 查询消息」的具体做法是Client 端从 MessageId 中解析出 Broker 的地址IP 地址和端口和 Commit Log 的偏移地址封装成一个 RPC 请求后通过 Remoting 通信层发送业务请求码VIEW_MESSAGE_BY_IDBroker 端走QueryMessageProcessorbroker/src/main/java/org/apache/rocketmq/broker/processor/QueryMessageProcessor.java其处理入口见 L78 对RequestCode.VIEW_MESSAGE_BY_ID的分发读取消息的过程是用其中的 commitLog offset 和 size 去 CommitLog 中找到真正的记录解析成一条完整的消息返回。6.2 按照 Message Key 查询消息「按照 Message Key 查询消息」主要基于 RocketMQ 的IndexFile 索引文件实现。IndexFile 的逻辑结构类似 JDK 中 HashMap 的实现具体结构如下IndexFile 文件的存储位置是$HOME/store/index/${fileName}文件名以创建时的时间戳命名文件大小固定等于40 500W*4 2000W*20 420000040字节如果消息的 properties 中设置了UNIQ_KEY属性就用topic # UNIQ_KEY 的 value作为 key 做写入操作如果消息设置了KEYS属性多个 KEY 以空格分隔也会用topic # KEY做索引。其中的索引数据包含Key Hash / CommitLog Offset / Timestamp / NextIndex offset四个字段一共 20 字节NextIndex offset即前面读出来的slotValue如果有 hash 冲突就可以用这个字段将所有冲突的索引用链表的方式串起来Timestamp记录的是消息storeTimestamp之间的差并不是一个绝对的时间40 字节的 Header 用于保存一些总的统计信息4 * 500W的 Slot Table 并不保存真正的索引数据而是保存每个槽位对应的单向链表的头20 * 2000W是真正的索引数据即一个 Index File 可以保存 2000W 个索引。「按照 Message Key 查询消息」主要通过 Broker 端的QueryMessageProcessor业务处理器来查询读取消息的过程就是用 topic 和 key 找到 IndexFile 索引文件中的一条记录根据其中的 commitLog offset 从 CommitLog 文件中读取消息的实体内容。结语从存储层的 CommitLog / ConsumeQueue / IndexFile 三件套到 Remoting 层自研协议与 Reactor 多线程模型再到客户端负载均衡、双阶段事务补偿与双维度消息查询RocketMQ 的设计文档勾勒出一套「顺序写盘 异步建索引 客户端均衡 最终一致」的完整工程体系。本文结合仓库源码store、remoting、broker、client、filter各模块对原设计文档进行了逐项印证与扩充希望读者能够以此为基础进一步阅读 docs/cn/design.md 原文档、深入对应模块源码将设计原理落到实际调优与二次开发中。【免费下载链接】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创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价