RocketMQ核心编程模型解析:消息发送、消费与参数调优实战
RocketMQ这个中间件在国内后端圈子里的存在感确实很强。我见过不少团队从Kafka迁到RocketMQ也有从RabbitMQ换过来的理由绕不开那几样事务消息、延迟消息以及更贴合业务场景的消费模型。但真正上手写代码的时候很多人会被它一套名词绕晕。Producer、Consumer还好理解Topic和Tag也不算难可一旦Queue、Offset、ConsumerGroup、Rebalance这些概念叠加在一起代码就开始跑偏了。这篇文章不打算复读官方文档而是从“核心编程模型”这个角度去拆RocketMQ你写一个生产者、一个消费者背后到底发生了什么为什么要这么设计真正落地的时候哪些参数必须调、哪些坑最好别踩。不管你是刚接触消息队列的新手还是已经在生产环境里被消息问题折腾过的老手这篇都可以当一份“编程模型速查手册”来读。1. 先看懂RocketMQ编程模型一条消息从发送到消费的完整链路1.1 四个核心角色各管什么RocketMQ的架构里一共四个角色NameServer、Broker、Producer、Consumer。我用大白话讲一下各自的分工。NameServer是“通讯录”维护着Broker的地址信息和Topic的路由信息。它不承担消息存储也不参与消息转发就是告诉客户端“你要找的Topic在哪个Broker上”。Broker是真正的“仓库”负责消息的存储、读写、副本同步。Producer是“发件人”把消息交到仓库。Consumer是“收件人”按需从仓库取消息。这些角色在编程模型里体现得非常直接。你写代码时用DefaultMQProducer和DefaultMQPushConsumer它们启动时会先从NameServer拉取路由表然后跟对应的Broker建立长连接消息就是这么通过一条链路流转起来的。提示NameServer本身是无状态的它不持久化任何消息数据。生产环境至少部署两台NameServer客户端会自动做故障切换。如果只部署一台宕机后Producer和Consumer还能继续用本地缓存的路由信息工作一段时间但新Topic创建、Broker变更这些操作就无法感知了。1.2 为什么说RocketMQ的编程模型对新手友好对比KafkaRocketMQ一个明显的差异是它把“消费”这件事包装得更接近业务直觉。Kafka的Consumer需要你手动管理partition的分配和offset的提交而RocketMQ的PushConsumer帮你做掉了大部分脏活消息到了自动推给回调方法offset自动提交分派逻辑由框架内部完成。这对业务的侵入非常小你只需要实现一个监听器处理收到的消息就行。但这不等于你可以完全不管底层。RocketMQ的“Push”本质上是“长轮询”Broker接收到请求后如果当前没有新消息会hold住这个请求一段时间等有消息了再返回。所以“Push”这个叫法更多是编程模型意义上的对你的代码来说消息像是被推过来的对底层来说其实是消费者主动去拉取。理解这一点你才能解释为什么消费端偶尔会出现短暂的消息延迟也能理解为什么消费线程数量需要手动控制。1.3 编程模型与实际存储结构的映射关系我建议所有写RocketMQ代码的人都把存储结构简单过一遍否则很多问题会变成玄学。RocketMQ的消息物理存储在一个叫CommitLog的文件里所有Topic的消息都混合写入这一个顺序文件。CommitLog是顺序写性能极高这是RocketMQ吞吐量的基石。但消息不能直接从CommitLog里按Topic捞出来所以Broker还会维护每个Topic下每个Queue对应的ConsumeQueue这个文件里存的是消息在CommitLog中的物理偏移量、消息长度、Tag HashCode这些轻量信息。消费者实际上是从ConsumeQueue上读的而不是直接扫CommitLog。这段逻辑在你写代码时是透明的但理解它有几个实际价值。第一Tag过滤之所以高效是因为ConsumeQueue里存了Tag的HashCode比对时不需要读完整消息体。第二消息堆积的实质是ConsumeQueue中的有效数据没有被消费掉堆积量的判断要以ConsumeQueue里的进度为准而不是看CommitLog的大小。第三如果你发现某个Topic的消息消费很慢不要想着去清理CommitLog那是所有Topic共享的只能从消费逻辑和队列数量上去优化。2. 编程模型里绕不开的核心概念与常见误区2.1 Topic、Tag、Queue这三个东西决定代码怎么写Topic是消息的一级分类想成“快递仓库里的一个货区”。你的系统里可能有“订单Topic”“支付Topic”“物流Topic”每个Topic承载一类业务消息。Broker在创建Topic时会规定它有多少个读写队列比如readQueueNums4、writeQueueNums4这个数量直接决定消息的并发度和消费的并行度。Tag是Topic下的二级分类想成“货区里的不同货架”。同一个订单Topic里可以分出TagOrderCreate、TagOrderPay、TagOrderCancel。生产者在发送消息时打上Tag消费者在订阅时指定Tag。这个机制最大的意义是让消费者可以精挑细选自己关心的消息避免所有消息都砸过来。Queue则是物理上的并行单位想成“货区里的几条传送带”。Producer发消息时默认会轮询写入多个QueueConsumer消费时每个Queue只会被同一个消费组里的一个消费者实例持有这就是并发的来源。如果你的Topic只有1个Queue那生产者的并发写入和消费者的处理速度再怎么调也上不去这是很多性能问题的根因。特别提醒一下Tag的使用误区Tag过滤是在Broker端基于HashCode粗筛的不是纯粹的字符串精确匹配所以不同Tag如果HashCode碰巧冲突概率很低但存在可能多拉取少量消息在消费者端你收到的消息里的Tag一定是准的你自己判断时只能用真正的Tag字符串去判断不能依赖看发送时的那个Tag来断言“绝对只有这些消息”。不过实际生产中遇到概率极低知道这层机制就行。2.2 消费组与位点集群消费的幕后黑手ConsumerGroup是一组消费者的集合几个消费者实例共享一个GroupName它们合起来逻辑上算一个消费集群。RocketMQ保证一条消息在一个消费组内只会被其中一个消费者实例消费。这就是“集群消费”的语义。这里的核心机制是Rebalance也就是Topic里的Queue会在消费组的多个实例之间重新分配。比如Topic有8个Queue消费组里有3个实例那么可能分配成3、3、2。分配完成后某个Queue只属于其中一台机器这台机器维护这个Queue的消费位点Offset。如果某个实例宕机了RocketMQ会在几秒内触发Rebalance把宕机实例名下的Queue分给其他存活实例。这里有一个非常关键的编程模型认知RocketMQ的消费并发度由“Queue数量”决定而不是由“消费者实例数量”决定。你开10个消费者进程但Topic只有4个Queue那么最多只有4个实例在消费剩下6个空转。反过来你有2个实例但Topic有16个Queue每个实例应该能分到8个Queue理论上可以起更多消费线程去吃。如果你发现消费吞吐上不去先看Topic的Queue数量够不够再看消费组实例数是否合理这两个是决定消费并行度的核心变量。Offset是消费位点记录每个Queue消费到哪里了。RocketMQ的Offset默认存储在Broker端由消费组维度维护。框架会自动定期提交但提交时机和业务处理结果有关。如果你用并发消费成功返回CONSUME_SUCCESS之后框架会异步提交位点如果业务逻辑抛出异常并返回RECONSUME_LATER位点不会前进这条消息会被重新投递。理解这个机制你就明白为什么顺序消费场景下不能乱用并发监听器。3. Producer生产端实战从发送一条消息看核心编程模型3.1 最小可运行的生产者代码先来一段能直接跑的生产者代码我把依赖和启动步骤都带上。dependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-client/artifactId version4.9.7/version /dependencypublic class ProducerDemo { public static void main(String[] args) throws Exception { DefaultMQProducer producer new DefaultMQProducer(demo_producer_group); producer.setNamesrvAddr(127.0.0.1:9876); producer.start(); for (int i 0; i 10; i) { Message message new Message( OrderTopic, TagOrderCreate, order_00 i, (create order, seq i).getBytes(StandardCharsets.UTF_8) ); SendResult result producer.send(message); System.out.printf(msgId%s, sendStatus%s, queueId%d%n, result.getMsgId(), result.getSendStatus(), result.getMessageQueue().getQueueId()); } producer.shutdown(); } }这段代码里最有信息量的其实是SendResult。getMsgId()返回的是客户端生成的消息全局唯一标识getQueueId()告诉你这条消息被写进了哪个Queue还能通过result.getMessageQueue().getBrokerName()看到落在哪个Broker节点上。把这些信息打出来对于后续排查消息去向有非常大的帮助。我习惯在测试环境把每条消息的MsgId、QueueId打到日志里一旦消息丢失或重复至少能判断是不是Producer侧就出了问题。3.2 同步、异步、单向三种发送方式该怎么选RocketMQ的Producer提供了三种发送方式选错会导致性能或可靠性问题。同步发送就是上面的代码调用producer.send(msg)后阻塞等待Broker返回写入结果。它适合对可靠性要求高的场景比如交易、支付、订单状态变更。同步发送的RT一般在几毫秒到几十毫秒吞吐量相比异步会低一些但能立刻知道成功还是失败方便做补偿。异步发送需要传入SendCallback回调producer.send(message, new SendCallback() { Override public void onSuccess(SendResult sendResult) { // 记录成功日志或处理后续逻辑 } Override public void onException(Throwable e) { // 失败处理比如记录日志、写入补偿表 } });异步发送适合流量大且可以容忍短暂延迟的场景。它不阻塞业务线程发送动作是异步的由框架内部线程池执行。千万注意回调方法里不要做耗时太长的操作否则会拖慢发送线程影响整体吞吐。回调里最合理的做法是记录日志、更新状态、或者把结果交给另一个线程池处理。单向发送最简单调用producer.sendOneway(message)。它的语义是“发出去就不管了”不等待Broker任何响应。适合日志采集、埋点数据这种丢了还能接受、但对发送时延极其敏感的场景。单向发送的时延最低因为省掉了等待RTT的时间但它无法得知消息是否真的写入了Broker。心得生产环境一般只在这三种模式里选同步和异步。单向发送用在对数据完整性要求不高的日志链路里。如果业务方问“能不能保证必达”不要用单向发送。3.3 生产端真正需要调的核心参数很多新手把Producer调通后就再也不动参数了但有几个默认值在生产环境是要改的。第一个是setSendMsgTimeout默认3000毫秒。如果你的消息比较大或者Broker端刷盘压力高3秒超时会频繁触发超时重试。我处理过一个案例业务方发的消息体里塞了一个5MB的JSON字符串同步发送经常超时后来把超时时间调到10000毫秒才稳定下来。当然更合理的做法是控制消息体大小RocketMQ默认限制单条消息最大4MB。第二个是setRetryTimesWhenSendFailed默认2次。这里有个数据正确性的陷阱同步发送失败重试时如果第一次请求实际已经写入Broker只是响应超时没收到那么重试会形成重复消息。RocketMQ通过msgId没法做服务端去重所以业务侧必须做幂等。我在生产上见过不少“明明消费者做了幂等还是出现重复数据”的案例最后定位到Producer端重试逻辑上。解决思路是给业务消息带一个全局唯一的业务主键消费者端对这个主键做去重。第三个是setMaxMessageSize默认4MB。如果业务确实需要发大消息可以调大但我不建议超过10MB。大消息会带来两个问题占用Broker大量内存、拖慢ConsumeQueue的消费效率。更合理的方案是把大对象存到对象存储或数据库里MQ里只发一个引用文件路径或对象ID。我之前在物流行业做过一个项目轨迹数据里有大量图片Base64直接把消息体干到8MB后来改成图片传OSS、MQ只发图片ID消费者再主动拉取整体链路稳定太多了。3.4 消息Key与顺序发送的小技巧生产者在构造Message时可以指定第三个参数也就是Key。这个Key会存进IndexFile索引里用于根据业务唯一键快速查询消息。比如订单号、用户ID、请求ID都可以作为Key。排查问题时用mqadmin queryMsgByKey命令能快速定位一条消息在哪个Broker的哪个文件位置。我强烈建议所有核心业务消息都带上业务主键作为Key这个习惯在线上排查时会省下大量时间。顺序发送这块也值得提一下。RocketMQ的顺序消息依赖MessageQueueSelector你在发送时指定把同一个业务维度比如同一个订单ID的消息发到同一个Queue里。框架提供了SelectMessageQueueByHash选择器按业务Key哈希取模选Queueproducer.send(message, new SelectMessageQueueByHash(), orderId);这里的关键是顺序的保证范围是“同一个Queue”跨Queue不可能全局有序。所以发送端的Hash策略必须稳定同一个订单的消息一定要落到同一个Queue中。如果Hash策略中途变更或者Queue数量变了顺序就会被打乱。扩容Topic的Queue数量时务必评估对顺序消息的影响。4. Consumer消费端实战消费编程模型的四种玩法与参数陷阱4.1 Push与Pull两种消费模型的编程差异RocketMQ同时提供了DefaultMQPushConsumer和DefaultMQPullConsumer两种消费客户端叫法上容易让人误会我解释一下。DefaultMQPushConsumer是绝大多数场景的选择。它看起来是服务端把消息推给你实际底层用长轮询在拉取但框架帮你封装好了线程模型、重试逻辑、消息等分逻辑。你用起来只需要做三件事设置消费组、订阅Topic和Tag、注册消息监听器。public class ConsumerDemo { public static void main(String[] args) throws Exception { DefaultMQPushConsumer consumer new DefaultMQPushConsumer(demo_consumer_group); consumer.setNamesrvAddr(127.0.0.1:9876); consumer.subscribe(OrderTopic, TagOrderCreate || TagOrderPay); consumer.registerMessageListener(new MessageListenerConcurrently() { Override public ConsumeConcurrentlyStatus consumeMessage(ListMessageExt msgs, ConsumeConcurrentlyContext context) { for (MessageExt msg : msgs) { String body new String(msg.getBody(), StandardCharsets.UTF_8); System.out.printf(msgId%s, queueId%d, body%s%n, msg.getMsgId(), msg.getQueueId(), body); } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } }); consumer.start(); } }subscribe的第一个参数是Topic第二个参数是Tag表达式支持||分隔多个Tag。如果传入*表示订阅Topic下所有消息。注意Tag是精确匹配不是正则只有||这种枚举和*通配。DefaultMQPullConsumer则需要自己写拉取循环手动获取Queue列表、手动更新Offset、手动处理分派。它的编程模型原始一些但给开发者最大的控制权。我在两个场景下才会用PullConsumer一是消费逻辑需要精确控制Offset比如从指定时间点补数据二是业务场景需要自定义拉取频率PushConsumer的节奏框架说了算。除此之外一律建议用PushConsumer别自己造轮子。4.2 集群消费与广播消费语义差异决定了消息是否重复DefaultMQPushConsumer默认是集群消费模式。一个消费组内的多个实例共同消费一个Topic每条消息只会被组内的一个实例消费这是业务中使用最多的模式。广播消费需要显式设置consumer.setMessageModel(MessageModel.BROADCASTING);广播模式下消费组内每个实例都会收到每条消息。适合的场景是“每个节点都想知道全量事件”比如配置变更、缓存刷新通知。广播模式的弊端是消费组内每个实例都维护自己的Offset消息量大的时候N个实例会重复处理N份全量数据成本成倍上涨。另外广播模式下Offset存储在本地文件更换机器或本地文件丢失会导致重复消费用的时候要做好幂等。4.3 并发消费与顺序消费回归编程模型的核心差异MessageListenerConcurrently对应并发消费适合大多数业务场景处理效率高。它对消息的可靠性要求是“尽量成功”如果业务处理失败返回RECONSUME_LATER这条消息会稍后重新投递但不能保证顺序。MessageListenerOrderly对应顺序消费consumer.registerMessageListener(new MessageListenerOrderly() { Override public ConsumeOrderlyStatus consumeMessage(ListMessageExt msgs, ConsumeOrderlyContext context) { // 处理消息 return ConsumeOrderlyStatus.SUCCESS; } });框架会为每个Queue加锁并且对消息做定时批次提交确保同一个Queue里的消息按顺序交给监听器处理。但注意顺序消费的吞吐量明显低于并发消费因为同一Queue内是串行的。如果你负责的业务真的要求严格有序那么发送端要按业务Key选择同一个Queue消费端用Orderly监听器两边都做好顺序才成立。只改消费端不改发送端顺序基本等于空谈。RocketMQ 5.x里新增了PopConsumer这种消费形态它把消费进度从Broker端进一步解耦支持轻量级消费者。但经典编程模型中的PushConsumer仍然兼容目前存量业务里绝大多数用的还是4.x风格的PushConsumer。新项目可以关注5.x老项目不用急着折腾。4.4 消费线程数与批量参数吞吐量的隐形开关消费端的性能问题很多时候不是代码逻辑慢而是参数没调。DefaultMQPushConsumer有几个参数直接影响消费效率。setConsumeThreadMin和setConsumeThreadMax控制消费线程池的大小默认各为20。这里有个容易被忽略的细节如果Topic的Queue数量少于消费线程数多出来的线程是闲置的因为每个Queue同一时间只能被一个线程消费。我之前给一个团队排查消费积压看配置开了30个消费线程但Topic只有8个Queue实际上只有8个线程在干活。先加Queue数量再调消费线程才有效。setConsumeMessageBatchMaxSize控制每次批量投递给监听器的消息条数默认是1也就是一次只处理一条。你可以调大到10甚至32减少回调次数提高吞吐。但要注意批量增大后监听器里的处理逻辑必须对“一批消息部分失败”有准备如果你返回RECONSUME_LATER这一整批都会被重新投递而不仅仅是失败的那条。这会导致消息重复幂等设计必须覆盖这种情况。setPullBatchSize是每次长轮询拉取的消息条数默认32。如果你的每条消息处理耗时较长可以把pullBatchSize调小避免大堆消息积压在消费端的本地队列里造成内存压力和消费超时。相反如果处理很快可以适当调大。注意消费端还有一个隐形约束叫consumeTimeout默认15分钟。如果一条消息处理超过15分钟Broker端会认为消费者处理不了会触发消息重新投递。你有两类选择一是把业务处理拆细别在监听器里做超长任务二是长任务丢给独立线程池异步化先返回成功后面再处理。但异步化的代价是消息确认了但实际可能没处理完掉电或重启会丢这个权衡要自己拿捏。5. 高频问题排查与避坑实录5.1 消息堆积了怎么定位、怎么处理消息堆积是使用RocketMQ最常遇到的“事故”。首先要会看堆积最简单的方式是RocketMQ Dashboard它会在Consumer页面展示每个消费组的“堆积量”。如果你是命令行爱好者用mqadmin consumerProgress -g 消费组名一样能看。这两个数字的逻辑是Broker端当前最大Offset减去消费组已提交的Offset差值就是堆积量。有了堆积量后按这个顺序排查先看消费组实例数是否正常。如果只有1个Consumer实例但Topic有16个Queue那么只有这1个实例承担全部消费压力堆积几乎是必然的。扩容消费组实例到合适的数量让每个实例分到尽量均衡的Queue。再看消费线程和消费逻辑耗时。如果消费组成员正常但消费TPS远低于生产TPS去看消费逻辑里有没有慢调用。我处理过一个堆积事故业务方在消息监听器里同步调用了第三方短信接口一条消息处理耗时3秒TPS自然上不去。把短信发送改成异步削峰之后堆积很快就消化了。最后看Queue数量。如果消费组实例数和消费逻辑都优化过了堆积还是消不掉说明Topic的Queue数量不够无法水平扩展消费并行度。可以扩容Queue数量但要注意两点第一Queue数量的变化会影响顺序消息的顺序保证第二扩容Queue之后存量堆积的消息还在老的Queue里新的Queue不会自动承接存量所以扩容对已有堆积的缓解是有限的核心还是提高单Queue的消费效率。5.2 重复消费为什么防不住幂等怎么做RocketMQ的消息投递语义是At Least Once翻译成人话就是“至少一次”所以重复消费不是bug而是一个需要被设计的常规情况。重复的来源主要有三个Producer端发送重试Consumer端消费成功后还没来得及提交Offset就宕机或Rebalance批量消费时某条消息失败导致整批重投。应对重复消费唯一可靠的办法是业务侧幂等。幂等设计常见方案有这么几种第一种是唯一主键去重。在消息里带上业务唯一键消费端先查数据库或Redis如果这个主键已经处理过了就跳过。比如订单号、支付流水号。这种方案实现简单只是每次消费多一次查询开销。第二种是利用数据库唯一索引。消费逻辑里insert的数据带唯一索引重复插入直接报冲突捕获冲突后当作成功处理。这种方案适合新增业务数据不适合更新场景。第三种是状态机校验。如果业务数据有明确状态流转比如“已创建-已支付-已发货”消费端先查当前状态如果当前状态已经大于等于目标状态说明消息重复了直接返回成功。这种方案不仅防重还能防乱序是我在交易类场景最推荐的姿势。幂等还有一个容易被忽略的细节幂等判断应该放在监听器入口也就是拿到消息最先做的事情不是处理完业务再判断。因为消费超时后消息会重投如果每个人都在处理完整个业务逻辑之后才做“已处理”标记那重复消费期间业务逻辑还是会跑两遍。先做去重判断再执行真正的业务逻辑才是正确的顺序。5.3 Windows部署与可视化工具避坑RocketMQ官方虽然主打Linux部署但很多个人开发机和测试环境是Windows。Windows下部署RocketMQ有一堆细节我踩过的坑值得拿出来说。首先JDK版本要匹配。RocketMQ 4.9.x建议JDK 8RocketMQ 5.x建议JDK 8以上。Windows的启动脚本在bin目录下启动NameServer用mqnamesrv.cmd启动Broker用mqbroker.cmd。但有个大坑脚本默认加载的配置和Linux不一致需要手动设置环境变量。在Windows上需要先设置ROCKETMQ_HOME指向解压目录否则脚本会找不到配置文件。我建议用命令行设置set ROCKETMQ_HOMED:\rocketmq-4.9.7 set NAMESRV_ADDR127.0.0.1:9876然后启动NameServermqnamesrv.cmd启动成功后会看到The Name Server boot success默认端口9876。接着启动Brokermqbroker.cmd -n 127.0.0.1:9876启动Broker时需要修改conf/broker.conf里的brokerIP1在Windows上设置为实际IP不能是默认的127.0.0.1否则其他机器连不上Broker。启动命令里用-c指定配置文件mqbroker.cmd -n 127.0.0.1:9876 -c D:\rocketmq-4.9.7\conf\broker.conf还有个经典问题Windows下启动Broker时提示内存不足或者脚本闪退。原因是runbroker.cmd里默认的JVM参数是-Xms8g -Xmx8g个人电脑一般没这么大内存。需要改成-Xms512m -Xmx512m。同理runserver.cmd里默认-Xms4g -Xmx4g也改成小内存。可视化工具有两种路径。老牌的rocketmq-dashboard是Web应用能看Topic、Consumer、消息轨迹功能全面。它有单独的GitHub仓库下载后改配置文件里的rocketmq.config.namesrvAddr然后用Maven打jar包运行。新版RocketMQ 5.x在发行包里自带了一个Dashboard模块路径在dashboard目录下部署更方便。Dashboard除了看堆积还能按MsgId或Key查消息轨迹我建议生产环境一定部署一套它比命令行排查效率高太多。给普罗米修斯接入RocketMQ监控这个需求也比较常见。核心思路是利用RocketMQ的metrics模块或者通过JMX暴露指标再用Prometheus采集。RocketMQ 4.x里Broker默认支持JMX端口是10911可以用jmx_exporter采集。5.x提供了更原生的metrics支持配置好后就等着Prometheus定期抓即可。重点监控的指标包括ConsumerLag消费堆积、PullLatency、SendRT、以及Broker端的CommitLogSize。5.4 消息选型速查什么时候坚持RocketMQ结合搜索结果里频繁出现的“Kafka、RabbitMQ、RocketMQ选型”话题我也简单聊一下我自己的判断逻辑。如果你需要的是大吞吐量的日志管道、数据总线Kafka是首选它的顺序写盘和零拷贝在日志场景下确实压制全场。如果你追求轻量、灵活路由、各种交换机模型RabbitMQ更好。但如果你在一个中大型业务系统里需要事务消息、延迟消息、消息轨迹、以及更接近业务语义的消费模式RocketMQ的综合体验是最好的。它能处理常规的并发量还能把死信队列、重试队列这些“企业级能力”直接给到位。很多团队选RocketMQ不是因为它比Kafka快而是因为“消息中间件要伺候业务”而RocketMQ在这一点上做得最顺手。从我的经验看刚上手RocketMQ最容易犯的错误就是堆概念、背API却没有把“消息从生产到消费走一遍完整链路”这件事在脑海里形成一个清晰图景。编程模型的本质就是这张链路图怎么落到代码上Producer发送时选择Topic、Tag、QueueConsumer订阅时关心Group、Offset、RebalanceBroker在中间把存储和路由处理妥当。你在纸上把这几个角色的职责画清楚再去写代码很多参数不用硬记也能猜到它是干什么的。最后再分享一个我用了很久的排查习惯把生产环境的SendResult日志和消费端监听器入口日志都打开关键业务消息在入口处打上业务主键和时间戳。这样消息一旦出问题顺着业务主键就能把生产端、Broker、消费端三段日志串起来绝大多数消息丢失、重复、积压问题都能在十分钟内定位到具体环节。消息中间件的排障本质上就是把这个链路每一环的日志抓准、抓全。