Apache RocketMQ 批量消息发送实战指南:4MiB 限制、ListSplitter 拆分与源码级原理

发布时间:2026/9/20 23:47:22

Apache RocketMQ 批量消息发送实战指南:4MiB 限制、ListSplitter 拆分与源码级原理 消息队列后端微服务流处理【免费下载链接】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点击查看免费下载批量发送是 Apache RocketMQ 生产端提升发送效率、抬高系统吞吐量的核心技术手段它把多条消息封装成一次网络请求提交给 Broker显著减少 RPC 往返次数。本文以官方文档《批量消息发送》为主线结合仓库中的生产端示例 SimpleBatchProducer.java、SplitBatchProducer.java 与客户端源码讲清批量发送的适用条件、4MiB 上限的由来、大消息拆分算法以及批量消息在客户端内部的编码与校验链路读完即可在自己的生产者代码中落地实现。一、批量消息的适用条件与核心约束在动手写代码之前必须先理解 RocketMQ 对同一批消息的硬性要求。这些约束不是文档建议而是客户端源码层面的强制校验违反会直接抛出UnsupportedOperationException同一批消息的 topic 必须一致批量消息在编码阶段会被拼装为一条复合消息其外层只能携带一个 topic同一批消息的waitStoreMsgOK属性必须一致该属性决定发送时是否等待 Broker 落盘确认同步刷盘/异步刷盘语义一批内混用会导致语义不明确批量消息不支持延迟消息无论是delayTimeLevel延迟级别、delayTimeMs、delayTimeSec还是deliverTimeMs只要设置了任何一个延迟相关属性批量发送都会被拒绝批量消息不支持重试主题Retry Topictopic 以重试组前缀开头时同样被拒绝。上述校验可以在 MessageBatch.java 的generateFromList方法中看到完整实现if (message.getDelayTimeLevel() 0 || message.getDelayTimeMs() 0 || message.getDelayTimeSec() 0 || message.getDeliverTimeMs() 0) { throw new UnsupportedOperationException(Delayed messages are not supported for batching); } if (message.getTopic().startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) { throw new UnsupportedOperationException(Retry Group is not supported for batching); } if (!first.getTopic().equals(message.getTopic())) { throw new UnsupportedOperationException(The topic of the messages in one batch should be the same); } if (first.isWaitStoreMsgOK() ! message.isWaitStoreMsgOK()) { throw new UnsupportedOperationException(The waitStoreMsgOK of the messages in one batch should the same); }除此之外还有一个硬性体积限制单次批量发送最多 4MiB。如果需要发送更大的消息官方建议将大消息拆分成多个不超过 1MiB 的小消息再分批发送。4MiB 限制的源码出处4MiB 并非随意约定而是生产端DefaultMQProducer的默认maxMessageSize/** * Maximum allowed message body size in bytes. */ private int maxMessageSize 1024 * 1024 * 4; // 4M参见 DefaultMQProducer.java。该值可通过producer.setMaxMessageSize(int)调整用于控制单条消息含批量复合消息允许携带的最大体积。实际运行中Broker 端的maxMessageSize配置会构成最终约束生产端默认值与 Broker 默认配置保持一致均为 4MiB因此建议不要在生产端与 Broker 端分别做不一致的放大调整否则可能出现客户端认为合法、Broker 拒收的情况。二、发送不超过 4MiB 的批量消息如果你一次发送的总数据量不超过 4MiB直接使用批处理 API 即可非常简单。以仓库中的官方示例 SimpleBatchProducer.java 为蓝本package org.apache.rocketmq.example.batch; import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.List; import org.apache.rocketmq.client.producer.DefaultMQProducer; import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.common.message.Message; public class SimpleBatchProducer { public static final String PRODUCER_GROUP BatchProducerGroupName; public static final String DEFAULT_NAMESRVADDR 127.0.0.1:9876; public static final String TOPIC BatchTest; public static final String TAG Tag; public static void main(String[] args) throws Exception { DefaultMQProducer producer new DefaultMQProducer(PRODUCER_GROUP); // 本地调试时取消注释并将地址改为你的 NameServer 地址 // producer.setNamesrvAddr(DEFAULT_NAMESRVADDR); producer.start(); String topic BatchTest; ListMessage messages new ArrayList(); messages.add(new Message(topic, TagA, OrderID001, Hello world 0.getBytes(StandardCharsets.UTF_8))); messages.add(new Message(topic, TagA, OrderID002, Hello world 1.getBytes(StandardCharsets.UTF_8))); messages.add(new Message(topic, TagA, OrderID003, Hello world 2.getBytes(StandardCharsets.UTF_8))); SendResult sendResult producer.send(messages); System.out.printf(%s, sendResult); } }其中Message构造参数依次为topic、tag、消息唯一键key可用于按 key 查询消息、消息体字节数组。官方文档中的原始示例使用Hello world 0.getBytes()仓库示例则显式指定了StandardCharsets.UTF_8在实际工程中建议同样显式指定字符集避免跨平台默认字符集不一致导致的乱码。批量发送的 API 家族DefaultMQProducer针对CollectionMessage提供了多组重载见 DefaultMQProducer.java覆盖不同发送场景方法签名说明SendResult send(CollectionMessage msgs)同步发送使用默认发送超时SendResult send(CollectionMessage msgs, long timeout)同步发送自定义超时SendResult send(CollectionMessage msgs, MessageQueue messageQueue)同步发送到指定队列void send(CollectionMessage msgs, SendCallback sendCallback)异步发送通过回调接收结果void send(CollectionMessage msgs, MessageQueue mq, SendCallback sendCallback, long timeout)异步发送到指定队列并自定义超时这些重载的内部实现都是先将CollectionMessage通过batch(msgs)封装成批量消息再委托给defaultMQProducerImpl.send(...)走与单条发送相同的链路例如public SendResult send(CollectionMessage msgs, long timeout) throws MQClientException, RemotingException, MQBrokerException, InterruptedException { return this.defaultMQProducerImpl.send(batch(msgs), timeout); }客户端内部如何封装批量消息batch(msgs)最终调用的是MessageBatch.generateFromList(messages)MessageBatch.java。MessageBatch继承自Message并实现了IterableMessage它做三件事逐条执行上一节所述的约束校验延迟消息、重试主题、topic 一致性、waitStoreMsgOK一致性取出第一条消息的 topic 与waitStoreMsgOK作为整批消息的外层属性setTopic(first.getTopic())、setWaitStoreMsgOK(first.isWaitStoreMsgOK())通过encode()方法调用MessageDecoder.encodeMessages(messages)将批内多条消息按 RocketMQ 的二进制协议顺序编码进同一个消息体中作为一条复合消息交给底层 remoting 发送。也就是说从网络传输角度看一次批量发送就是一次单条消息的发送只是消息体内部包含了多条子消息的编码数据这正是吞吐量提升的根本原因批内消息数量越多节省的请求往返与协议头开销越明显。三、超过 4MiB 的大批量消息使用 ListSplitter 拆分当待发送的消息总大小不确定、或明确可能超过 4MiB 时直接producer.send(messages)会触发大小校验失败。此时官方推荐的做法是将大列表拆分成多个不超过 1MiB 的小批量逐个发送。1MiB 而不是 4MiB 的拆分粒度是为了给消息在传输、编码过程中产生的额外开销协议头、属性、日志开销等预留足够的余量。官方文档给出了一份ListSplitter实现仓库中的 SplitBatchProducer.java 是其可直接运行的增强版本修复了文档示例中curIndex与getStartIndex的变量名笔误并补充了单条消息超过上限时的防死循环保护class ListSplitter implements IteratorListMessage { private static final int SIZE_LIMIT 1000 * 1000; // 1MiB private final ListMessage messages; private int currIndex; public ListSplitter(ListMessage messages) { this.messages messages; } Override public boolean hasNext() { return currIndex messages.size(); } Override public ListMessage next() { int nextIndex currIndex; int totalSize 0; for (; nextIndex messages.size(); nextIndex) { Message message messages.get(nextIndex); int tmpSize message.getTopic().length() message.getBody().length; MapString, String properties message.getProperties(); for (Map.EntryString, String entry : properties.entrySet()) { tmpSize entry.getKey().length() entry.getValue().length(); } // 为日志/编码开销预留 20 字节 tmpSize tmpSize 20; if (tmpSize SIZE_LIMIT) { // 单条消息本身超过上限属异常情况这里放行以免阻塞拆分流程 if (nextIndex - currIndex 0) { nextIndex; } break; } if (tmpSize totalSize SIZE_LIMIT) { break; } else { totalSize tmpSize; } } ListMessage subList messages.subList(currIndex, nextIndex); currIndex nextIndex; return subList; } Override public void remove() { throw new UnsupportedOperationException(Not allowed to remove); } }拆分算法的关键设计点消息体积估算公式topic 长度 body 长度 所有属性 key/value 长度之和 20 字节。其中 20 字节用于补偿协议/日志开销log overhead。注意这里Message.getProperties()返回的属性映射已包含 tag、key、系统属性等自动附加的键值因此估算结果基本覆盖了消息在存储与传输中的实际体积。贪心累加从currIndex开始向后累加消息体积直到加入下一条会超过 1MiB 上限为止将[currIndex, nextIndex)区间切为一个子列表。单条超限保护如果某条消息单独就超过SIZE_LIMIT仓库版本做了特殊处理——当nextIndex - currIndex 0当前子列表还没有任何元素时强制nextIndex把这条超限消息单独发出去避免hasNext()恒真导致的死循环否则直接跳出。内存效率拆分使用List.subList()视图而非复制元素不会产生额外的大列表拷贝。使用拆分器发送ListSplitter splitter new ListSplitter(messages); while (splitter.hasNext()) { try { ListMessage listItem splitter.next(); producer.send(listItem); } catch (Exception e) { e.printStackTrace(); // 处理失败可记录失败的子列表稍后重试或转入单条发送 } }仓库示例 SplitBatchProducer.java 中构造了100 * 1000十万条消息的大批量用上述拆分器循环分批发送并打印每次的SendResult可以直接作为压力验证脚本使用public static final int MESSAGE_COUNT 100 * 1000; // ... ListMessage messages new ArrayList(MESSAGE_COUNT); for (int i 0; i MESSAGE_COUNT; i) { messages.add(new Message(TOPIC, TAG, OrderID i, (Hello world i).getBytes(StandardCharsets.UTF_8))); } ListSplitter splitter new ListSplitter(messages); while (splitter.hasNext()) { ListMessage listItem splitter.next(); SendResult sendResult producer.send(listItem); System.out.printf(%s, sendResult); }四、运行前提与工程化建议运行前置条件示例默认使用127.0.0.1:9876作为 NameServer 地址本地调试时需先启动 NameServer 与 Broker并取消代码中producer.setNamesrvAddr(...)的注释改为实际地址主题BatchTest需要提前创建可通过mqadmin updateTopic或管理控制台创建生产者的PRODUCER_GROUPBatchProducerGroupName在集群中应保持唯一命名避免与其他业务组冲突。工程化建议失败子列表的重试策略示例中 catch 后仅打印堆栈。生产环境建议把失败的listItem暂存结合退避策略重试多次失败后再降级为逐条发送避免整批数据丢失体积估算与上限对齐拆分粒度建议保持官方推荐的 1MiB。若通过producer.setMaxMessageSize()放大上限请同步确认 Broker 端maxMessageSize配置两端不一致可能导致发送端校验通过而 Broker 拒收异步批量发送若对延迟敏感可改用send(CollectionMessage, SendCallback)系列异步接口在回调中处理成功/失败配合本地缓冲攒批如按时间窗或条数阈值触发能进一步摊薄 RPC 开销避免混用延迟消息任何需要延迟投递的消息都不要放入批量发送改用单条发送并设置delayTimeLevel等延迟属性MessageBatch.java 会直接抛异常拒绝。五、总结批量消息发送是 RocketMQ 生产端最直接的吞吐优化手段小批量≤4MiB直接调用producer.send(CollectionMessage)大批量则借助ListSplitter按 1MiB 粒度拆分后循环发送。理解三条硬约束topic 一致、waitStoreMsgOK一致、不支持延迟消息与 4MiB 上限的来源生产端默认maxMessageSize 4M见 DefaultMQProducer.java以及MessageBatch.generateFromList的封装校验逻辑就能在享受吞吐提升的同时规避踩坑。可运行示例位于 example/src/main/java/org/apache/rocketmq/example/batch读者可直接对照源码加深理解。赞分享消息队列后端微服务流处理【免费下载链接】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点击查看免费下载相关推荐Apache RocketMQ 批量消息发送实战4MiB 限制、ListSplitter 大消息拆分与底层原理Apache RocketMQ 批量消息发送实战4MiB 限制、ListSplitter 大消息拆分与底层原理 批量发送是 Apache RocketMQ 生消息队列流处理后端CANN/ge LLM集群连接API link\_clusters 产品支持情况 Atlas A3 训练系列产品/Atlas A3 推理系列产品支持 Atlas A2 推理系列产品支持 At消息队列流处理后端Apache RocketMQ 批量消息发送实战指南从 4MiB 单批限制到大数据量自动切分Apache RocketMQ 批量消息发送实战指南从 4MiB 单批限制到大数据量自动切分 批量发送是 Apache RocketMQ 生产端提升吞吐的关键消息队列后端微服务流处理上一篇如何用WeChatMsg永久保存微信聊天记录从数据碎片到数字记忆的完整指南下一篇如何永久保存微信聊天记录3步实现数据自主的终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/9/20 23:47:22

