ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

RocketMQ中的broker

RocketMQ中的broker Broker 是 RocketMQ 最核心服务节点承担消息存储、接收发送、投递消息、元数据管理、HA 主从、定时延时、事务消息等能力底层基于 Netty Remoting 对外提供服务。1、消息接收与存储接收生产者发送消息做校验、解析。写入 CommitLog支持同步刷盘 / 异步刷盘。异步构建ConsumeQueue消费队列索引、IndexFile消息 key 索引根据 key 查询消息。支持普通消息、延时消息、事务消息、批量消息。2、消息投递处理消费者请求处理消费者PULL 拉消息请求PushConsumer 底层也是 pull。根据 ConsumeQueue 索引查找消息数据返回给客户端。维护消费偏移量 offsetoffset 可以存储在 Broker默认也可以存储在客户端本地。处理消费结果处理重试消息把失败消息写入重试队列 % RETRY% 消费组名重试耗尽转入死信队列%DLQ%消费组名。3、元数据管理功能Topic 管理创建、更新 Topic维护 Topic 与 MessageQueue 映射。向所有 NameServer定时上报心跳与 Topic 路由元数据30s 一次。维护消费组信息接收消费者心跳保存当前消费组在线消费者列表返回给消费者用于客户端 Rebalance 重平衡。4、主从高可用 HA 能力Master 接收写请求Slave 不接收写默认只负责读。HA 模块主从复制。支持异步复制、同步复制同步双写。Slave 通过HAClient主动 pull 拉取 Master 的 CommitLog 数据做同步。5.x DLedger‑Broker 支持自动选主、主节点故障自动切换。5、特殊消息处理① 延时消息内部维护SCHEDULE_TOPIC_XXXX系统延时 topic后台定时任务扫描到期消息重新投递到业务 Topic。② 事务消息处理半消息Prepared 消息半消息存储在 CommitLog。事务回查定时向生产者发送 check 回调查询本地事务状态Commit / Rollback。根据生产者回执把半消息提交或者回滚。③ 顺序消费队列锁为顺序消费MessageListenerOrderly提供MessageQueue 分布式锁锁续约、锁超时释放防止重平衡期间多实例并发消费同一个队列。6、存储文件管理CommitLog 文件滚动默认 1GConsumeQueue、IndexFile 文件管理。消息文件过期清理按照配置的保留时间删除过期 CommitLog 物理文件。PageCache 友好支持预分配文件。所有 Topic 的消息全部顺序写入同一个 CommitLog 逻辑文件不是每个 Topic 单独建日志文件再配合 ConsumeQueue 做索引。CommitLog 不是单个物理文件是一组固定大小默认1G的物理文件滚动生成逻辑上看成一个巨大文件。核心目的是充分利用磁盘顺序写性能高的特点结构举例CommitLog真实消息体全部消息顺序追加msgA(topic‑order)offset0msgB(topic‑pay)offset200msgC(topic‑order)offset450ConsumeQueue// topic‑order‑0索引文件只存偏移索引0:offset0len2001:offset450lenxxxConsumeQueue// topic‑pay‑0索引文件0:offset200lenxxx写消息只追加写 CommitLog异步构建各个 ConsumeQueue 索引。消费消息读 ConsumeQueue 拿到 CommitLog 物理 offset再根据 offset 去 CommitLog 读真实消息内容。IndexFile文件生产者发送消息时可以设置keybroker存储消息时会为设置了keys的消息在indexFile中新增一条index索引条目便于后续通过key在commitLog中查询具体的消息。例如设置keymsg.setKeys(order_10086)查询sh mqadmin queryMsgByKey -n 127.0.0.1:9876 -t topic -k order_10086存储位置${ROCKETMQ_HOME}/store/index/文件名示例20260827201245000文件名是该 IndexFile 创建时间戳。固定大小400M 固定长度文件内存映射 MappedFile。分为两大部分Header 头部 Hash 槽位 (slot) Hash 索引条目1、Header 头部40 字节记录元信息beginTimestamp该索引文件第一条消息时间endTimestamp该索引文件最后一条消息时间beginPhyOffset对应 CommitLog 起始物理偏移endPhyOffset对应 CommitLog 结束物理偏移hashSlotCount已经使用的 hash 槽数量indexCount已经写入的索引条目数量2、Hash 槽位 slots 500 万个槽每个 4 字节总大小500W × 4B ≈ 20MB每个 slot 存第一个匹配该 hash 的 Index 条目的下标。对message keys做 hash 取模定位到槽位实现 hash 链表。3、Index 条目列表 index item每个 20 字节一条索引 item 固定 20 字节表格字段字节说明keyHash4 字节对消息keys字段做 hash 值phyOffset8 字节CommitLog 物理偏移量最重要timeDiff4 字节消息时间 - IndexFile 的 beginTimestampnextIndex4 字节冲突链表下一个同 hash 的 item 下标hash 冲突拉链一条消息设置 keys就会生成一条 Index 索引条目如果一条消息设置多个 keys用空格隔开会生成多条 Index 索引。通过业务 key 查找 CommitLog 物理偏移完整流程计算 keyHash → 定位 hash 槽 → 遍历 hash 冲突链表 → 拿到 phyOffset → 使用 phyOffset 读取 CommitLog步骤 1对查询的 key 做 hash对传入的 key 字符串做 hash 算法得到keyHash。步骤 2hash 取模定位到 slot 槽位slotIndex keyHash % hashSlotNum hashSlotNum 5000000通过 slotIndex在 slots 数组拿到index item 的下标itemIndex slots[slotIndex]如果itemIndex 0代表这个 hash 槽下没有任何索引直接返回查无数据。如果不为 0拿到第一条 index 条目的位置。步骤 3遍历 hash 冲突链表取出下标itemIndex对应的IndexItem对比 item 的keyHash和查询 key 算出的keyHash✅相等找到候选记录取出里面的phyOffset8 字节CommitLog 物理偏移量记录下来。读取该 item 的nextIndex4 字节它是同 hash 下一条 item 的下标。如果nextIndex !0把nextIndex赋值给itemIndex循环继续取下一条 item直到nextIndex0链表结束。⚠️hash 冲突不同业务 key 算出来 hash 值一样会挂在同一个槽位的链表上必须全部遍历对比 keyHash。注意IndexFile 只存 hash 值不存原始 key 字符串所以会出现 hash 碰撞返回多条上层还需要再过滤真实消息 keys 字段。步骤 4拿着 phyOffset 读取 CommitLog拿到phyOffset就是消息在 CommitLog 文件中的绝对物理偏移。根据 phyOffset 定位到对应的 CommitLog 映射文件直接从该偏移读取消息实体。从 CommitLog 读出完整消息之后还要比对消息实际的 keys 字段防止 hash 碰撞带来的错误结果。
返回列表