如何在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)
- 实现自定义消息转换器,将纯文本消息转换为
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); } }
- 修改接收端点配置,添加自定义转换器:
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; }); });
- 从RabbitMQ控制台发送消息时,务必设置ContentType为
text/plain,消息内容填写纯文本(如abcds)。
方案2:直接消费string类型消息
如果不需要严格绑定SampleMessage,可以修改消费者直接接收纯文本:
- 修改消费者实现:
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; } }
- 修改服务注册配置,添加字符串消息的转换器:
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); }); });
- 从RabbitMQ控制台发送消息时,ContentType设置为
application/json,消息内容用双引号包裹(如"abcds"),或者设置为text/plain并配合转换器使用。
内容的提问来源于stack exchange,提问作者Abhishek Patnaik
相关产品推荐
相关产品推荐