OBV能量潮改选股公式:捕捉主力资金启动前夜

简介:面向股票技术分析与通达信指标使用者,这份教程性质资源给出了OBV能量潮改造的选股公式源码,并围绕其编写思路与实战含义展开讲解。文档先介绍OBV指标衡量买卖压力与资金流向的基本原理,再逐步拆解公式中的关键节点&#xff1…

2026/9/20 23:47:22

用Python和Playwright实现头条自动发文:从登录到发布的自动化实战

简介:面向熟悉 Python,希望通过爬虫与自动化脚本提升内容发布效率的开发者,提供一套今日头条自动发文项目源码。项目综合运用爬虫技术与浏览器自动化,从新闻 API、知乎热榜等渠道抓取内容,并采用 PyQt5 构建可视化操作…

2026/9/21 0:47:25

RAG技术优化:检索增强生成系统的关键策略与实践

1. RAG技术体系概述检索增强生成(Retrieval-Augmented Generation)作为当前NLP领域的前沿技术,通过将信息检索与文本生成相结合,有效解决了传统大语言模型的知识固化问题。我在实际项目中发现,标准的RAG流程通常包含四…

2026/9/21 0:47:25

Claude Code 桌面版接入 DeepSeek 与离线 Skills 安装全攻略

1. 为什么我要折腾这套组合:Claude Code 桌面版 DeepSeek 离线 Skills先说清楚这套东西到底是什么。Claude Code 是 Anthropic 推出的一个命令行 AI 编程助手,它跟普通聊天式 AI 最大的区别在于:它能直接读写你本地的项目文件、执行终端命令…

