- 消息队列
- 后端
- 微服务
- 流处理
【免费下载链接】rocketmq
Apache RocketMQ is a cloud native messaging and streaming platform, making it simple to build event-driven applications.
批量发送是 Apache RocketMQ 生产端提升发送效率、抬高系统吞吐量的核心技术手段:它把多条消息封装成一次网络请求提交给 Broker,显著减少 RPC 往返次数。本文以官方文档《批量消息发送》为主线,结合仓库中的生产端示例 SimpleBatchProducer.java、SplitBatchProducer.java 与客户端源码,讲清批量发送的适用条件、4MiB 上限的由来、大消息拆分算法,以及批量消息在客户端内部的编码与校验链路,读完即可在自己的生产者代码中落地实现。
一、批量消息的适用条件与核心约束
在动手写代码之前,必须先理解 RocketMQ 对"同一批消息"的硬性要求。这些约束不是文档建议,而是客户端源码层面的强制校验,违反会直接抛出UnsupportedOperationException:
- 同一批消息的 topic 必须一致:批量消息在编码阶段会被拼装为一条复合消息,其外层只能携带一个 topic;
- 同一批消息的
waitStoreMsgOK属性必须一致:该属性决定发送时是否等待 Broker 落盘确认(同步刷盘/异步刷盘语义),一批内混用会导致语义不明确; - 批量消息不支持延迟消息:无论是
delayTimeLevel(延迟级别)、delayTimeMs、delayTimeSec还是deliverTimeMs,只要设置了任何一个延迟相关属性,批量发送都会被拒绝; - 批量消息不支持重试主题(Retry Topic):topic 以重试组前缀开头时同样被拒绝。
上述校验可以在 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"; List<Message> 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针对Collection<Message>提供了多组重载(见 DefaultMQProducer.java),覆盖不同发送场景:
| 方法签名 | 说明 |
|---|---|
SendResult send(Collection<Message> msgs) | 同步发送,使用默认发送超时 |
SendResult send(Collection<Message> msgs, long timeout) | 同步发送,自定义超时 |
SendResult send(Collection<Message> msgs, MessageQueue messageQueue) | 同步发送到指定队列 |
void send(Collection<Message> msgs, SendCallback sendCallback) | 异步发送,通过回调接收结果 |
void send(Collection<Message> msgs, MessageQueue mq, SendCallback sendCallback, long timeout) | 异步发送到指定队列并自定义超时 |
这些重载的内部实现都是先将Collection<Message>通过batch(msgs)封装成批量消息,再委托给defaultMQProducerImpl.send(...)走与单条发送相同的链路,例如:
public SendResult send(Collection<Message> msgs, long timeout) throws MQClientException, RemotingException, MQBrokerException, InterruptedException { return this.defaultMQProducerImpl.send(batch(msgs), timeout); }客户端内部如何封装批量消息
batch(msgs)最终调用的是MessageBatch.generateFromList(messages)(MessageBatch.java)。MessageBatch继承自Message并实现了Iterable<Message>,它做三件事:
- 逐条执行上一节所述的约束校验(延迟消息、重试主题、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 Iterator<List<Message>> { private static final int SIZE_LIMIT = 1000 * 1000; // 1MiB private final List<Message> messages; private int currIndex; public ListSplitter(List<Message> messages) { this.messages = messages; } @Override public boolean hasNext() { return currIndex < messages.size(); } @Override public List<Message> 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; Map<String, String> properties = message.getProperties(); for (Map.Entry<String, 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; } } List<Message> 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 { List<Message> listItem = splitter.next(); producer.send(listItem); } catch (Exception e) { e.printStackTrace(); // 处理失败:可记录失败的子列表,稍后重试或转入单条发送 } }仓库示例 SplitBatchProducer.java 中构造了100 * 1000(十万条)消息的大批量,用上述拆分器循环分批发送并打印每次的SendResult,可以直接作为压力验证脚本使用:
public static final int MESSAGE_COUNT = 100 * 1000; // ... List<Message> 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()) { List<Message> 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_GROUP(BatchProducerGroupName)在集群中应保持唯一命名,避免与其他业务组冲突。
工程化建议
- 失败子列表的重试策略:示例中 catch 后仅打印堆栈。生产环境建议把失败的
listItem暂存,结合退避策略重试,多次失败后再降级为逐条发送,避免整批数据丢失; - 体积估算与上限对齐:拆分粒度建议保持官方推荐的 1MiB。若通过
producer.setMaxMessageSize()放大上限,请同步确认 Broker 端maxMessageSize配置,两端不一致可能导致发送端校验通过而 Broker 拒收; - 异步批量发送:若对延迟敏感,可改用
send(Collection<Message>, SendCallback)系列异步接口,在回调中处理成功/失败,配合本地缓冲攒批(如按时间窗或条数阈值触发)能进一步摊薄 RPC 开销; - 避免混用延迟消息:任何需要延迟投递的消息都不要放入批量发送,改用单条发送并设置
delayTimeLevel等延迟属性(MessageBatch.java 会直接抛异常拒绝)。
五、总结
批量消息发送是 RocketMQ 生产端最直接的吞吐优化手段:小批量(≤4MiB)直接调用producer.send(Collection<Message>),大批量则借助ListSplitter按 1MiB 粒度拆分后循环发送。理解三条硬约束(topic 一致、waitStoreMsgOK一致、不支持延迟消息)与 4MiB 上限的来源(生产端默认maxMessageSize = 4M,见 DefaultMQProducer.java),以及MessageBatch.generateFromList的封装校验逻辑,就能在享受吞吐提升的同时规避踩坑。可运行示例位于 example/src/main/java/org/apache/rocketmq/example/batch,读者可直接对照源码加深理解。
- 消息队列
- 后端
- 微服务
- 流处理
【免费下载链接】rocketmq
Apache RocketMQ is a cloud native messaging and streaming platform, making it simple to build event-driven applications.
相关推荐
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 生产端提升吞吐的关键
消息队列后端微服务流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考