C# Send<T>()发送RabbitMQ前获取消息大小并拆分大消息咨询
解决RabbitMQ消息大小超限的高效拆分方案
核心逻辑
RabbitMQ的消息大小限制针对的是序列化后的字节总数,而非内存中CLR对象的大小。因此必须基于你实际使用的消息序列化器(如JSON、Protobuf)计算最终字节数,再通过"固定开销+批量Item"的方式拆分消息,既能规避超限,又能最大化单条消息的Item数量、减少发送次数。
具体实现步骤
1. 实现序列化大小计算工具
根据你项目使用的序列化方式,编写对应的字节数计算方法:
场景1:使用Newtonsoft.Json序列化
public static long GetJsonSerializedSize<T>(T obj) { using var memoryStream = new MemoryStream(); var serializer = new JsonSerializer(); using var streamWriter = new StreamWriter(memoryStream); using var jsonWriter = new JsonTextWriter(streamWriter); serializer.Serialize(jsonWriter, obj); return memoryStream.Length; }
场景2:使用Protobuf序列化(推荐,字节效率更高)
public static long GetProtobufSerializedSize<T>(T obj) where T : IMessage<T> { return obj.CalculateSize(); // Protobuf内置方法直接计算序列化大小 }
2. 计算消息固定部分的开销
先创建不含Item的空消息,算出这部分的固定字节数(每条消息都会包含这部分内容):
var emptyMessage = new Message { UserID = originalMessage.UserID, UserToken = originalMessage.UserToken, datetime = originalMessage.datetime, Itens = new List<T>() }; // 替换为你实际使用的序列化计算方法 long fixedSize = GetJsonSerializedSize(emptyMessage); // Protobuf场景则用:GetProtobufSerializedSize(emptyMessage)
3. 批量拆分并发送消息
基于固定开销和单个Item的大小,批量组装消息直到接近阈值后发送:
const long MAX_RABBITMQ_LIMIT = 134217728; // RabbitMQ限制的128MB const long SAFETY_MARGIN = 5 * 1024 * 1024; // 预留5MB余量,避免意外超限 long maxAllowedItemTotalSize = MAX_RABBITMQ_LIMIT - SAFETY_MARGIN - fixedSize; var currentBatch = new Message { UserID = originalMessage.UserID, UserToken = originalMessage.UserToken, datetime = originalMessage.datetime, Itens = new List<T>() }; long currentBatchItemSizeSum = 0; foreach (var item in originalMessage.Itens) { // 计算当前Item的序列化大小 long itemSize = GetJsonSerializedSize(item); // Protobuf场景则用:GetProtobufSerializedSize(item) // 加入当前Item后会超出限制,先发送现有批次 if (currentBatchItemSizeSum + itemSize > maxAllowedItemTotalSize) { await sendEndpoint.Send(currentBatch); // 重置新批次 currentBatch = new Message { UserID = originalMessage.UserID, UserToken = originalMessage.UserToken, datetime = originalMessage.datetime, Itens = new List<T>() }; currentBatchItemSizeSum = 0; } currentBatch.Itens.Add(item); currentBatchItemSizeSum += itemSize; } // 发送剩余的最后一批数据 if (currentBatch.Itens.Count > 0) { await sendEndpoint.Send(currentBatch); }
4. 优化建议
- 缓存Item大小:如果所有Item的结构、数据量差异极小,可以预先计算样本平均大小,后续用平均值估算,减少重复序列化的开销。
- 优先用Protobuf:Protobuf序列化后的字节数远小于JSON,能大幅提升单条消息可容纳的Item数量,减少拆分次数。
- 动态调整固定开销:若UserToken等固定字段长度波动较大,可在每次创建新批次时重新计算固定大小,避免余量不足。
注意事项
- 必须保证计算逻辑和实际发送时的序列化逻辑完全一致,否则会出现计算值与实际发送值不符,导致仍然超限。
- 若Item包含大字段(如长文本、二进制数据),需单独处理这类特殊Item,避免单个Item就超过阈值。
内容的提问来源于stack exchange,提问作者Fernando
相关产品推荐
相关产品推荐