2026/9/21 0:47:25

QGIS等时圈分析实战:ORS插件Key申请与参数设置避坑指南

1. 等时圈分析与ORS插件到底在做什么等时圈分析这件事,说白了就是回答一个很朴素的问题:从某个点出发,在给定时间内,我到底能走到哪些地方。做城市规划的要拿它评估公共服务覆盖范围,做商业选址的要拿它算门店辐射半径…

2026/9/21 0:47:25

普通人用AI变现,第一个工具到底该怎么选?

我见过太多人,一听说AI能变现,第一反应就是到处问:现在哪个AI工具最强?哪个能不限次数白嫖?哪个生成的内容最像真人?然后就开始了一场漫长的工具测评之旅。各种官网、教程、对比帖收藏了上百篇,…

2026/9/21 0:47:25

JDK 17.0.8免安装版Windows配置指南:从下载到环境变量

简介:JDK 17.0.8 Windows免安装版为Java开发者提供开箱即用的开发环境,无需经过复杂安装流程,解压配置环境变量即可使用。作为长期支持(LTS)版本,它包含javac编译器、Java运行环境、javadoc文档生成器、jdb…

2026/9/21 0:42:24

Xilinx 7系列FPGA入门:从选型架构到时序约束实战要点

简介:面向FPGA初学者与嵌入式开发者的Xilinx 7系列FPGA入门介绍文档,以简明方式梳理系列整体定位与核心技术要点。内容涵盖Spartan-7、Artix-7、Kintex-7、Virtex-7四个子系列的适用场景、性能参数与功耗优势,详细对比单位功耗性价比、成本削…

