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

如何在C#中使用MassTransit消费RabbitMQ纯文本消息

问题:MassTransit消费RabbitMQ纯文本消息抛出序列化异常

当前配置代码

服务注册配置:

services.AddMassTransit(x =>
{
    x.AddConsumer<MessageConsumer>();
    x.UsingRabbitMq((context, cfg) =>
    {
        cfg.Host("localhost", hostConfigurator =>
        {
            hostConfigurator.Username(rabbitMqConfigurations.Username);
            hostConfigurator.Password(rabbitMqConfigurations.Password);
        });
        cfg.ReceiveEndpoint("testQ", configurator =>
        {
            configurator.ConfigureConsumer<MessageConsumer>(context);
        });
        cfg.ConfigureEndpoints(context);
    });
});

消费者实现

public class MessageConsumer : IConsumer<SampleMessage>
{
    public Task Consume(ConsumeContext<SampleMessage> context) 
    {
        Console.WriteLine("Hello World");
        return Task.CompletedTask;
    }
}

消息定义

public sealed record SampleMessage(string Data);

抛出异常

Exception thrown: 'System.Text.Json.JsonException' in System.Text.Json.dll
Exception thrown: 'System.Runtime.Serialization.SerializationException' in MassTransit.dll
Exception thrown: 'System.Runtime.Serialization.SerializationException' in System.Private.CoreLib.dll

使用版本:MassTransit 8.2.4、MassTransit.RabbitMQ 8.2.4、MassTransit.Newtonsoft 8.2.4


解决方案

原因

MassTransit默认期望消息携带类型元数据(如__TypeId头部),或消息内容是符合目标类型(SampleMessage)结构的JSON。直接发送纯文本无法被反序列化为SampleMessage对象,因此抛出序列化异常。

方案1:自定义消息转换器(将纯文本转为SampleMessage)

  1. 实现自定义消息转换器,将纯文本消息转换为SampleMessage:
public class PlainTextToSampleMessageConverter : IMessageConverter
{
    public bool CanConvert(MessageContext context, Type messageType)
    {
        // 仅处理text/plain类型的消息转为SampleMessage
        return messageType == typeof(SampleMessage) 
               && context.ContentType?.MediaType.Equals("text/plain", StringComparison.OrdinalIgnoreCase) == true;
    }

    public Task<MessageBody> ConvertToBody<T>(T message, MessageContext context) where T : class
    {
        throw new NotImplementedException("本转换器仅支持从纯文本到SampleMessage的反序列化");
    }

    public Task<T?> ConvertToMessage<T>(MessageBody body, MessageContext context) where T : class
    {
        var plainText = body.GetText();
        var sampleMessage = new SampleMessage(plainText);
        return Task.FromResult<T?>(sampleMessage as T);
    }
}
  1. 修改接收端点配置,添加自定义转换器:
cfg.ReceiveEndpoint("testQ", configurator =>
{
    configurator.ConfigureConsumer<MessageConsumer>(context);
    // 注册自定义转换器
    configurator.AddMessageConverter<SampleMessage>(new PlainTextToSampleMessageConverter());
    // 禁用默认消费拓扑,避免与默认Exchange绑定冲突
    configurator.ConfigureConsumeTopology = false;
    // 手动绑定队列到Exchange(根据你的RabbitMQ配置调整)
    configurator.Bind("testQ", x =>
    {
        x.RoutingKey = "testQ";
        x.ExchangeType = ExchangeType.Direct;
    });
});
  1. 从RabbitMQ控制台发送消息时,务必设置ContentType为text/plain,消息内容填写纯文本(如abcds)。

方案2:直接消费string类型消息

如果不需要严格绑定SampleMessage,可以修改消费者直接接收纯文本:

  1. 修改消费者实现:
public class MessageConsumer : IConsumer<string>
{
    public Task Consume(ConsumeContext<string> context) 
    {
        Console.WriteLine($"Received plain text: {context.Message}");
        // 可选:转为SampleMessage处理
        var sampleMsg = new SampleMessage(context.Message);
        // ... 后续业务逻辑
        return Task.CompletedTask;
    }
}
  1. 修改服务注册配置,添加字符串消息的转换器:
services.AddMassTransit(x =>
{
    x.AddConsumer<MessageConsumer>();
    x.UsingRabbitMq((context, cfg) =>
    {
        cfg.Host("localhost", hostConfigurator =>
        {
            hostConfigurator.Username(rabbitMqConfigurations.Username);
            hostConfigurator.Password(rabbitMqConfigurations.Password);
        });
        cfg.ReceiveEndpoint("testQ", configurator =>
        {
            configurator.ConfigureConsumer<MessageConsumer>(context);
            // 使用原始JSON反序列化器处理纯文本字符串
            configurator.UseRawJsonDeserializer();
            configurator.AddMessageConverter<string>(new RawJsonMessageConverter());
        });
        cfg.ConfigureEndpoints(context);
    });
});
  1. 从RabbitMQ控制台发送消息时,ContentType设置为application/json,消息内容用双引号包裹(如"abcds"),或者设置为text/plain并配合转换器使用。

内容的提问来源于stack exchange,提问作者Abhishek Patnaik

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 08:28:16