尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

RocketMQ核心概念详解:队列模型、消费位点与可靠性机制

RocketMQ核心概念详解:队列模型、消费位点与可靠性机制 最近帮朋友排查一个RocketMQ问题发现他对几个核心概念的误解相当典型一直把Topic当成和Kafka的Topic一样的“容器”把消费组当成随机分配结果消息堆积和重复消费反复出现。我忽然觉得RocketMQ的核心概念其实值得单独写一篇把这些基础但关键的机制讲透。这篇会从业务痛点、角色分工、队列模型、消息类型、消费机制、存储与可靠性一路讲下来最后给出选型和部署时容易踩的坑。无论你是刚接触RocketMQ的开发者还是用了很长时间但对MessageQueue、消费位点、刷盘机制等概念依然模糊的运维这篇都能帮你把底层逻辑理顺。1. 从“为什么需要RocketMQ”开始削峰、异步与解耦很多人一上来就钻API结果学完了还说不清到底哪里能用。我建议反过来先看业务痛点。消息队列能解决的绕不开三件事削峰、异步、解耦。1.1 削峰填谷先讲所有消息队列最核心的存在意义给系统装上缓冲。假设一个秒杀系统上游瞬时请求量1万TPS直接打到数据库数据库大概率秒死。引入RocketMQ之后生产端把请求先转成消息发进Broker下游按自己能承受的速度消费。消费者如果每秒能处理2000条消息就在Broker里排队慢慢消化。这是所有中间件削峰的基本形态只不过RocketMQ在吞吐量和堆积能力上做得很好单机写入能到十万级TPS的级别堆积几千万条消息也不会立刻把Broker拖垮。我常打一个比方高峰期餐厅客人不会一股脑涌进后厨而是前台排队叫号后厨按自己的节奏做菜。消息队列就是那个“排队区域”。但这里有个容易被忽视的设计细节削峰不是把峰值消失而是把压力延后、摊平。所以下游系统一定要按自己的消费能力设计不要因为消息中间件能扛就把消费端写得很重否则峰值过去后堆积消息一次性压过来照样会拖垮数据库。1.2 异步化另一个价值是异步。以电商下单为例下单接口后面通常要扣库存、加积分、发短信、通知财务如果全同步执行任何一个环节卡顿用户都要一直等着。把非核心链路改成订阅消息后下单接口只做两件事写订单数据、发一条“订单创建”消息。下游各自消费消息处理。从用户角度看下单响应时间可能从几秒降到几十毫秒。要注意异步不是减少工作量而是把工作量迁移到合适的时机。比如发短信失败可以通过重试机制解决用户不会因为短信晚了几秒而感知异常。这里有一个设计习惯我很推荐消息体要表达“发生了什么”而不是“去做什么”。写成“订单已支付”比写成“调用积分系统”更合理因为新加入的订阅方不会影响已有系统语义也更稳定。1.3 解耦第三个价值是解耦。没有消息队列时订单系统要同步调用库存、积分、短信等多个系统每次新增一个下游系统订单系统就要改代码。引入RocketMQ以后协作模式变成基于事件驱动的订单系统只往Topic发消息不关心谁订阅。积分系统想参与自己写一个消费者订阅这个Topic就行订单系统一行代码都不用改。不过解耦也不是没有成本。Topic的消息体、字段含义、版本演进都要提前设计好。多个团队各自订阅同一份消息消息结构一旦随意调整很可能一个团队改了字段另一个团队没收到通知消费端直接解析失败。所以解耦之后消息规范反而要比接口规范更严格。2. 入门最先接触的五个角色NameServer、Broker、Producer、Consumer与消费位点RocketMQ的角色模型不算复杂但很多人只记住了四个角色把消费位点漏了。实际上消费位点才是各种诡异问题的来源必须放在核心位置来理解。2.1 NameServer一个类似“电话本”的存在NameServer用来保存路由信息。Broker启动之后把自己的地址、Topic路由表上报给NameServerProducer和Consumer启动之后会定期从NameServer拉取路由信息。它不参与消息读写不存储消息数据节点之间也不互相通信。这一点和Kafka依赖的ZooKeeper很不一样RocketMQ的NameServer更轻量设计目标就是简单、无状态。可以把它理解成电话本想给某个Topic发消息先从NameServer查出它在哪几个Broker上然后去连接对应的Broker。因为NameServer无状态所以部署两台甚至更多都没有主从关系客户端会随机轮询请求把某台NameServer挂掉的概率降到很低。有一个实际经验即便NameServer全部挂了已经拿过路由的客户端短时间内仍能继续发消息因为路由信息已在本地缓存。但新Topic的发现、Broker上下线的感知就会失效。生产环境至少部署两台NameServer成本极低收益很明显。2.2 Broker真正干活的节点Broker是真正存消息、发消息的节点。一个Broker进程可以承载多个Topic一个Topic也可以分布在多个Broker上。Broker有主从之分Master负责读写请求Slave负责备份。生产环境通常把Master和Slave放在不同机器防止单机故障。这里经常有新手踩配置的坑Broker在云主机或Docker里启动后如果不设置brokerIP1注册到NameServer的可能是一个容器内网IP客户端通过这个IP根本连不上。正确做法是显式指定brokerIP1为宿主机实际可达的内网或外网IP。这种问题排查起来很隐蔽因为Broker进程本身启动成功日志也正常但客户端就是报“connect to xxxxx failed”。2.3 Producer与ProducerGroupProducer是发送消息的一端。发送方式有三种同步发送、异步发送、单向发送。同步发送会等待Broker返回确认结果可靠性最高异步发送带回调函数吞吐更高适合对延迟敏感但需要感知结果的场景单向发送只管发、不关心结果适合日志这类允许丢失的高频场景。选择哪种方式取决于业务对可靠性的要求。ProducerGroup是同一类生产者的集合。这个概念比ConsumerGroup轻很多主要和事务消息相关。同一ProducerGroup的实例在事务消息回查时会被认为承担同一类业务用于确定从哪里恢复本地事务状态。实际项目中通常一个应用定义一个ProducerGroup即可没必要为一个Topic单独建一个组。2.4 Consumer与ConsumerGroupConsumer负责拉取消息并处理业务逻辑。一对多关系是这样的同一个消费组内一条消息只会被该组内的一个消费者实例消费不同消费组之间互不干扰各自消费同一份消息。所以同一个Topic可以被多个业务方订阅比如“订单Topic”订单消费者组和积分消费者组都在消费但各自保留各自的消费位点。用“订阅报纸”来理解也很直观Topic是报纸ConsumerGroup是订户家庭一份报纸送到一个家庭家里有几个人不重要只看一份不同家庭各自订各自的报纸。所以当你新增了一个下游业务系统只需新建一个ConsumerGroup去订阅老系统完全不需要改动。2.5 消费位点容易被忽略却最容易出问题的角色消费位点的完整含义是某个消费组在某条MessageQueue上已经消费到哪一条消息。RocketMQ的Push模式虽然由客户端主动拉消息但消费位点的管理在Broker端。消费者每消费一条消息都会上报当前位点Broker持久化到本地文件里。如果消费者宕机重启会从上次提交的位点继续消费。很多“消息凭空消失”的假象其实都和位点有关。比如位点提交过早消费逻辑还没执行完就上报一旦进程崩溃消息就不会再重新投递。又比如处理消息时抛异常但没有返回重试状态还继续向后提交位点消息就会真正丢掉。所以理解位点的更新时机比理解API本身更重要。3. Topic与MessageQueue一个Topic里真正干活的是队列很多文档会把Topic和MessageQueue混在一起说但我建议把它们分开Topic是业务上的逻辑分类MessageQueue才是物理存储和并发的基本单位。3.1 Topic是逻辑分类MessageQueue是物理队列Topic帮你把消息按业务分类比如订单Topic、日志Topic、积分Topic。但消息不会凭空悬浮在一个Topic里而是落到这个Topic底下的具体MessageQueue中。可以回想一下我们项目的实际目录结构每个Broker上的Topic都会创建若干个MessageQueue消息写入时按路由规则选择其中一个队列。MessageQueue本身是FIFO结构先进先出。同一个MessageQueue内的消息是有先后顺序的但跨队列之间没有顺序保证。这意味着如果你想让某个业务键的消息严格有序必须让它们进同一个队列。默认情况下生产者发送消息采用轮询策略散到各个队列里顺序自然无法保证。3.2 队列数和并发的关系理解MessageQueue数量的意义比记任何API都重要。消息消费的并行度取决于MessageQueue数量而不是消费者实例数量。假如一个Topic有8个队列哪怕你起了20个消费者实例也最多只有8个实例能消费到消息其余12个实例会一直空闲。这个坑我见过太多次线上消息堆积团队第一反应是加消费者实例结果实例从4个加到20个堆积一点没缓解。原因就是Topic队列数只有4个。所以提前评估队列数量很重要。经验算法是先估算峰值生产TPS再算出单队列的理论消费TPS用生产速率大致定队列数。更简单的做法是让队列数等于消费者实例数的1到2倍给未来扩容留一点余地。但队列数也不是越多越好队列太多会导致单队列消息稀疏消费端频繁拉取却拉不到数据反而增加无谓的网络开销。3.3 读写队列配置RocketMQ允许对Topic单独设置writeQueueNums和readQueueNums。默认两者一致一般不需要修改。什么时候会需要不一致比如做流量迁移或灰度发布时可以先把写队列数调小让新消息只写入部分队列而读队列保持原样保证旧消息还能被消费到。但我要泼一盆冷水除非你非常清楚自己在做什么否则不要随便动读写队列配置。如果写队列和读队列严重不一致可能造成一部分队列里的消息永远不被消费或者部分消费者分配不到队列。而且调整队列数后客户端路由信息不会立即更新需要等一段时间期间可能出现消费分配不均。3.4 和Kafka Partition对比Kafka的Topic下有PartitionRocketMQ的Topic下有MessageQueue本质上都是分区机制都是为了并行和顺序控制。但两者底层存储差异很大Kafka每个Partition对应独立的日志文件目录RocketMQ则让所有MessageQueue共享一个CommitLog文件顺序写。由此带来的直接结果是Kafka在Topic数量很多的时候会产生大量小文件随机写性能下降明显RocketMQ无论Topic有多少写入都是往同一个CommitLog文件后面追加所以它能承载数千个Topic写入性能依然稳定。这一点也是很多团队从Kafka迁到RocketMQ的原因之一尤其是业务Topic数量多、单Topic流量又不太大的场景。4. 消息类型解析普通、顺序、延迟、事务消息各自的坑RocketMQ官方把消息分成几类但很多开发者只是简单知道有这些类型用的时候才发现各类的边界和限制很容易踩坑。我把每类都展开讲一遍。4.1 普通消息普通消息是最基础的消息类型没有顺序要求没有延迟要求发送后立刻可以被消费。大部分业务场景比如通知类、日志类、数据同步类用普通消息就够了。普通消息也有两个容易被忽视的点一是重复投递是常态消费端必须做幂等二是消费失败需要返回RECONSUME_LATER状态让Broker触发重试不能把异常吞掉。幂等性怎么设计最简单的是利用业务唯一键比如订单号、流水号在消费端先查一遍是否处理过处理过就直接返回成功。也可以用Redis、数据库唯一索引做去重。总之RocketMQ本身只保证消息不丢但做不到恰好一次业务幂等必须自己做。4.2 顺序消息顺序消息分两种全局有序和分区有序。全局有序要求一个Topic只有一个MessageQueue所有消息都进这个队列吞吐量被限制得很低一般只在特别小流量的场景才考虑。分区有序则是在普通多队列状态下通过业务键把同一类消息路由到同一个队列。发送端用MessageQueueSelector来选队列比如按订单号hash取模消费端用MessageListenerOrderly监听。典型场景是订单状态流转订单创建、订单支付、订单完成这三条消息如果散到不同队列消费者看到的顺序就可能错乱。顺序消息的坑在于消费失败后的重试策略。顺序消费者的重试方式和并发消费者不一样它不会简单地把消息重新投递给别的线程而是会挂起队列等待当前消息处理成功后才继续消费后续消息。这样做保证了顺序但会阻塞整个队列的处理。如果一条坏消息一直处理失败后面的消息都会被堵住。所以业务上通常会在消息体内加入状态或版本号消费时判断当前消息是否允许被处理如果因为前一条未完成还不允许处理可以返回稍后重试但要有最大次数保护避免死循环。4.3 延迟消息延迟消息是RocketMQ的一个特色功能。但很多人第一次用就踩坑RocketMQ并不支持任意秒数的延迟而是预置了18个延迟级别1s、5s、10s、30s、1m、2m、3m、4m、5m、6m、7m、8m、9m、10m、20m、30m、1h、2h。发送时通过message.setDelayTimeLevel来指定级别不能直接写“15分钟”。实际项目里“订单30分钟未支付自动关闭”这个场景刚好落在30m级别所以很好用。但如果你需要15分钟预置级别里只有10m和20m可选就比较尴尬。两种处理思路一是消费到延迟消息后判断剩余时间不够就再发一条延迟消息继续等二是使用定时任务扫描数据库兜底。我更推荐后者因为延迟消息的实现机制是把消息临时写入一个内部Topic由定时线程轮询到时间后再放入目标Topic大量延迟消息堆积时Broker内部的时间轮、调度线程压力会明显增加。大促期间尤其要注意。4.4 事务消息事务消息是RocketMQ区别于Kafka、RabbitMQ的最大亮点用来解决“本地业务操作”和“发送消息”两者之间的一致性问题。场景很典型用户下单既要写订单表又要发一条“订单已创建”的消息。如果先写库再发消息发消息失败会导致下游收不到如果先发消息再写库本地事务失败会让下游消费到假消息。事务消息的流程分四步生产者发送一条“半消息”half messageBroker存储该消息但消费者不可见。半消息发送成功后生产者执行本地事务。本地事务执行成功生产者向Broker发送commit执行失败则发送rollback。commit之后消息才对消费端可见。如果步骤3因为网络或进程故障超时Broker会回调生产者实现的checkLocalTransaction接口生产者根据本地事务的实际状态决定返回commit还是rollback。这里有两个坑建议重点注意。第一checkLocalTransaction接口里不能只查内存状态因为Broker可能在半消息发送后等很久才回查进程可能已经重启内存状态早丢了。必须从数据库或者业务表里反查事务是否成功。第二事务消息不能使用发送并忘掉sendOneway的方式因为拿不到发送结果后续流程根本没法走。整个机制说到底就是通过“半消息回查”保证要么业务成功了消息也可见要么业务失败消息被丢弃不存在中间状态。5. 消费模型与消息堆积Push/Pull、消费组、重试、死信消费模型直接关系到消息堆积、重复消费和位置管理。很多问题排查到最后都会归到这一章。5.1 Push模式本质上是长轮询RocketMQ默认提供的DefaultMQPushConsumer名字里虽然有Push但底层不是Broker主动推消息而是客户端持续向Broker拉取。客户端内部有RebalanceService和PullMessageService线程负责分配队列和拉取消息。Broker端做了长轮询优化当队列暂时没有新消息时请求不会立刻结束而是挂起一段时间等新消息到了再返回。这样一来从使用效果上看就很像“推”。理解这一点对排错很有帮助。比如消费者实例多了以后会触发重平衡Rebalance重平衡期间部分队列可能暂停消费消息堆积短时上涨这未必是消费速度不行也可能是重平衡本身带来的抖动。如果你需要完全控制拉取速度和位点可以用DefaultLitePullConsumer主动pull这在处理特定数据源导入、批量处理场景时更灵活。5.2 消费组与广播模式消费组决定了消息的分发方式。默认是集群消费同一个消费组内的一条消息只会被一个实例消费。广播模式则相反组内每个实例都会消费到全量消息。广播模式适合每个节点都需要拿到同一份数据的场景比如应用配置同步、本地缓存预热。但广播模式有几个明显的风险。第一位点保存在本地进程文件里换一台机器消费位点就丢了没法从Broker恢复第二实例扩容时新实例会从最新位点开始消费大概率丢失历史数据。所以核心业务链路不建议使用广播模式。我看到很多人把广播模式用在实时推荐、规则同步这类场景虽然可行但一定要清楚这些限制。5.3 消费位点与堆积排查堆积的本质是生产速率大于消费速率或者消费卡住不动了。实际上排查堆积有相对固定的套路我总结过一遍打开Dashboard的Consumer页面观察diffTotal是否持续上涨。短时上涨可能是重平衡持续上涨才是真堆积。查看该消费组的消费者实例数。如果实例数等于0说明应用挂了这是最常见的“堆积”原因。查看消费日志里有没有大量失败重试的异常。很多堆积不是消费不过来而是一直消费失败。检查消费方法有没有慢调用比如RPC超时、查数据库慢、锁等待等。这通常是单条消息处理时间过长导致的。如果消费逻辑没问题再看Topic队列数是否够。队列数就是并发上限实例再多也突破不了这个上限。这里还要提一个常见误判Dashboard显示的diffTotal包含当前未消费、未拉取以及重试队列里的消息不一定代表所有消息都处理不动。结合消费速率和队列长度才能定性。5.4 重试队列和死信队列消费失败的默认处理方式是返回RECONSUME_LATER消息会进入一个重试主题然后在后续时间点重新投递给消费者。默认重试次数上限是16次重试完还是失败消息就转入死信队列。死信队列的命名规则是%DLQ%消费组名里面全是人工处理的消息。死信队列怎么处理一般可以直接订阅%DLQ%前缀的Topic把捞出来的消息做补偿也可以在Dashboard里浏览死信消息确认业务数据是否需要手工恢复。我见过不少团队完全不关注死信队列直到对账才发现丢了一大片。所以建议生产环境把死信队列的消息数量、积压情况纳入告警。还有一个实战建议如果某个消费逻辑依赖外部系统而外部系统正在故障不加判断地连续重试16次意义不大反而会把消息全部打到死信。更合理的做法是设置合理的重试次数同时把失败消息写入一张本地重试表由定时任务按业务策略补偿避免全部涌入死信。6. 存储与可靠性CommitLog、ConsumeQueue、刷盘、主从RocketMQ为什么快为什么能扛住万亿级消息这些答案都在存储设计里。理解存储才算真正理解RocketMQ的核心概念。6.1 CommitLog高性能的基石CommitLog是一个物理文件所有Topic的所有MessageQueue写操作都追加到这个文件的末尾。顺序写磁盘的吞吐量远高于随机写这就是RocketMQ单机能到十万级TPS的原因之一。一个CommitLog文件默认1GB写满之后自动创建新文件文件名用起始偏移量命名。这里可以类比写日记一次只在一个本子上顺序往下写而不是在很多本子之间来回跳。RocketMQ这种“所有队列共享一个文件”的设计让它即使Topic数量很多也不会出现严重的小文件随机写。这也是它和Kafka在存储模型上最大的差别。6.2 ConsumeQueue逻辑队列的索引ConsumeQueue是每个MessageQueue对应的逻辑队列索引。消息真实内容在CommitLog里ConsumeQueue里只存定位信息CommitLog偏移量、消息大小、Tag哈希码每条记录固定20字节。消费时先读ConsumeQueue找到偏移量再去CommitLog里读取真实消息。可以把它理解成书的目录目录很薄告诉你内容在第几页但真正的正文在正文区。因为消息写入CommitLog后是异步构建ConsumeQueue索引的所以极端情况下刚写入的消息可能Consumer查询索引时还没及时构建好但通常这个过程非常快几十毫秒内完成实际使用中很少感知到。6.3 IndexFile按Key查消息IndexFile是按消息Key建立的索引文件用于根据业务编号快速查询消息。比如客服反馈一笔订单没有消息记录我们可以根据订单号作为Key通过工具或Dashboard查询这条消息当时的存储位置、发送结果、消费状态。IndexFile的定位和ConsumeQueue不同ConsumeQueue是消费路径上的索引每个消费者都会用到IndexFile主要是后台运维和管理场景用数据量太大时不建议频繁使用按Key扫消息会影响Broker性能。6.4 刷盘策略刷盘策略有两种异步刷盘和同步刷盘。异步刷盘是默认配置消息写入操作系统的PageCache就算发送成功由操作系统决定何时落盘。这种模式性能高但机器突然断电时可能会丢少量消息。同步刷盘则要求消息必须真正写入物理磁盘后才返回成功可靠性高但吞吐量会有所下降。我的建议是涉及资金、交易、对账等关键业务至少开同步刷盘日志、通知等非关键业务用异步刷盘就足够。另外刷盘策略和主从复制是两套独立机制不要混在一起配置。一个常见误区是以为只要主从异步复制开着就万事大吉其实主从复制只解决Broker宕机后从副本恢复的问题不能解决断电时操作系统页缓存未落盘导致的消息丢失。6.5 主从同步与高可用Broker主从部署时Master负责读写请求Slave负责备份。同步复制模式下Master写入消息后要等Slave确认才向Producer返回成功数据安全性最高但发送耗时增加异步复制模式则Master写完就返回Slave异步追赶可能主节点故障时丢掉少量最新消息。这里有一个容易被团队忽略的点RocketMQ原生的主从模式并不提供自动故障转移。Master挂掉之后消费者可以切换到Slave继续读但生产者无法写入新消息整个服务的能力会退化。如果希望实现自动选主需要启用Dledger模式通过Raft协议在多个副本之间自动选主这是4.5版本之后支持的部署方式。很多业务团队以为配一个Master一个Slave就自动高可用了直到Master挂掉才发现写不了消息这点必须提前评估。6.6 消息不丢的完整链路消息不丢是三个环节共同配合的结果单靠Broker刷盘并不能解决全部问题。生产者端用同步发送或异步发送回调发送失败要重试。Broker端采用同步刷盘重要消息配合同步复制。消费者端业务逻辑处理完成后再返回消费成功并提交位点。不要在业务没处理完就提交也不要用try-catch吞掉异常后假装成功。除此之外因为RocketMQ在整个投递链路里至少会投递一次消息消费端幂等处理是不可省的一步。现实中很多丢消息事件都不是中间件丢的而是消费者端位点提交过早、异常被吞、广播模式换机器等各种使用层面的问题。理解了整条链路才能定位到具体是哪个环节出了问题。7. 选型对比与部署实践Kafka、RabbitMQ、RocketMQ怎么选以及安装、Dashboard、监控最后把选型和部署运维常见的坑集中说一遍。这些内容来自我接触过的真实项目也覆盖了大家搜索时最常遇到的问题。7.1 三种MQ选型对比维度RocketMQKafkaRabbitMQ开发语言JavaScala/JavaErlang吞吐量很高单机十万级极高单机百万级中等万级延迟低较低极低微秒级消息可靠性高支持事务消息、同步刷盘高但事务能力弱一些高AMQP协议经典可靠消息顺序队列内顺序分区内顺序单队列内顺序延迟消息内置18个级别原生不支持需自研通过TTL死信队列模拟限制较多事务消息内置支持机制成熟一般不用支持但配置复杂生态与社区国内互联网公司使用广泛中文资料多大数据生态强流计算标配中小规模系统、传统企业多典型场景业务系统解耦、削峰、事务一致性、延迟任务日志采集、大数据管道、流处理复杂路由、低延迟请求响应、经典企业集成如果团队是Java技术栈需要事务消息、延迟消息同时希望有可靠消费和顺序消息能力RocketMQ是很容易落地的选择。如果核心是大数据管道、实时计算、海量日志Kafka更合适。如果是复杂路由、需要极低延迟但吞吐量要求不高RabbitMQ仍然是好选择。我们当时的项目选RocketMQ原因很直接团队全是Java业务里需要事务消息来保证交易和积分系统的最终一致还需要现成的延迟消息做超时关闭功能Kafka和RabbitMQ都不能这么顺滑地满足。7.2 Windows安装与Linux部署的几个坑很多入门者都在Windows电脑上装RocketMQ。流程并不复杂下载release包、设置JAVA_HOME、进入bin目录双击启动mqnamesrv.cmd再启动mqbroker.cmd。默认端口是9876NameServer和10911Broker。启动成功后本地Java程序配置namesrvAddr为127.0.0.1:9876即可测试。Windows跑RocketMQ没问题但生产环境基本不会这么做只适合本地学习和调试。真正部署到Linux服务器时有几个坑很常见默认JVM参数很大小规格服务器启动直接内存不足。需要修改runserver.sh和runbroker.sh里的JAVA_OPT把堆内存降到合适范围。云服务器必须设置brokerIP1为对外可达IP否则注册到NameServer的是内网IP客户端连接失败。安全组和防火墙要放行9876、10911等端口很多“启动成功但连不上”的问题其实都是端口没放行。宝塔面板之类的傻瓜面板里部署时更要把JVM调小宝塔自带Java环境版本可能比较乱最好独立安装一个明确的JDK版本再部署。7.3 Dashboard与可视化工具排查消息堆积、查询消息轨迹最常用的工具是rocketmq-dashboard。它是一个Spring Boot应用下载后配置namesrvAddr启动就能在浏览器里看到Topic、Consumer、Message、Producer等页面。Dashboard能直接看到每个消费组的堆积数量也能通过按时间或Key查询消息内容非常方便。不过要明确一点Dashboard只是观察工具不能解决堆积。很多人开了Dashboard看到红色数字就慌了然后去调各种消费端参数其实应该先按照第5章的排查链路去分析。另外Dashboard本身会访问NameServer和Broker不要把Dashboard部署在核心生产节点上避免占用资源影响线上服务。7.4 监控接入与日志消费生产环境建议把RocketMQ接入Prometheus监控。常用方案是rocketmq-exporter它从Broker和NameServer采集指标暴露metrics给Prometheus再由Grafana做可视化。重点监控指标包括消息堆积量、重试队列大小、Broker是否存活、存储磁盘使用率、CommitLog文件占用趋势。日志场景中Logstash输出到RocketMQ也是常见需求Logstash作为生产者把日志写入MQ日志分析系统作为消费者订阅处理。需要注意插件版本和RocketMQ服务端版本兼容性日志流量和业务消息最好分离部署避免日志量大时影响业务消息的写入延迟。我个人在实际项目中的体会是RocketMQ的核心概念并不算多但它们之间是连锁关系。没理解队列数和消费位点堆积和重复消费就会反复出现没理解刷盘和主从就对“会不会丢消息”没有感知没理解事务消息的半消息和回查机制又容易把本地事务和发消息做成两段简单拼凑。很多人用不好RocketMQ很多时候不是功能不够而是这些基础概念没有真正串起来。先把这一层理顺后面的使用和扩展基本就是顺水推舟的事。
返回列表