Apache RocketMQ DefaultMQProducer 完全指南:核心字段、构造方法、消息发送 API 与源码级实践详解 📅 发布时间:2026/9/20 19:38:55 👁 浏览次数: 消息队列后端微服务流处理【免费下载链接】rocketmqApache RocketMQ is a cloud native messaging and streaming platform, making it simple to build event-driven applications.项目地址https://gitcode.com/gh_mirrors/ro/rocketmq点击查看免费下载DefaultMQProducer是 Apache RocketMQ Java 客户端中应用投递消息的统一入口类它开箱即用封装了同步、异步、Oneway 三种发送模式均支持批量并对外暴露了从创建 Topic 到事务消息、消息轨迹在内的一整套生产者 API。本文以 DefaultMQProducer 官方 API 参考文档 为骨架结合 DefaultMQProducer.java 与 DefaultMQProducerImpl.java 的源码实现逐字段、逐方法讲解每个配置参数与每个 send 变体的真实行为、默认值与适用场景帮助读者在阅读完本文后能够根据业务需求正确选型、配置并使用生产者 API。类简介生产者的核心入口DefaultMQProducer的类声明如下public class DefaultMQProducer extends ClientConfig implements MQProducer该类的核心职责与设计要点对应源码注释见 DefaultMQProducer.java#L58-L70应用投递消息的入口通过无参构造即可快速创建生产者负责消息的发送支持同步/异步/Oneway 三种发送方式且这三种方式均支持批量发送。参数可调可以通过该类提供的 getter/setter 方法调整发送者的参数同时该类提供了多个send方法每个方法语义略有不同使用前务必仔细了解其意图。线程安全在配置并启动完成后该类可以在多个线程间安全共享。源码 javadoc 明确标注了这一点这也意味着一个生产者实例可以被整个应用的多个业务线程复用。下面给出一个最基础的生产者示例与官方示例代码一致public class Producer { public static void main(String[] args) throws MQClientException { // 创建指定分组名的生产者 DefaultMQProducer producer new DefaultMQProducer(ProducerGroupName); // 启动生产者 producer.start(); for (int i 0; i 128; i) try { // 构建消息 Message msg new Message(TopicTest, TagA, OrderID188, Hello world.getBytes(RemotingHelper.DEFAULT_CHARSET)); // 同步发送 SendResult sendResult producer.send(msg); // 打印发送结果 System.out.printf(%s%n, sendResult); } catch (Exception e) { e.printStackTrace(); } producer.shutdown(); } }该示例的完整版可在仓库的 example/simple 目录下找到其中还包含同步、异步、Oneway、批量发送、顺序消息等更丰富的发送示例。字段摘要生产者的关键配置属性DefaultMQProducer的核心字段及其含义如下表所示类型字段名称描述DefaultMQProducerImpldefaultMQProducerImpl生产者的内部默认实现StringproducerGroup生产者分组StringcreateTopicKey在发送消息时自动创建服务器不存在的 topicintdefaultTopicQueueNums创建 topic 时默认的队列数量intsendMsgTimeout发送消息的超时时间intcompressMsgBodyOverHowmuch压缩消息体的阈值intretryTimesWhenSendFailed同步模式下内部尝试发送消息的最大次数intretryTimesWhenSendAsyncFailed异步模式下内部尝试发送消息的最大次数booleanretryAnotherBrokerWhenNotStoreOK是否在内部发送失败时重试另一个 brokerintmaxMessageSize消息体的最大长度TraceDispatchertraceDispatcher基于 RPCHook 实现的消息轨迹插件字段详细信息producerGroup生产者分组private String producerGroup;含义生产者的分组名称。相同的分组名称表明生产者实例在概念上归属于同一分组。这对事务消息十分重要如果原始生产者在事务之后崩溃那么 broker 可以联系同一生产者分组的不同生产者实例来提交或回滚事务。对于非事务消息只要保证每个进程内唯一即可。默认值DEFAULT_PRODUCER该常量定义于 MixAll.java#L76MixAll.DEFAULT_PRODUCER_GROUP。命名约束由数字、字母、下划线、横杠-、竖线|或百分号%组成不能为空长度不能超过 255。defaultMQProducerImpl内部默认实现protected final transient DefaultMQProducerImpl defaultMQProducerImpl;生产者的内部默认实现在构造生产者时自动初始化见 DefaultMQProducer.java#L316提供了本类绝大部分方法的内部实现包括默认 topic 路由选择、消息压缩、发送超时控制与失败重试等逻辑。createTopicKey自动创建 topic 的 Keyprivate String createTopicKey TopicValidator.AUTO_CREATE_TOPIC_KEY_TOPIC;含义在发送消息时自动创建服务器不存在的 topic 需要指定 Key该 Key 可用于配置发送消息所在 topic 的默认路由。默认值TBW102即AUTO_CREATE_TOPIC_KEY_TOPIC。建议仅用于测试或 demo生产环境下不建议打开自动创建配置。defaultTopicQueueNums默认队列数private volatile int defaultTopicQueueNums 4;创建 topic 时默认的队列数量默认值为 4。该字段被声明为volatile可在运行时安全调整。sendMsgTimeout发送超时private int sendMsgTimeout 3000;默认值3000单位毫秒。建议不建议修改该值该值应该与 broker 配置中的sendTimeout保持一致。发送超时时可临时修改该值但根治方案是提升 broker 集群的 TPS 以解决超时问题。compressMsgBodyOverHowmuch压缩阈值private int compressMsgBodyOverHowmuch 1024 * 4;默认值1024 * 44KB单位字节。大于 4K 的消息体将默认进行压缩。底层实现压缩动作发生在发送主流程中见 DefaultMQProducerImpl.java#L1112-L1136 的tryToCompressMessage方法当消息体字节长度 compressMsgBodyOverHowmuch时调用defaultMQProducer.getCompressor().compress(body, compressLevel)对消息体进行压缩成功后替换原消息体。注意批量消息MessageBatch当前不支持压缩。压缩参数可通过DefaultMQProducerImpl.setZipCompressLevel方法设置压缩率默认为 5可选范围 [0, 9]可通过DefaultMQProducerImpl.tryToCompressMessage方法测试出 compressLevel 与 compressMsgBodyOverHowmuch 的最佳组合值。从源码看压缩类型与级别还支持通过系统属性rocketmq.message.compress.level与rocketmq.message.compress.type覆盖默认 ZLIB见 DefaultMQProducer.java#L213-L223。retryTimesWhenSendFailed / retryTimesWhenSendAsyncFailed重试次数private int retryTimesWhenSendFailed 2; private int retryTimesWhenSendAsyncFailed 2;含义分别在同步模式、异步模式下在返回发送失败之前内部尝试重新发送消息的最大次数。默认值均为 2即默认情况下一条消息最多会被投递 3 次。注意在极端情况下这可能会导致消息的重复需要应用层自行做幂等处理。retryAnotherBrokerWhenNotStoreOK跨 broker 重试private boolean retryAnotherBrokerWhenNotStoreOK false;含义同步模式下消息保存失败时是否重试另一个 broker。默认值false。注意此配置关闭时非投递时产生异常的情况下会忽略retryTimesWhenSendFailed配置。maxMessageSize最大消息体private int maxMessageSize 1024 * 1024 * 4;默认值1024 * 1024 * 44MB单位字节。含义当消息体的字节数超过 maxMessageSize 时发送失败。该校验在 Validators.checkMessage 中执行——发送路径如 DefaultMQProducerImpl.java#L688在真正投递前都会先调用该校验方法超出大小上限的消息会直接抛出MQClientException。traceDispatcher消息轨迹private TraceDispatcher traceDispatcher null;含义开启消息轨迹后该类通过 hook 的方式把消息生产者、消息存储的 broker 和消费者消费消息的信息像链路一样记录下来。初始化时机在构造生产者时根据构造入参enableMsgTrace决定是否创建该对象。在start()方法中若开启轨迹会创建AsyncTraceDispatcher并注册SendMessageTraceHookImpl、EndTransactionTraceHookImpl等 hook见 DefaultMQProducer.java#L380-L405轨迹数据默认写入RMQ_SYS_TRACE_TOPIC可通过customizedTraceTopic自定义。构造方法七种创建方式方法名称方法描述DefaultMQProducer()由默认参数值创建一个生产者DefaultMQProducer(final String producerGroup)使用指定的分组名创建一个生产者DefaultMQProducer(final String producerGroup, boolean enableMsgTrace)使用指定的分组名创建一个生产者并设置是否开启消息轨迹DefaultMQProducer(final String producerGroup, boolean enableMsgTrace, final String customizedTraceTopic)使用指定的分组名创建一个生产者并设置是否开启消息轨迹及追踪 topic 的名称DefaultMQProducer(RPCHook rpcHook)使用指定的 hook 创建一个生产者DefaultMQProducer(final String producerGroup, RPCHook rpcHook)使用指定的分组名及自定义 hook 创建一个生产者DefaultMQProducer(final String producerGroup, RPCHook rpcHook, boolean enableMsgTrace, final String customizedTraceTopic)使用指定的分组名及自定义 hook 创建一个生产者并设置是否开启消息轨迹及追踪 topic 的名称从源码看无参构造内部调用this(MixAll.DEFAULT_PRODUCER_GROUP)最终所有构造器都会汇聚到最完整的构造器DefaultMQProducer.java#L309-L317在其中初始化defaultMQProducerImpl new DefaultMQProducerImpl(this, rpcHook)。仓库中还存在带topics列表用于事务生产者初始化路由以及带 namespace 的构造器标注Deprecated说明 namespace 能力已并入ClientConfig。各构造器入参说明汇总参数名类型是否必须缺省值描述producerGroupString是DEFAULT_PRODUCER生产者的分组名称rpcHookRPCHook否null每个远程命令执行后会回调 rpcHookenableMsgTraceboolean是false是否开启消息轨迹customizedTraceTopicString否RMQ_SYS_TRACE_TOPIC消息轨迹 topic 的名称使用方法核心 API 详解生命周期管理start 与 shutdownstart()public void start() throws MQClientException启动生产者实例。在发送或查询消息之前必须调用此方法。它执行了许多内部初始化比如检查配置、与 namesrv 建立连接、启动一系列心跳等定时任务等。从源码看DefaultMQProducer.java#L373-L406start()会先给 producerGroup 补全 namespace然后调用defaultMQProducerImpl.start()完成客户端实例的启动若开启了消息轨迹还会在此创建AsyncTraceDispatcher并注册发送 hook、启动轨迹分发器。shutdown()public void shutdown()关闭当前生产者实例并释放相关资源。会依次关闭内部实现defaultMQProducerImpl.shutdown()、自动批量累加器produceAccumulator与轨迹分发器traceDispatcher见 DefaultMQProducer.java#L411-L420。Topic 与消息队列管理createTopic两个重载public void createTopic(String key, String newTopic, int queueNum) public void createTopic(String key, String newTopic, int queueNum, int topicSysFlag)在 broker 上创建一个 topic。参数说明参数名类型是否必须默认值值范围说明keyString是访问密钥newTopicString是新建 topic 的名称。由数字、字母、下划线_、横杠-、竖线|或百分号%组成长度小于 255不能为 TBW102 或空queueNumint是0(0, maxIntValue]topic 的队列数量topicSysFlagint是0保留字段暂未使用fetchPublishMessageQueuespublic ListMessageQueue fetchPublishMessageQueues(String topic)获取 topic 的消息队列即生产者可投递的队列列表底层委托给defaultMQProducerImpl.fetchPublishMessageQueues(withNamespace(topic))。消息查询类 API方法签名功能关键返回值/异常long earliestMsgStoreTime(MessageQueue mq)查询最早的消息存储时间单位毫秒long maxOffset(MessageQueue mq)查询给定消息队列的最大物理偏移量最大物理偏移量long minOffset(MessageQueue mq)查询给定消息队列的最小物理偏移量最小物理偏移量long searchOffset(MessageQueue mq, long timestamp)查找指定时间的消息队列的物理偏移量指定时间的物理偏移量timestamp 单位毫秒QueryResult queryMessage(String topic, String key, int maxNum, long begin, long end)按关键字查询消息返回查询到的消息集合key 可空begin/end 为毫秒时间戳MessageExt viewMessage(String offsetMsgId)根据给定的 msgId 查询消息返回MessageExt包含 topic 名称、消息体、消息 ID、消费次数、生产者 host 等信息MessageExt viewMessage(String topic, String msgId)根据给定的 msgId 查询消息并指定 topic同上这些查询类方法在生产者状态非 Running、未找到 broker、broker 返回失败、网络异常或线程中断时均会抛出MQClientException等异常viewMessage在 msgId 非法时也会抛出MQClientException。消息发送核心 API同步、异步、OnewayDefaultMQProducer提供了大量send方法按发送模式可划分为三大类每一类都有一组重载。选型要点如下同步发送SendResult send(...)同步发送仅在发送过程完全完成后返回是默认推荐的可靠性发送方式。常见变体包括方法签名特点SendResult send(Message msg)同步单条发送内部失败重试次数由retryTimesWhenSendFailed控制未指定队列时默认轮询策略SendResult send(Message msg, long timeout)指定超时时间超时抛出RemotingTooMuchRequestExceptionSendResult send(Message msg, MessageQueue mq)向指定的消息队列同步发送单条消息SendResult send(Message msg, MessageQueue mq, long timeout)指定队列 指定超时SendResult send(Message msg, MessageQueueSelector selector, Object arg)通过队列选择器计算目标队列此方式发送失败内部不会重试SendResult send(Message msg, MessageQueueSelector selector, Object arg, long timeout)选择器 超时SendResult send(CollectionMessage msgs)同步批量发送集合内消息必须属于同一个 topic默认轮询队列SendResult send(CollectionMessage msgs, long timeout)批量 超时SendResult send(CollectionMessage msgs, MessageQueue messageQueue)向指定队列批量发送SendResult send(CollectionMessage msgs, MessageQueue messageQueue, long timeout)指定队列 超时MessageQueueSelector的典型用途通过自实现该接口将某一类消息发送至固定队列。例如将同一个订单的状态变更消息始终投递到固定队列从而配合消费端实现顺序消费。仓库示例 example/ordermessage 中提供了基于MessageQueueSelector的完整顺序消息发送示例。异步发送void send(..., SendCallback sendCallback, ...)异步发送调用后立即返回并在发送成功或异常时回调sendCallback方法签名特点void send(Message msg, SendCallback sendCallback)异步发送单条消息内部重试次数由retryTimesWhenSendAsyncFailed控制void send(Message msg, SendCallback sendCallback, long timeout)异步 超时超时后回调收到RemotingTooMuchRequestExceptionvoid send(Message msg, MessageQueue mq, SendCallback sendCallback)指定队列异步发送void send(Message msg, MessageQueue mq, SendCallback sendCallback, long timeout)指定队列 超时void send(Message msg, MessageQueueSelector selector, Object arg, SendCallback sendCallback)选择器 异步void send(Message msg, MessageQueueSelector selector, Object arg, SendCallback sendCallback, long timeout)选择器 异步 超时重要提醒异步发送时sendCallback参数不能为 null否则在回调时会抛出NullPointerException。Oneway 发送void sendOneway(...)public void sendOneway(Message msg) public void sendOneway(Message msg, MessageQueue mq) public void sendOneway(Message msg, MessageQueueSelector selector, Object arg)以 oneway 形式发送消息时broker 不会响应任何执行结果行为与 UDP 类似——具有最大的吞吐量但消息可能会丢失。适用于消息量大、追求高吞吐量且允许消息丢失的场景。源码中sendOneway通过CommunicationMode.ONEWAY模式直接投递见 DefaultMQProducerImpl.java#L1212-L1218。发送路径上的统一校验从源码可以看到所有发送入口在投递前都会先执行Validators.checkMessage(msg, this.defaultMQProducer)见 DefaultMQProducerImpl.java#L688、L745、L1232 等校验消息 topic 合法性、消息体大小maxMessageSize等随后进入压缩、选队列、发送的主流程。因此maxMessageSize、compressMsgBodyOverHowmuch等字段的取值会直接作用于每一次发送。事务消息sendMessageInTransactionpublic TransactionSendResult sendMessageInTransaction(Message msg, LocalTransactionExecuter tranExecuter, final Object arg) public TransactionSendResult sendMessageInTransaction(Message msg, final Object arg)该类不做默认实现始终抛出RuntimeException异常。第一个重载中的LocalTransactionExecuter已过期将在 5.0.0 版本中移除请勿使用该方法。事务消息应使用TransactionMQProducer类实现。仓库中 client/src/main/java/org/apache/rocketmq/client/producer/TransactionMQProducer.java 提供了完整实现example/transaction 下有事务消息的完整示例包含本地事务执行器与事务回查监听器的实现。发送异常速查同步与批量发送可能抛出的异常包括MQClientExceptionbroker 不存在或未找到、namesrv 地址为空、未找到 topic 的路由信息等客户端异常RemotingException网络异常MQBrokerExceptionbroker 发生错误InterruptedException发送线程中断RemotingTooMuchRequestException发送超时超时时间由sendMsgTimeout或各方法入参timeout指定。异步与 Oneway 发送的异常列表与同步类似但不包含MQBrokerException。生产实践建议综合上文字段与方法的语义在实际生产使用中建议遵循以下要点生产者复用与生命周期DefaultMQProducer是线程安全的一个 JVM 进程内建议复用少量生产者实例遵循start()→ 多次send→shutdown()的生命周期不要在每次发送消息时都新建生产者。可靠性与吞吐的权衡核心业务使用同步发送send(msg)必要时指定超时与MessageQueueSelector保证顺序可容忍少量丢失的高吞吐场景使用 Oneway需要快速返回且能接受异步回调的场景使用异步发送。三种模式的重试次数分别由retryTimesWhenSendFailed与retryTimesWhenSendAsyncFailed控制注意重试可能带来消息重复应用层需做好幂等。参数校准sendMsgTimeout应与 broker 的sendTimeout保持一致maxMessageSize默认 4MB与 broker 侧允许的最大消息大小需匹配开启自动创建 topiccreateTopicKey仅限测试环境。消息轨迹通过new DefaultMQProducer(group, true)或new DefaultMQProducer(group, true, 自定义轨迹Topic)开启消息轨迹即可在RMQ_SYS_TRACE_TOPIC中观测消息从生产、存储到消费的完整链路。扩展阅读DefaultMQProducer 源码字段默认值、构造器聚合逻辑与 start/shutdown 实现DefaultMQProducerImpl 源码发送主流程、消息压缩tryToCompressMessage、路由选择与重试逻辑Validators.checkMessage发送前的消息合法性校验生产者示例合集example/simple同步/异步/Oneway/批量、example/ordermessage顺序消息、example/transaction事务消息消费端对应 APIAPI_Reference_DefaultPullConsumer.md赞分享消息队列后端微服务流处理【免费下载链接】rocketmqApache RocketMQ is a cloud native messaging and streaming platform, making it simple to build event-driven applications.项目地址https://gitcode.com/gh_mirrors/ro/rocketmq点击查看免费下载相关推荐Cornell McRay核心技术解析如何集成kajiya、physx-rs和dolly三大引擎Cornell McRay核心技术解析如何集成kajiya、physx rs和dolly三大引擎 Cornell McRay是一个基于Rust语言开发的快速游消息队列后端微服务流处理Apache RocketMQ批量消息发送实践指南Apache RocketMQ批量消息发送实践指南 批量消息发送概述 在分布式消息系统中批量消息发送是一种重要的性能优化手段。Apache RocketMQ作消息队列后端微服务流处理三步免费解锁Wand专业版告别游戏时间限制的终极方案三步免费解锁Wand专业版告别游戏时间限制的终极方案 还在为Wand原WeMod的2小时限制而烦恼吗想要免费享受AI游戏指南、自定义配置保存等专业版特权消息队列后端微服务流处理上一篇Pinpoint Agent插件开发脚手架从零到一的完整指南 下一篇UUID与数据库索引ramsey/uuid提升查询性能的终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考