2026/9/20 0:04:49

GAMP 5 基于风险的计算机化系统验证:软件分类与审计追踪实践

简介:《A Risk-Based Approach to Compliant GxP Computerized Systems》即业内熟知的GAMP 5指南,面向制药企业质量与IT合规人员、验证工程师及计算机化系统管理者,用于解决GxP法规环境下系统合规性难以科学落地的问题。文档以风险管理为主线…

2026/9/20 0:04:49

安全托管MSSP实战:从静态防御到人机协同的攻防运营与应急响应

简介:这份PPT围绕互联网业务安全托管服务展开,面向企业安全负责人、IT运维人员及关注MSSP/MSS选型的读者,重点回应传统安全过度依赖人工、碎片化静态防御难以对抗产业化攻击等痛点。资源共1个pptx文件,包体约30.63MB,以…

2026/9/21 0:02:23

OpenResearch:构建可复现的开放式研究工作流

第一次看到“OpenResearch”这个名字,我脑子里冒出的不是某个具体软件,而更像一种研究方式的宣言:开放、可复现、可验证。这三件事放在一起,其实比大多数人想象中难得多。过去几年我一直在折腾自己的研究工作流,从纯纸…

2026/9/20 4:54:47

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/20 5:01:23

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/20 5:09:33

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码