ConsumeQueue 构建流程源码剖析(ReputMessageService)
ConsumeQueue 是 CommitLog 的索引文件由 ReputMessageService 这个后台线程负责构建。ReputMessageService 的核心机制Broker 启动时会开启一个线程每毫秒执行一次 doReput() 方法它有一个属性 reputFromOffset记录消息重放的偏移量每次执行时从 reputFromOffset 开始读取 CommitLog 中的消息逐条解析消息构建 ConsumeQueue 和 IndexFileConsumeQueue 条目格式每个 ConsumeQueue 条目固定 20 个字节字段 长度 说明CommitLog 偏移量 8 字节 消息在 CommitLog 中的物理位置消息长度 4 字节 消息的字节数Tag 哈希码 8 字节 Tag 的哈希值用于消息过滤工作流程图有新数据无新数据Broker 启动启动 ReputMessageService 线程每 1ms 执行一次 doReput从 reputFromOffset检查 CommitLog 是否有新数据逐条解析消息构建 ConsumeQueue 条目offset size tagHash构建 IndexFileKey → Offset 映射更新 reputFromOffset短暂休眠Rebalance 机制源码剖析Rebalance重平衡是消费者组内 Queue 分配的动态调整机制。触发入口RebalanceService 是一个线程任务在消费者客户端启动时被调用。广播模式集群模式RebalanceService.rundoRebalance遍历 consumerTable获取每个 MQConsumerInnerrebalanceByTopic执行重平衡消费模式简化版分配所有 Queue 分配给所有 Consumer核心分配逻辑根据分配策略计算AllocateMessageQueueStrategy平均分配 / 环形分配更新 ProcessQueue创建/释放 PullRequest完成核心实现负载均衡的最终执行逻辑在 rebalanceByTopic() 方法中。广播模式和集群模式的实现不同集群模式是核心实现。分配策略策略 实现类 特点平均分配 AllocateMessageQueueAveragely 尽可能均匀分配默认策略环形分配 AllocateMessageQueueAveragelyByCircle 轮流分配类似发牌事务消息源码剖析半消息与回查事务消息是 RocketMQ 最复杂的特性之一其核心是半消息Half Message 和事务回查Transaction Check 机制。半消息的存储半消息发送成功后会进入 RocketMQ 内部的 RMQ_SYS_TRANS_HALF_TOPIC 的 ConsumeQueue 中。这个消息只存储在 CommitLog 中但在 ConsumeQueue 中不可见因此消费者无法消费到它。事务消息的完整流程业务应用BrokerProducer业务应用BrokerProducer进入回查流程loop[每 60 秒回查]alt[COMMIT][ROLLBACK][UNKNOWN]发送半消息写入 RMQ_SYS_TRANS_HALF_TOPICCommitLog 可见ConsumeQueue 不可见半消息发送成功执行本地事务返回事务状态6a. 提交事务7a. 将消息转发到目标 Topic6b. 回滚事务7b. 删除半消息8c. 发起回查请求9c. 检查本地事务状态10c. 返回状态11c. 提交最终事务状态事务回查机制如果 Producer 在发送半消息后未能及时通知 Broker 提交或回滚Broker 会定期默认 60 秒回查 Producer 的事务状态。check 方法会对半消息进行过滤如超过 72 小时的事务消息算作过期只保留符合条件的半消息进行回查。消息重试与死信源码剖析消息重试机制RocketMQ 的消费者在消息处理失败时会根据返回状态码自动触发重试默认情况下消息最多重试 16 次并发消费每次重试的时间间隔呈指数增长如 10s → 30s → 1min → 2min…重试与死信的流转是否是否消息消费消费成功提交 Offset进入重试队列%RETRY%ConsumerGroup延迟重试时间逐次递增重试次数≤ 16 次进入死信队列%DLQ%ConsumerGroup需要人工介入处理并发消费与顺序消费的重试差异并发消费失败的消息进入重试队列不影响其他消息的消费顺序消费一条消息失败会阻塞该 Queue 的后续消息直到重试成功或进入 DLQ死信队列当消息重试次数超过最大重试次数默认 16 次消息会被放入死信主题%DLQ%ConsumerGroup。死信队列中的消息需要人工介入处理。