消息中间件进阶:RocketMQ 消息的生产与发送(二) 六、消息的生产与发送Producer 启动流程与路由获取我们先从 Producer 的启动开始。一个 Producer 启动的时候它可不是傻乎乎地直接就开始发消息了——它得先弄清楚“把消息发到哪里去”。Producer 启动的核心步骤创建 Producer 实例设置 ProducerGroup 名称同一个业务使用同一个 Group设置 NameServer 地址列表指定集群注册中心的地址配置其他参数发送超时时间、重试次数等调用 start 方法启动 Producer初始化 MQClientInstance客户端核心实例负责网络通信和路由管理启动 Netty 通信客户端建立与 NameServer 的连接从 NameServer 拉取所有 Topic 的路由信息启动定时任务每 30 秒更新一次路由信息Producer 就绪等待发送消息关键点解析ProducerGroup一个业务标识同一个业务里的 Producer 实例归属于同一个 Group。在事务消息中同一个 Group 的 Producer 可以互相回查事务状态。MQClientInstance这是客户端最核心的实例所有 Producer 和 Consumer 共用一个 MQClientInstance同一个 JVM 里。它负责管理与 NameServer 的连接管理与 Broker 的连接维护本地路由缓存统一的心跳和网络通信路由获取Producer 启动时会从 NameServer 拉取所有 Topic 的路由信息不仅仅是某个特定 Topic这样后续发消息时就不需要再等路由查询了。定时更新路由信息在本地缓存后会有一个定时任务每隔 30 秒 从 NameServer 拉取最新路由保证路由信息的时效性。消息发送的三种方式同步、异步、单向在《入门认知篇》里我们简单提过三种发送方式现在我们从底层实现的角度再来看一遍。同步发送Sync这是最常用、最简单的方式。发送消息后线程阻塞等待 Broker 返回响应收到响应后才继续执行。SendResult sendResult producer.send(msg);// 阻塞等待直到收到 Broker 的响应BrokerProducerBrokerProducer线程阻塞等待线程恢复⏱️ 整个过程中线程处于阻塞状态构造消息发送请求同步返回响应成功/失败根据结果处理业务逻辑适用场景关键业务——比如下单成功后的订单消息必须确认 Broker 收到了才能继续。异步发送Async发送消息后线程不阻塞立即返回。等 Broker 响应回来后通过回调函数来处理结果。producer.sendAsync(msg, new SendCallback() {Overridepublic void onSuccess(SendResult sendResult) {// 处理成功逻辑}Overridepublic void onException(Throwable e) {// 处理失败逻辑}});// 这里立即返回不阻塞BrokerProducer 异步线程Producer 主线程BrokerProducer 异步线程Producer 主线程主线程立即返回可以继续其他工作✅ 主线程未被阻塞⏱️ 响应处理在异步线程中构造消息提交异步发送任务发送请求返回响应回调 SendCallback适用场景对延迟敏感但需要知道结果——比如前端请求发消息不能阻塞用户操作但失败时需要通知用户。单向发送Oneway只发不管连响应都不等。最轻量、最快但最不可靠。producer.sendOneway(msg);// 不管结果直接继续BrokerProducerBrokerProducer发完就走不等待任何响应 最快但没有任何可靠性保证构造消息发送请求单向继续执行后续代码适用场景日志上报、监控数据等——丢了就丢了业务上能接受。 小贴士同步发送和异步发送虽然看起来差异很大但在底层它们都使用了异步网络通信Netty。同步发送只是在异步通信的基础上用 CountDownLatch 等同步工具做了一个“阻塞等待”的封装。本质上RocketMQ 的网络通信模型是全异步的。消息发送的负载均衡策略当一个 Topic 有多个 MessageQueue分布在不同的 Broker 上时Producer 如何选择将消息发到哪个 Queue这就涉及负载均衡策略。RocketMQ 默认提供了两种内置策略你也可以自定义实现策略一轮询Round Robin—— 默认策略依次轮流选择 Queue确保消息均匀分布。在顺序消息场景下会结合消息 Key 做哈希保证同一个 Key 的消息落到同一个 Queue。MessageQueue消息序列第 1 条第 2 条第 3 条第 4 条第 5 条第 6 条消息 1消息 2消息 3消息 4消息 5消息 6Queue 0Queue 1Queue 2按顺序轮流分配1→Q1, 2→Q2, 3→Q3, 4→Q1…策略二一致性哈希Consistent Hash根据消息 Key 的哈希值通过哈希环来决定消息发到哪个 Queue。保证同一个 Key 的消息始终落在同一个 Queue 上对于顺序消息特别重要。消息 Key哈希值映射到环上哈希值映射到环上哈希值映射到环上顺时针查找第一个节点顺时针查找第一个节点顺时针查找第一个节点一致性哈希环节点 1Queue 0节点 2Queue 1节点 3Queue 2节点 4Queue 3节点 5Queue 0节点 6Queue 1order_123order_456order_789位置 A位置 B位置 C同一 Key 始终落到同一 Queue新增节点时只影响局部 Key 的重新分配自定义策略实现 MessageQueueSelector 接口根据自己的业务逻辑选择 Queue比如按订单 ID 取模。消息发送的重试机制与故障规避消息发送不是总能一次成功的。网络抖动、Broker 故障、磁盘满了……各种原因都可能导致发送失败。RocketMQ 有一套完善的重试和故障规避机制。重试机制是否否是是否是否发送消息第 1 次尝试成功返回成功是否为可重试的异常直接返回失败不重试第 2 次尝试等待 50ms成功… 最多重试 retryTimesWhenSendFailed 次最后一次成功重试的关键配置参数 默认值 说明retryTimesWhenSendFailed 2 同步发送失败时的重试次数不含第一次retryTimesWhenSendAsyncFailed 2 异步发送失败时的重试次数不含第一次retryAnotherBrokerWhenNotStoreOK false 同步模式下如果写入失败是否换 Broker 重试哪些异常会触发重试✅ 可重试网络超时、连接异常、Broker 返回 SYSTEM_BUSY、SERVICE_NOT_AVAILABLE 等❌ 不重试消息体超限4MB、Topic 不存在、消息格式错误等重试也没用故障规避机制故障延迟机制这是 RocketMQ 一个非常巧妙的设计。它的核心思想是某个 Broker 如果发送失败了在接下来的一段时间内尽量不要再往这个 Broker 上发消息。Broker B 故障后不发送ProducerBroker B ❌Broker A ✅Broker C ✅在故障延迟时间内Broker B 被加入黑名单自动规避正常情况ProducerBroker ABroker BBroker C具体工作方式每次发送失败时记录当前系统时间和失败的 Broker在接下来的 sendLatencyFaultEnable 时间内默认 30 秒该 Broker 会被加入黑名单Producer 选择 Queue 时会过滤掉黑名单中的 Broker30 秒后自动恢复重新尝试往该 Broker 发送这个机制可以避免 Producer 持续向故障 Broker 发送消息从而减少不必要的超时等待提升整体的发送成功率。消息的路由队列选择算法消息从 Producer 到 Broker 的过程实际上经历了两层选择需要发送消息根据 Topic 获取本地缓存的路由信息从路由信息中获取所有可用的MessageQueue 列表过滤掉处于故障延迟期的 Broker 上的 Queue根据负载均衡策略选择一个 MessageQueue根据选中的 Queue 获取对应的 Broker 地址构造消息发送请求通过 Netty 发送到目标 Broker等待响应完整的选择决策流程包含以下因素路由信息从本地缓存中获取 Topic 对应的所有 Broker 和 Queue 信息Broker 可用性过滤掉 NameServer 已剔除的、或故障延迟期内的 Broker负载均衡策略轮询或一致性哈希Queue 状态如果某个 Queue 满了可能被临时剔除虽然 RocketMQ 很少出现这种情况因为文件是滚动的消息钩子MessageHook的使用MessageHook 是什么MessageHook 是 RocketMQ 提供的一个扩展点允许开发者在消息发送的前后插入自定义逻辑。典型用途消息发送钩子发送前executeBefore消息发送发送后executeAfter 链路追踪埋点记录耗时 权限校验检查是否有发送权限 流量染色记录调用方信息 审计日志记录消息发送记录使用方式// 实现 MessageHook 接口public class CustomSendHook implements MessageHook {Overridepublic String hookName() {return “CustomSendHook”;}Override public void executeBefore(SendMessageContext context) { // 发送前记录开始时间、校验权限等 System.out.println(开始发送Topic: context.getMessage().getTopic()); } Override public void executeAfter(SendMessageContext context) { // 发送后记录耗时、审计日志等 System.out.println(发送结束耗时: context.getCostTime()); }}// 注册钩子producer.getDefaultMQProducerImpl().registerMessageHook(new CustomSendHook());典型应用场景全链路追踪在发送前生成 TraceId发送后记录日志性能监控统计消息发送的耗时分布流量控制在发送前进行限流或权限校验消息发送的超时处理与异常场景消息发送过程中可能遇到各种异常我们来看 RcoketMQ 的异常处理是如何设计的超时处理同步发送有一个超时时间默认 3 秒。如果在这个时间内没收到 Broker 的响应Producer 会抛出 RemotingTimeoutException。发送请求启动超时计时器未超时收到响应3秒到了ProducerBrokerTimer 3s等待响应✅ 返回成功❌ 抛出超时异常⚠️ 注意超时异常发生时你无法确定 Broker 是否收到了消息。可能消息根本没到 Broker也可能 Broker 已经写入了只是网络响应慢了。所以遇到超时异常时业务方需要根据消息 Key 做幂等处理避免重复消费。常见异常场景异常类型 可能原因 是否重试 处理建议RemotingTimeoutException 网络慢或 Broker 响应慢 ✅ 是 适当增大超时时间或考虑异步发送RemotingConnectException Broker 连接不上 ✅ 是 检查 Broker 是否存活故障规避机制会处理MQClientException Topic 不存在或消息体超限 ❌ 否 检查 Topic 配置和消息大小MQBrokerException Broker 返回业务错误 看情况 根据错误码判断是否可重试RemotingSendRequestException 网络发送失败 ✅ 是 检查网络连接消息发送的请求响应流程Netty 通信RocketMQ 的底层网络通信基于 Netty 实现。一条消息从 Producer 到 Broker经过的请求响应流程是标准化的存储层Broker 业务处理器Netty Server(Broker 端口 10911)Netty ClientProducer存储层Broker 业务处理器Netty Server(Broker 端口 10911)Netty ClientProducerBroker 接收请求调用发送 API通过 Netty Channel 发送请求码: SEND_MESSAGE根据请求码路由到SendMessageProcessor解析请求头Topic、Queue、消息体等调用存储层写入 CommitLog返回写入结果和偏移量构造响应体通过 Netty 返回响应响应返回客户端回调/唤醒等待线程处理发送结果关键设计细节请求码路由Netty 收到请求后根据请求码SEND_MESSAGE、PULL_MESSAGE 等将请求分发到不同的处理器单线程模型RocketMQ 的 Broker 使用单线程处理消息写入通过 putMessageLock保证 CommitLog 的顺序写入异步化Broker 处理完写入后通过 Netty 的异步回调机制返回响应不阻塞 Netty 的 I/O 线程超时控制客户端的超时是基于 Netty 的 ResponseFuture 机制实现的批量消息发送的实现与限制批量消息发送允许将多条消息打包成一条网络请求发送从而显著提升吞吐量。批量发送1 次请求消息1批量消息消息1消息2消息3消息2消息3Broker1 次网络请求单条发送请求1请求2请求3消息1Broker消息2消息33 次网络请求批量发送的使用限制同一 Topic批量消息必须属于同一个 Topic同一 Broker批量消息的 Queue 列表必须落在同一个 Broker 上因为一次网络请求只能发给一个 Broker大小限制批量消息的总大小不能超过 maxMessageSize默认 4MB不能有延迟/事务批量消息不支持延迟消息和事务消息使用示例List messages new ArrayList();messages.add(new Message(“topic_test”, “消息1”.getBytes()));messages.add(new Message(“topic_test”, “消息2”.getBytes()));messages.add(new Message(“topic_test”, “消息3”.getBytes()));SendResult sendResult producer.send(messages); // 批量发送延迟消息的实现原理延迟等级RocketMQ 的延迟消息非常有意思——它不是精确到秒的任意延迟而是基于预定义的延迟等级来实现的。延迟等级RocketMQ 内置了 18 个延迟等级messageDelayLevel每个等级对应一个固定的延迟时间等级 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18延迟 1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h底层实现原理延迟消息处理存储层发送延迟消息设置 delayTimeLevel3计算投递时间消息写入 CommitLog异步构建索引写入 SCHEDULE_TOPIC_XXXX定时任务扫描到达投递时间消息被 Consumer 拉取ProducerBroker投递时间 当前时间 10sCommitLogConsumeQueue延迟 Topic 的 Queue按延迟等级分队列后台定时任务每秒扫描一次重新投递到目标 TopicConsumer 消费实现逻辑原理层面Producer 发送消息时设置 delayTimeLevel 参数比如 3 表示延迟 10 秒Broker 收到消息后不直接把消息写入目标 Topic 的 ConsumeQueue而是先写入一个内部的 SCHEDULE_TOPIC_XXXX这个 Topic 有 18 个 Queue对应 18 个延迟等级Broker 有一个后台定时任务ScheduleMessageService每秒扫描一次检查哪些延迟消息到了投递时间到时间的消息被重新写入目标 Topic 的 CommitLog 和 ConsumeQueueConsumer 正常消费——对消费者来说这个消息就像是刚刚发送的一样 小贴士这种设计的核心是——延迟消息没有特殊的存储格式就是普通消息加上一个“定时投递”的机制。所有的延迟逻辑都在 Broker 内部完成Consumer 完全无感知。所以如果你需要更灵活的定时精度也可以基于这个机制自行扩展。5.x 版本的变化RocketMQ 5.x 对延迟消息做了增强支持 任意时间点的定时消息不再是只有 18 个等级通过 Timing Wheel 算法实现了更灵活的定时投递。但核心原理和上面的流程是一样的只是内部的定时机制从“固定等级扫描”升级到了“任意时间点调度”。消息轨迹Message Trace的开启与使用消息轨迹是什么消息轨迹记录了一条消息从生产到消费的完整链路信息。你可以通过消息轨迹清楚地知道这条消息是谁发的、什么时候发的、发到了哪个 Broker这条消息被谁消费了、什么时候消费的、消费结果如何轨迹数据消息轨迹Producer 发送Broker 存储Consumer 消费发送时间发送方 IP消息 Key存储位置消费时间消费方 IP消费状态开启方式方式一通过配置文件开启推荐在 broker.conf 中配置traceTopicEnabletruemsgTraceTopicNameRMQ_SYS_TRACE_TOPIC方式二通过客户端显式开启// Producer 端开启消息轨迹DefaultMQProducer producer new DefaultMQProducer(“producer_group”);producer.setEnableMsgTrace(true);producer.setMsgTraceTopic(“RMQ_SYS_TRACE_TOPIC”);消息轨迹的实现原理消息发送/消费通过 MessageHook 和 ConsumeHook拦截关键事件构建轨迹数据发送时间、IP、状态等将轨迹数据作为普通消息发送到 Trace TopicBroker 存储轨迹消息用户通过控制台或 API查询轨迹信息核心设计无侵入通过 Hook 机制实现业务代码不需要改造异步上报轨迹数据的发送是异步的不影响主消息发送的性能独立存储轨迹数据存储在独立的 RMQ_SYS_TRACE_TOPIC 中与业务数据隔离可查询通过 RocketMQ 的 Dashboard 或 API 可以按消息 ID/Key 查询轨迹实际使用中的注意事项开启消息轨迹会增加少量性能开销需要额外发送轨迹消息生产环境根据实际需求开启轨迹消息本身也占用存储空间需要适当设置 fileReservedTime对于超高吞吐量的场景可以采样记录轨迹而不是每条都记录小结这篇文章我们完整走通了消息从生产到发送的全链路通过 8 张流程图搞清楚了Producer 启动时如何初始化并获取路由信息三种发送方式同步、异步、单向的底层实现和适用场景负载均衡策略轮询、一致性哈希如何选择目标 Queue重试机制和故障规避机制如何在发送失败时保证高可用