You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 07:27:34