ARTICLE DETAIL

资讯详情

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

RocketMQ 知识体系

RocketMQ 知识体系 文章目录一、 核心概念与领域模型 (Core Concepts Architecture)1. 基础消息模型2. 核心组件角色二、 存储与底层机制 (Storage Low-Level Mechanisms)1. 核心存储文件结构2. 刷盘与主从复制机制三、 高级特性与高并发机制 (Advanced Features Mechanisms)1. 核心高级消息类型2. 流量控制与高可用设计四、 生产环境常见问题与高阶运维 (Production Issues Operations)1. 消息消费典型问题2. 运维调优与监控构建一个完整的RocketMQ 知识体系(参考类比)(各层联系)可以从底层核心概念、架构设计、高可用与高并发机制、生产环境核心问题以及高级特性五个维度来进行系统化拆解。以下是 Apache RocketMQ 的全景知识图谱一、 核心概念与领域模型 (Core Concepts Architecture)1. 基础消息模型Message (消息): 业务数据的载体包含 Topic、Tag、Key、Body 及自定义属性。Topic (主题): 消息的逻辑分类逻辑上承载一类业务消息。Tag (标签): Topic 的细分类别用于在同一个 Topic 下过滤细粒度消息服务端不解析 Tag由消费者客户端过滤。Key (业务主键): 消息的唯一业务标识便于在控制台或运维时通过 Key 查询消息轨迹。Queue (队列/分区): Topic 的物理分区Kafka 中称为 Partition。一个 Topic 包含多个 QueueRocketMQ 默认一个 Topic 在每个 Broker 上有 4 个读写队列。2. 核心组件角色NameServer (名字服务):轻量级的服务注册与发现中心类似于 ZooKeeper但无状态且节点间不进行数据同步。Broker 定时向所有 NameServer 汇报心跳Producer/Consumer 从 NameServer 获取 Topic 的路由信息。Broker (代理服务器):消息存储、转发、查询的核心组件。负责接收 Producer 发来的消息、持久化消息、响应 Consumer 的拉取请求。分为Master和SlaveMaster 负责读写Slave 只负责读或在同步/异步复制下进行容灾。Producer (生产者): 负责生产消息并发送到 Broker支持同步、异步、单向Oneway三种发送方式。Consumer (消费者):分为PushConsumer服务端推动实际底层也是长轮询拉取和PullConsumer客户端主动拉取。消费模式分为集群消费 (Clustering)负载均衡一条消息只会被同组内一个消费者消费和广播消费 (Broadcasting)同组内每个消费者都能收到全量消息。二、 存储与底层机制 (Storage Low-Level Mechanisms)1. 核心存储文件结构RocketMQ 采用极其独特的混合型存储架构主要由以下三类文件组成CommitLog:消息存储的物理文件默认大小 1GB所有 Topic 的消息顺序写入同一个 CommitLog 中。实现了极高的磁盘写入性能顺序 I/O PageCache。ConsumeQueue (消费队列):逻辑消费队列相当于 CommitLog 的索引文件。记录了消息在 CommitLog 中的物理偏移量CommitLog Offset、消息大小Size和 Tag 的 Hash 值消费者通过它来寻找消息。IndexFile (索引文件):基于 Hash 索引键Message Key的快速检索文件支持通过 Key 或时间范围快速查询 CommitLog 中的消息。2. 刷盘与主从复制机制刷盘机制:同步刷盘 (Sync Flush): 消息写入 PageCache 且成功持久化到磁盘后才向 Producer 返回成功。数据安全性高吞吐量较低。异步刷盘 (Async Flush): 消息写入 PageCache 即可返回成功由后台线程异步刷盘。吞吐量极高机器宕机可能丢失少量未刷盘数据。主从复制机制:同步复制 (Sync Master-Slave): Master 和 Slave 都写成功后才返回成功俗称“双写”。异步复制 (Async Master-Slave): Master 写入成功即返回异步将数据同步给 Slave存在极短的主备延迟。三、 高级特性与高并发机制 (Advanced Features Mechanisms)1. 核心高级消息类型延时消息 / 定时消息:支持固定等级的延时消息如 1s, 5s, 10s… 2h。底层原理RocketMQ 内部会将延时消息临时存储在特定的系统 TopicSCHEDULE_TOPIC_XXXX中通过定时器Timer/TimerWheel到期后再投递到真实 Topic。事务消息 (Transactional Message):用于解决分布式事务最终一致性基于两阶段提交2PC 定时反查机制。流程发送半消息→ \rightarrow→执行本地事务→ \rightarrow→提交/回滚事务→ \rightarrow→若 Broker 未收到明确指令则主动向生产者发起回查 (Check)。顺序消息 (Ordered Message):保证局部顺序如创建订单→ \rightarrow→支付→ \rightarrow→发货。实现方式生产者通过自定义MessageQueueSelector将同一业务 ID如订单号的消息发送到同一个 Queue 中消费者端通过加锁ConsumeMessageConcurrentlyService或ConsumeMessageOrderlyService单线程/加锁消费该 Queue。2. 流量控制与高可用设计消费端限流与重平衡 (Rebalance):当 Consumer 数量变化或 Topic 队列数变化时触发 Rebalance重新分配消费队列。支持消费端限流通过pullThresholdForQueue控制每个队列的最大缓存消息数或字节数。死信队列 (DLQ - Dead Letter Queue):当一条消息消费重试超过最大次数默认 16 次依然失败时RocketMQ 会将其自动投入死信队列%DLQ%ConsumerGroup供人工排查和处理。四、 生产环境常见问题与高阶运维 (Production Issues Operations)1. 消息消费典型问题消息积压 (Message Accumulation):排查: 检查消费端逻辑性能、是否有死循环、是否数据库瓶颈、线程池满。解决: 临时扩容消费者实例、优化消费逻辑、若允许可编写临时程序将积压消息转移到新 Topic 加速消费。重复消费与幂等性保证 (Message Idempotency):原因: 网络闪断、Consumer 宕机重启等原因导致 ACK 失败Broker 会进行消息重投。解决: 消费端必须做幂等设计如利用数据库唯一主键、Redis 分布式锁、业务状态机校验、去重表。消息丢失排查场景:Producer 端未捕获异步发送异常、未处理返回值。Broker 端异步刷盘/异步复制 机器突然断电宕机。Consumer 端自动提交 Offset 模式下业务还没处理完代码报错或宕机。2. 运维调优与监控NameServer 与 Broker 监控指标:Broker 读写 TPS、磁盘使用率达到 85% 触发强制写保护cleanResourceImmediately、PageCache 繁忙程度、主备延时。集群部署模式:多 Master 模式: 简单、无单点但单个 Master 宕机期间该机器上的队列消息无法消费直到恢复。多 Master 多 Slave 异步复制/同步双写: 高可用标准生产部署方案。DLedger 模式 (Raft 选主): 类似 Kafka 的 Controller 或 etcd 机制通过 Raft 协议实现 Broker 自动主备切换减少人工介入。
返回列表