如何使用MassTransit读取RabbitMQ现有队列中的Base64编码消息
使用MassTransit读取RabbitMQ中Base64编码的消息
你当前的错误是因为配置了UseRawJsonSerializer,但队列中的消息是Base64编码的字符串而非JSON格式,导致MassTransit尝试将Base64字符串解析为JSON时失败。以下是两种可行的解决方案:
方案一:直接消费原始字节并手动解码Base64
修改ReceiveEndpoint配置,使用UseRawSerializer接收原始字节流:
cfg.ReceiveEndpoint("test-queue", e => { e.ConfigureConsumeTopology = false; e.ClearMessageDeserializers(); // 使用RawSerializer接收原始字节 e.UseRawSerializer(); e.Consumer<TestQueueConsumer>(context); e.Durable = true; });
实现对应消费者,处理byte[]类型消息并解码Base64:
using System.Text; using MassTransit; public class TestQueueConsumer : IConsumer<byte[]> { public async Task Consume(ConsumeContext<byte[]> context) { // 将字节数组转换为Base64字符串(假设消息以UTF-8编码) var base64Content = Encoding.UTF8.GetString(context.Message); // 解码Base64得到原始消息 var decodedBytes = Convert.FromBase64String(base64Content); var originalMessage = Encoding.UTF8.GetString(decodedBytes); // 这里添加你的业务逻辑处理 Console.WriteLine($"解码后的原始消息:{originalMessage}"); await Task.CompletedTask; } }
方案二:自定义反序列化器映射到自定义消息类型
如果希望将解码后的消息封装到自定义类型中,可以创建自定义反序列化器:
第一步:定义消息实体类
public class QueueMessage { public string Content { get; set; } }
第二步:实现自定义反序列化器
using System.Text; using MassTransit; public class Base64QueueMessageDeserializer : IMessageDeserializer { // 匹配队列消息的Content-Type,若发送时未指定,默认text/plain public ContentType ContentType => new ContentType("text/plain"); public string[] ContentTypeAliases => Array.Empty<string>(); public async Task<ConsumeContext> Deserialize(ReceiveContext receiveContext) { // 读取消息体字节 var bodyBytes = await receiveContext.GetBodyStream().ReadToEndAsync(); // 转换为Base64字符串并解码 var base64String = Encoding.UTF8.GetString(bodyBytes); var decodedBytes = Convert.FromBase64String(base64String); var originalContent = Encoding.UTF8.GetString(decodedBytes); // 封装到自定义消息类型 var message = new QueueMessage { Content = originalContent }; // 创建消费上下文并返回 return new ConsumeContext<QueueMessage>(receiveContext, message); } }
第三步:修改ReceiveEndpoint配置
cfg.ReceiveEndpoint("test-queue", e => { e.ConfigureConsumeTopology = false; e.ClearMessageDeserializers(); // 注册自定义反序列化器 e.AddMessageDeserializer(new Base64QueueMessageDeserializer()); e.Consumer<TestQueueConsumer>(context); e.Durable = true; });
第四步:实现对应自定义类型的消费者
using MassTransit; public class TestQueueConsumer : IConsumer<QueueMessage> { public async Task Consume(ConsumeContext<QueueMessage> context) { // 处理解码后的消息 Console.WriteLine($"解码后的消息内容:{context.Message.Content}"); await Task.CompletedTask; } }
注意事项
- 确保发送到队列的消息是标准Base64编码字符串,若存在编码格式差异(如UTF-16),需调整解码时的
Encoding类型。 - 如果发送消息时指定了特定的
Content-Type,自定义反序列化器的ContentType需与之匹配,否则MassTransit不会使用该反序列化器。
内容的提问来源于stack exchange,提问作者MrMohande3
相关产品推荐
相关产品推荐

