ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

深入理解RocketMq架构原理及案例

深入理解RocketMq架构原理及案例 rocketMq架构讲解https://www.doubao.com/thread/xfb3b8b73d229807591caf746c83f84551.集群消费模式Cluster下同一个 MessageQueue消息队列/分区在同一个消费者组ConsumerGroup内任意时刻最多只能有一台消费者实例在消费。❌ 误区纠正不是「topic 只能一个消费者消费」不是「整个 broker 只能一个消费者消费」限制粒度是Queue分区 消费组1-1.为什么要这么设计核心目的保证顺序消费基础如果一个队列同时被多个消费者拉取消息会被两台机器并发处理局部顺序消息彻底失效。RocketMQ 顺序消息原理同一队列内消息串行消费前提就是队列只能绑定一个消费者。避免同组内消息重复消费集群消费业务诉求一条消息组内只消费一次。如果队列分给多个消费者消息会被多个实例同时拉取出现重复。Offset 位点管理消费位点offset是绑定在「消费组 queue」维度。一个队列同时多个客户端消费offset 提交会错乱。1-2.广播消费模式下规则失效广播消费同一个组内所有消费者都会收到全部消息。 此时一个队列会同时被组内所有消费者消费每条消息组内每台机器都收到适合配置同步、通知推送场景。1-3.重平衡Rebalance和这条规则的关系当消费者上下线会触发重平衡Broker / 客户端重新分配队列给组内实例。分配完成后再次遵守单个 Queue 仅归属一个实例。在重平衡瞬间可能短暂出现旧消费者还在处理消息新消费者开始拉取 → 极短窗口存在重复消费业务要做好幂等。2.广播消费模式下一个队列同时多个客户端消费offset 提交会错乱么?广播消费并不会出现 offset 错乱根本原因两种模式 Offset 的存储模型完全不同。1、集群消费Cluster为什么必须限制一个 Queue 同一时间只能一个消费者消费集群消费✅Offset 存在 Broker同一个消费组共享一份【groupqueue】的 offset假设放开限制允许多个消费者同时消费同一个 Queue消费者 A 消费到 offset100提交到 Broker → 位点更新为 100消费者 B 同时消费到 offset90提交到 Broker → 位点被回退到 90 offset 互相覆盖、错乱、消息重复 / 无限循环消费所以 RocketMQ 用「队列独占分配」来规避这个冲突。一句话集群模式共享 offset→ 不能多实例同时消费同一个队列。2、广播消费Broadcast如何规避这个问题广播模式设计上直接换了一套方案✅Offset 不保存在 Broker每个消费者实例独立把 offset 写到【本地磁盘文件】路径~/.rocketmq_offsets/clientId/groupName/offsets.json关键点消费者 A 的 offset 和 消费者 B 的 offset 互相隔离互不感知、互不覆盖A 更新自己本地文件B 更新自己本地文件大家不需要争抢同一个位点不需要 Broker 统一维护消费进度自然不存在 offset 互相覆盖错乱的问题所以广播模式允许同一个队列同时被组内所有消费者消费。3、一张对比表把逻辑打通表格模式Offset 存储位置同组消费者 offset 关系是否允许多实例同时消费同一个 Queue约束来源集群消费Broker 远端共享同一份 offset❌ 禁止会 offset 互相覆盖错乱Rebalance 队列独占分配广播消费客户端本地文件每个实例独立 offset互不干扰✅ 允许无队列分配全部实例消费全部队列3.普通消息和顺序消息含义消息本身没有 “天生属性”是【发送方式 消费方式】共同决定这条消息能不能实现顺序消费。不是消息打上标签区分普通 / 顺序而是你调用不同发送 API、路由策略决定。1. 普通消息默认发送方式发送特点java运行// 普通发送不指定队列轮询分发 producer.send(msg);RocketMQ 默认策略多条消息轮询发到不同队列。Topic 有 4 个队列消息依次往 queue0、queue1、queue2、queue0… 分发。 后果同一业务的多条消息分散在不同队列天然无法保证全局有序。就算你消费者单线程消费也解决不了创建订单在 queue0支付消息在 queue1顺序没法保证。2. 顺序消息分区有序发送特点java运行// 传入selector相同orderId路由到同一个队列 producer.send(msg, selector, orderId);关键相同业务标识订单 ID、用户 ID的消息强制发送到同一个 MessageQueue。示例订单 1001创建、支付、发货 三条消息 → 全部进入 queue2订单 1002创建、支付消息 → 全部进入 queue3配套约束缺一不可否则依然乱序发送相同业务 key 路由到同一个队列消费该队列只分配给一个消费进程集群模式自带保证消费端单线程串行消费该队列消息三个条件全部满足才能实现业务有序。补充一个形象例子地铁Topic有 4 条轨道4 个 Queue普通消息乘客随机分配到 4 条轨道顺序消息同一个旅行团所有人强制走同一条轨道即便同一个轨道只有一辆列车一个消费进程如果列车内部多个工作人员并行处理乘客多线程依然有可能先处理后面上来的乘客。想要有序必须轨道内排队依次处理单线程消费。4.消息的并发如何实现的?核心先行并发分为三层跨队列并发多实例、多机器最重要同一个队列,同实例内部多线程并发单机并行消息批量拉取并发限制同一个队列无法跨实例并发但可以单机内部并发普通消息一、跨队列并发水平扩展真正提升吞吐量规则回顾Topic 拆分多个 MessageQueue分区。集群消费下重平衡会把不同 Queue 分配给组内不同消费者实例。举例Topic 有 8 个队列消费者组启动 2 台实例实例 1Queue0、1、2、3实例 2Queue4、5、6、7✅两个实例同时消费互不干扰真正并行 这就是 RocketMQ 水平扩容的基础。关键约束重中之重一个队列同一时间只能分配给组内一个消费者实例所以队列数量 最大有效并行消费者实例上限例如 Topic 4 个队列你启动 10 个消费者最多只有 4 个实例分到队列剩下 6 个空闲无法提升并发。想要扩容先增加队列数量。二、单实例内部消费线程并发consumerThreadNum当一个消费者进程拿到若干队列后开启多线程处理消息。这里区分普通消息 VS 顺序消息1普通消息 ✅ 支持多线程并发消费者一次性从 Broker 批量拉取多条消息丢到线程池并行处理。缺陷同一个队列的多条消息处理顺序无法保证前面讲过。优势单机吞吐量大幅提高。伪流程plaintext消费者拉取 Queue0 [msg1,msg2,msg3,msg4] 线程1 处理 msg1 线程2 处理 msg2 线程3 处理 msg3 线程4 处理 msg42顺序消息 ❌ 禁止同一个队列多线程并发框架强制同一个队列的消息串行处理一条消费成功才会取下一条所以顺序消息单机只能单线程跑单个队列并发能力弱。想要提升顺序消息并发增加队列数量把不同业务哈希分发到不同队列依靠【跨队列并发】。三、批量拉取辅助提升并发消费者不是一条一条拉默认批量拉取多条消息减少网络 IO。配合多线程进一步放大处理能力。四、一张全景图理解整体并发模型假设Topic 8 队列2 台消费者实例consumerThreadNum5普通消息实例 A 负责 Queue0~3内部 5 个线程并发处理这 4 个队列消息实例 B 负责 Queue4~7内部 5 个线程并发处理这 4 个队列消息✅ 整体并发 跨机器并发 单机线程池并发五、几个高频误区澄清误区 1多加消费者实例就能无限提高并发不行。有效并发上限 ≤ Topic 队列数量队列数 4消费者开到 10 台最多 4 台干活其余闲置。误区 2顺序消息调高消费线程数提升性能没用。顺序消息框架内部对单个队列串行执行线程再多同一个队列消息依旧排队多个队列之间可以并行。误区 3能不能让同一个队列被多个实例并发消费提升性能集群消费模式做不到原因共享 offset多实例同时消费同一队列会覆盖位点消息重复、死循环。如果你非要多实例消费同一队列新建不同消费组组之间相互独立但这不是分担压力是一份消息重复消费多遍。六、两种场景如何设计并发方案场景 1普通消息无顺序要求追求高吞吐设置合理队列数量预估最大消费者实例数多台消费者实例水平部署适当调大 consumerThreadNum单机多线程并发场景 2分区顺序消息订单、流程类消息提高队列数量并发扩容依靠队列相同业务 ID 哈希路由到同一个队列消费端线程数保持默认不要指望单队列多线程提速极简总结跨机器并发依靠多个队列队列分配给不同实例单机并发依靠消费线程池仅普通消息可用同一队列无法跨实例并发集群消费设计约束防止 offset 混乱顺序消息牺牲单机内部并发换取时序只能依靠增加队列实现整体扩容。5.如何保持消息不丢5.1生产者保持消息不丢核心风险消息发出去网络闪断、Broker 没收到生产者以为发送成功。方案同步发送 SyncSend生产者阻塞等待 Broker 返回确认 ACK收到成功响应才算发送完成异步发送不能实时感知失败容易丢消息。发送重试机制同步发送失败超时、网络异常Producer 自带重试⚠️ 注意重试会产生重复消息下游必须幂等。事务消息解决本地事务 MQ 消息一致性防止本地事务成功、消息没发出或者消息发出、本地事务回滚。补充边界即使同步发送 重试如果 Broker 宕机、消息只写入内存还未落盘依然有丢失风险需要 Broker 侧可靠性。如果生产者发送给重试依旧失败呢https://www.doubao.com/thread/xS4O77Bxyu2MArFTV代码案例https://www.doubao.com/thread/xve5fmy6JeKoNl5Y5核心结论先讲清楚生产者开启重试重试耗尽仍然发送失败 → 这条消息就会在生产者侧丢失。RocketMQ 本身不会帮你保存这条待发消息没有内置 “消息本地持久化排队补发” 能力。我们分层拆解整套流程、风险、解决方案。1. 先搞懂Producer 重试范围默认同步发送send()参数retryTimesWhenSendFailed同步发送重试次数默认 2 次首次发送 2 次重试一共尝试 3 次。什么情况会触发生产者重试网络超时、连接异常、Broker 返回系统繁忙等可重试异常。什么情况不会重试Broker 返回业务错误topic 不存在、权限不足、消息超限、非法参数直接失败不重试。当所有重试次数全部用完依然发送失败send()方法直接抛出异常。消息仅此一份在内存没有落地程序不处理 消息消失。2.典型错误代码消息丢失现场java运行try { producer.send(msg); } catch (Exception e) { // 仅仅打印日志啥也不干 log.error(发送失败,e); } // 消息直接丢掉3. 有哪些兜底方案按推荐优先级排序方案 1本地失败消息表工业标准方案首选发送消息前先把消息存入业务库【消息发送表】状态 待发送1. DB insert msg_record(消息体,唯一bizId,status待发送) 2. 调用mq发送 3. 发送成功 → update status已发送 4. 发送失败 → 不更新状态额外启动定时任务扫描待发送消息循环重试投递。优势宕机重启后定时任务依然能补偿完全自主可控。这也是事务消息 half 思路的简化思想。了解方案 2本地磁盘文件兜底不推荐生产主流程发送失败把消息写入本地文件。缺点机器宕机、磁盘损坏文件丢失多实例不好统一管理不适合集群多机器部署。适合小型项目临时应急不能当做核心保障。了解案 3外部消息缓冲队列Redis List发送失败写入 Redis独立消费线程不断从 Redis 取出重试。风险点Redis 宕机也会丢需要做好持久化。4.延伸面试高频问题事务消息能解决这个问题吗不能直接解决事务消息解决的是「本地事务 和 MQ 消息 最终一致性」场景本地事务执行成功保证消息一定发出本地事务失败消息不投递。但是事务消息的 half 消息发送阶段如果网络问题多次重试依然发送失败half 消息无法存入 Broker一样丢失。事务消息依然无法规避【生产者无法和 Broker 建立通信】的极端故障极端场景依然需要外部补偿表兜底。5. 极简总结版面试答题可用生产者重试次数耗尽仍然发送失败消息仅存在生产者内存RocketMQ 不会保存消息若不手动处理直接丢失内置重试只是短时故障自救无法应对长时间 Broker 不可用标准可靠方案发送前落库【消息发送记录表】成功后更新状态定时任务扫描待发送消息持续补偿注意区分生产者发送失败无死信消费失败超限才会进入 Broker 死信队列事务消息无法完全规避该问题极端网络故障下依然需要外部补偿机制。6.代码案例https://www.doubao.com/thread/xgB77M58wArUG2i875.2防 Broker 服务端消息丢失消息已经到达 Broker消息抵达 Broker 后存在内存还没持久化 / 同步副本Broker 宕机就丢失。刷盘策略同步刷盘消息写入磁盘成功才返回 ACK 给生产者可靠性最高吞吐低异步刷盘写入 PageCache 就返回后台异步落盘宕机存在丢消息风险默认主从复制策略同步主从SYNC_MASTER主节点必须同步给从节点成功再响应生产者异步主从ASYNC_MASTER主写完直接返回后台同步从节点主宕机未同步数据丢失。✅ 最高可靠组合同步刷盘 同步主从金融级场景性能损耗大普通业务异步刷盘 异步主从配合重试取舍性能。防 Broker 服务端消息丢失消息已经到达 Brokerhttps://www.doubao.com/thread/xWIgltZxbNw43cujw先说明重点同步刷盘、同步主从属于 Broker 服务端配置*不在业务生产者代码里控制*生产者代码无法开启刷盘策略、主从复制策略。生产者能控制的只有发送模式同步 / 异步、消息发送等待确认级别。一、先分清两个概念Broker 侧配置运维在 broker.conf 配置无代码properties# 刷盘策略 flushDiskTypeSYNC_FLUSH # 同步刷盘 / ASYNC_FLUSH 异步刷盘(默认) # 主从复制策略 brokerRoleSYNC_MASTER # 同步主从 / ASYNC_MASTER 异步主从✅ 组合最高可靠SYNC_FLUSH SYNC_MASTER含义消息落盘成功 同步复制到从节点成功才返回成功给生产者。生产者代码能控制waitStoreMsgOK存储等待确认默认waitStoreMsgOKtrue生产者等待 Broker 存储结果如果设为 falseBroker 写入 pagecache 直接返回不等待落盘风险更高。1. 生产者代码示例开启等待存储确认推荐java运行DefaultMQProducer producer new DefaultMQProducer(PRODUCER_GROUP); producer.setNamesrvAddr(127.0.0.1:9876); // 关键等待Broker存储完成再返回默认true建议显式写出 producer.setWaitStoreMsgOK(true); producer.start(); Message msg new Message(TEST_TOPIC, TAG, 消息内容.getBytes(StandardCharsets.UTF_8)); try { // 同步发送 SendResult result producer.send(msg); System.out.println(发送结果 result.getSendStatus()); } catch (Exception e) { log.error(发送异常, e); }waitStoreMsgOK 作用trueBroker 执行完刷盘、主从同步逻辑之后才回复 ACKfalseBroker 收到消息写入内存直接返回不等持久化。⚠️ 只是 “等待结果”不能改变 Broker 是同步刷盘还是异步刷盘。二、发送时指定消息级别可选高可靠场景设置RocketMQ 消息支持消息级别可以单独对某条消息强制要求同步刷盘不依赖 Broker 全局配置java运行Message msg new Message(TEST_TOPIC, TAG, body.getBytes()); // 设置消息存储要求同步刷盘 msg.setWaitStoreMsgOK(true); SendResult result producer.send(msg);注意老版本部分 RocketMQ该参数最终受制于 Broker 全局配置无法单独覆盖全局 flushDiskType。生产统一推荐 Broker 全局配置管理。三、区分容易混淆的面试考点高频问题 1我代码怎么开启同步刷盘答代码无法开启同步刷盘必须修改 broker.conf重启 Broker。生产者只能控制是否等待存储结果不能修改 Broker 持久化策略。问题 2同步发送 ≠ 同步刷盘plaintext同步发送生产者阻塞等待响应 ≠ 同步刷盘Broker消息落地磁盘 ≠ 同步主从Broker同步复制给从节点三者独立可以同步发送 Broker 异步刷盘也可以异步发送 Broker 同步刷盘。四、模拟风险场景演示直观理解丢消息风险场景Broker 使用默认配置ASYNC_FLUSH 异步刷盘生产者同步发送收到成功响应Broker 仅写入 PageCache还没后台线程刷入磁盘Broker 进程立即宕机、服务器断电。 PageCache 数据丢失 → 消息永久丢失。解决手段Broker 开启 SYNC_FLUSH 同步刷盘。五、完整运维配置示例 broker.conf高可用方案properties# 主节点配置 brokerNamebroker-a brokerId0 brokerRoleSYNC_MASTER # 同步主从 flushDiskTypeSYNC_FLUSH # 同步刷盘 namesrvAddr127.0.0.1:9876从节点配置propertiesbrokerNamebroker-a brokerId1 brokerRoleSLAVE namesrvAddr127.0.0.1:9876六、生产落地建议总结普通业务Broker 默认ASYNC_FLUSH ASYNC_MASTER追求吞吐量依靠生产者补偿表兜底。金融、订单资金等高可靠场景运维调整 BrokerSYNC_FLUSH SYNC_MASTER生产者代码开启waitStoreMsgOKtrue默认已开启使用同步发送。七、面试标准答题模板Broker 侧消息丢失防护分为运维配置 生产者配合Broker 配置开启同步刷盘SYNC_FLUSH消息写入磁盘成功后再响应开启同步主从SYNC_MASTER消息同步复制到从节点成功再返回生产者代码层面启用 waitStoreMsgOKtrue等待 Broker 存储结果 ACK注意刷盘、主从策略属于 Broker 配置业务代码无法动态修改只能运维调整。5.3防消费端消息丢失消息投递到消费者根本原则先消费成功再提交 OffsetRocketMQ 默认机制Broker 推送消息给消费者不会立刻移动 offset消费者执行业务逻辑处理完成主动返回 CONSUME_SUCCESSBroker 收到成功响应才更新消费位点代表这条消息处理完毕。反面错误场景消息丢失经典坑先提交 offset再执行业务业务中途崩溃 → 消息永久丢失消费逻辑抛出异常开发者直接捕获异常不返回重试状态框架认为消费完成推进 offset。失败重试机制消费返回失败状态RECONSUME_LATER普通消息消息退回 Broker延迟一段时间重新投递超过最大重试次数进入死信队列。顺序消息不会跳过消息持续阻塞重试保证有序不被破坏。重点死信队列也要定时巡检否则消息堆积在 DLQ 无人处理等同于业务丢失防消费端消息丢失 核心原理核心铁律业务处理成功 → 再返回消费成功 不要先提交 offset再执行业务RocketMQ PushConsumer 模式消费者返回CONSUME_SUCCESSBroker 才更新消费位点offset。如果业务执行异常返回RECONSUME_LATERBroker 会重新投递这条消息。❌ 致命错误捕获所有异常直接返回成功消息直接被丢弃。代码案例https://www.doubao.com/thread/xIzbwYm1ClXeo0NGrhttps://www.doubao.com/thread/xeoO5DuQnXDOsIPjajava 版一、标准安全消费代码推荐防止消息丢失java运行public class SafeConsumerDemo { public static void main(String[] args) throws MQClientException { DefaultMQPushConsumer consumer new DefaultMQPushConsumer(ORDER_CONSUMER_GROUP); consumer.setNamesrvAddr(127.0.0.1:9876); consumer.subscribe(ORDER_TOPIC, *); // 消息监听器 consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) - { // RocketMQ 批量消息这里只取第一条演示生产注意批量处理逻辑 MessageExt msg msgs.get(0); try { String body new String(msg.getBody(), StandardCharsets.UTF_8); log.info(开始消费消息{}, body); // 执行业务逻辑 handleBusiness(body); // // ✅ 业务正常执行完毕返回成功Broker推进offset return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } catch (Exception e) { log.error(消费业务异常消息需要重试, e); // ✅ 业务失败通知Broker稍后重新投递消息 return ConsumeConcurrentlyStatus.RECONSUME_LATER; } }); consumer.start(); log.info(消费者启动成功); } // 模拟业务处理 private static void handleBusiness(String body) { // 数据库操作、调用第三方接口等业务逻辑 } }二、几种典型【消息丢失错误代码】面试高频坑错误 1捕获异常直接返回成功最常见消息永久丢失java运行// ❌ 禁止这么写 consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) - { MessageExt msg msgs.get(0); try { handleBusiness(new String(msg.getBody())); } catch (Exception e) { log.error(消费异常, e); // 异常依旧返回成功 → Broker认为消费完成offset前进消息丢失 } return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; });错误 2手动异步处理消息监听器线程直接返回成功java运行// ❌ 危险代码 consumer.registerMessageListener((msgs, context) - { MessageExt msg msgs.get(0); // 丢给异步线程当前线程立刻返回成功 executor.submit(() - handleBusiness(new String(msg.getBody()))); // 异步线程还没执行offset已经提交异步线程报错消息无法重试 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; });结论监听器同步执行业务不要把消息丢到外部线程池异步处理。三、补充两个关键机制重试与死信RECONSUME_LATER 不是立刻重发Broker 会延迟投递并且重试次数递增延迟每条消息有最大重试次数默认 16 次。超过最大重试次数消息自动转入死信队列 % DLQ% 消费组名死信消息不再自动重试需要单独监听 DLQ 告警人工处理否则等同于消息丢失。监听死信队列示例兜底java运行DefaultMQPushConsumer dlqConsumer new DefaultMQPushConsumer(ORDER_CONSUMER_GROUP); dlqConsumer.setNamesrvAddr(127.0.0.1:9876); // 订阅当前消费组对应的死信Topic dlqConsumer.subscribe(%DLQ%ORDER_CONSUMER_GROUP, *); dlqConsumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) - { MessageExt msg msgs.get(0); log.error(收到死信消息需要人工排查:{}, new String(msg.getBody())); // 告警钉钉/短信 return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; }); dlqConsumer.start();四、顺序消费者特殊写法顺序消息不能并发重试顺序消息监听器使用MessageListenerOrderlyjava运行consumer.registerMessageListener((MessageListenerOrderly) (msgs, context) - { MessageExt msg msgs.get(0); try { handleBusiness(new String(msg.getBody())); return ConsumeOrderlyStatus.SUCCESS; } catch (Exception e) { log.error(顺序消息消费失败, e); // 顺序消息SUSPEND_CURRENT_QUEUE_A_MOMENT暂停队列持续重试保证有序 return ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT; } });注意顺序消息不要返回 RECONSUME_LATER会破坏消息顺序五、PushConsumer 核心要点总结答题可用业务逻辑同步执行在监听器内部正常执行返回CONSUME_SUCCESS异常返回RECONSUME_LATER禁止吞异常直接返回成功避免消息丢失禁止将消息提交外部线程异步处理消息重试超限进入死信队列必须监听死信并告警顺序消息使用独立监听器失败返回暂停队列保证时序。springboot ymlSpringBoot RocketMQ 防消费丢失完整方案包含yml 配置 Consumer 代码 避坑说明、死信监听、顺序消费者示例基于rocketmq-spring-boot-starter官方 starter1. Maven 依赖xmldependency groupIdorg.apache.rocketmq/groupId artifactIdrocketmq-spring-boot-starter/artifactId version2.2.3/version /dependencyapplication.yml 核心消费配置yamlrocketmq: name-server: 127.0.0.1:9876 consumer: group: ORDER_CONSUMER_GROUP # 批量消费消息数量默认1根据业务调整 consume-message-batch-max-size: 1 # 消费线程数量普通消息可调大顺序消息建议设1 consume-thread-number: 8 # 消息重试间隔(毫秒)慎用broker侧有重试阶梯策略一般不自定义 # suspend-time-millis: 3000 # 最大重试次数默认16次超限进入死信队列 max-reconsume-times: 16⭐关键底层机制starter 封装无需代码手动控制RocketMQ Push 模式监听器业务执行成功框架才通知 broker 提交 offset只要业务异常抛出、不捕获吞异常 → 自动触发重试3. 普通消息消费者【标准安全写法防止消息丢失】java运行import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.stereotype.Service; Slf4j Service RocketMQMessageListener( topic ORDER_TOPIC, consumerGroup ORDER_CONSUMER_GROUP ) public class OrderConsumer implements RocketMQListenerString { /** * 重点 * 1. 方法正常执行结束无异常 → 框架自动返回CONSUME_SUCCESSoffset推进 * 2. 方法抛出任意异常 → 框架自动返回RECONSUME_LATER消息重试 */ Override public void onMessage(String message) { log.info(收到订单消息{}, message); try { // 执行业务逻辑更新数据库、调用接口 handleBusiness(message); } catch (Exception e) { log.error(消费业务异常消息触发重试 message{}, message, e); // ✅ 关键点捕获异常后**重新抛出**不要吞异常 throw new RuntimeException(消费失败, e); } } private void handleBusiness(String message) { // 模拟业务处理 } }❌ 错误示范千万不要这样写消息丢失经典坑java运行Override public void onMessage(String message) { try { handleBusiness(message); } catch (Exception e) { log.error(异常,e); // 只打印日志不抛出异常框架认为消费成功offset提交消息永久丢失 } }❌ 第二个大坑内部异步线程处理java运行Override public void onMessage(String message) { // 丢进线程池异步执行当前方法立刻结束 threadPool.execute(()- handleBusiness(message)); // 方法正常退出offset已经提交异步线程报错无法重试 → 消息丢失 }结论业务逻辑必须同步在 onMessage 内部执行4. 顺序消息消费者SpringBoot 写法顺序消息必须使用RocketMQOrderListener失败不能普通重试需要暂停队列防止乱序java运行import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQOrderListener; import org.apache.rocketmq.spring.support.RocketMQConsumeMode; import org.springframework.stereotype.Service; import org.apache.rocketmq.client.consumer.ConsumeOrderlyStatus; import org.apache.rocketmq.client.producer.MessageOrderly; Slf4j Service RocketMQMessageListener( topic ORDER_SEQ_TOPIC, consumerGroup ORDER_SEQ_CONSUMER_GROUP, consumeMode RocketMQConsumeMode.ORDERLY // 开启顺序消费模式 ) public class OrderSeqConsumer implements RocketMQOrderListener { Override public ConsumeOrderlyStatus onMessage(MessageOrderly message) { String body new String(message.getBody()); try { log.info(顺序消息{}, body); handleSeqBusiness(body); return ConsumeOrderlyStatus.SUCCESS; } catch (Exception e) { log.error(顺序消息消费失败,e); // 暂停当前队列持续重试保证顺序 return ConsumeOrderlyStatus.SUSPEND_CURRENT_QUEUE_A_MOMENT; } } private void handleSeqBusiness(String body) { // 有序业务逻辑 } }5. 死信队列监听兜底方案防止消息进入死信无人处理消息重试达到 max-reconsume-times 后转入死信 topic%DLQ%消费组名java运行Slf4j Service RocketMQMessageListener( topic %DLQ%ORDER_CONSUMER_GROUP, consumerGroup ORDER_DLQ_CONSUMER_GROUP // 使用新消费组 ) public class OrderDeadLetterConsumer implements RocketMQListenerString { Override public void onMessage(String message) { log.error(【死信消息】多次消费失败请人工排查{}, message); // 此处触发钉钉/短信告警 } }注意死信消费者务必使用全新消费组不要和业务消费组共用6. 核心要点汇总面试直接背SpringBoot Starter 规则onMessage无异常正常返回→ 自动提交 offset消费成功onMessage抛出异常→ 自动通知 Broker 重试RECONSUME_LATER防消费丢失编码规范不能捕获异常后不抛出不能把业务丢到外部异步线程池配置层面max-reconsume-times设置最大重试次数超限进入死信务必监听死信并告警顺序消息单独使用RocketMQOrderListener失败返回SUSPEND_CURRENT_QUEUE_A_MOMENT7. 配套延伸消费幂等解决重复消费消费拿到消息先根据业务唯一 id 判断是否已处理防止重试带来重复问题。示例极简片段java运行Override public void onMessage(String message) { MessageExt msg (MessageExt) MessageAccessor.getOriginalMessage(message); String bizId msg.getKeys(); // 发送时存入keys作为唯一业务编号 // 1. 查询数据库/redis判断是否已处理 boolean exist checkProcessed(bizId); if(exist){ log.info(消息已处理直接跳过 bizId{},bizId); return; } // 2. 执行业务事务内写入处理记录 handleBusiness(message); // 3. 标记已处理 markProcessed(bizId); }6.如果消费者业务消费耗时很长会一直等待执行完成才提交offset么
返回列表