ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

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),仅供参考
返回列表