如何在MassTransit中使用自定义JsonConverter消费原始JSON?
在MassTransit中触发自定义JsonConverter消费带外层节点的JSON
问题根源
你当前配置中调用了UseRawJsonDeserializer(isDefault: true),这会让MassTransit启用原始JSON反序列化器。该反序列化器的逻辑是直接将接收到的完整JSON字符串反序列化为目标消息类型(MyClass),但它完全独立于ConfigureJsonSerializerOptions配置的System.Text.Json选项链——也就是说你添加的MyClassMessageConverter根本不会被这个反序列化器调用,再加上你的JSON外层多了data节点,直接反序列化必然失败。
解决方案
方案一:移除RawJson配置,使用默认System.Text.Json序列化器
直接去掉RawJson相关的三行配置,保留自定义转换器的配置,让MassTransit使用默认的System.Text.Json序列化器,此时转换器会被正常触发。
配置代码修改:
services.AddMassTransit(opt => { opt.AddConsumer(consumerType); opt.UsingAzureServiceBus((context, cfg) => { // 移除RawJson相关的ClearSerialization、UseRawJsonSerializer、UseRawJsonDeserializer调用 cfg.ConfigureJsonSerializerOptions(options => { options.Converters.Insert(0, new MyClassMessageConverter()); return options; }); cfg.Host(config.GetConnectionString("ServiceBus")); cfg.Message<MyClass>(x => x.SetEntityName(config["Messaging:MembershipEndpoint"]!)); cfg.SubscriptionEndpoint<MyClass>(config["Messaging:MembershipEndpoint"]!, e => { e.ConfigureConsumer(context, consumerType); }); cfg.ConfigureEndpoints(context); }); });
自定义转换器示例(需正确处理外层data节点):
public class MyClassMessageConverter : JsonConverter<MyClass> { public override MyClass Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options) { using var doc = JsonDocument.ParseValue(ref reader); var dataElement = doc.RootElement.GetProperty("data"); return JsonSerializer.Deserialize<MyClass>(dataElement.GetRawText(), options); } public override void Write(Utf8JsonWriter writer, MyClass value, JsonSerializerOptions options) { // 如果需要发送消息时也包装为带外层data的结构,实现此方法 writer.WriteStartObject(); writer.WritePropertyName("data"); JsonSerializer.Serialize(writer, value, options); writer.WriteEndObject(); } }
方案二:自定义RawJson反序列化器(若必须保留RawJson模式)
如果你因为某些原因必须使用RawJson序列化器,需要自定义一个IMessageDeserializer,先提取data节点的内容,再交给原始RawJson反序列化器处理。
自定义反序列化器代码:
public class CustomRawJsonDeserializer : IMessageDeserializer { private readonly RawJsonDeserializer _innerDeserializer; public CustomRawJsonDeserializer(RawJsonDeserializer innerDeserializer) { _innerDeserializer = innerDeserializer; } public ContentType ContentType => _innerDeserializer.ContentType; public MessageBody Deserialize(ReceiveContext context) { var rawBody = context.GetBody(); using var doc = JsonDocument.Parse(rawBody); var dataElement = doc.RootElement.GetProperty("data"); var dataJsonBytes = Encoding.UTF8.GetBytes(dataElement.GetRawText()); // 构造模拟ReceiveContext传递处理后的JSON var mockContext = new MockReceiveContext(dataJsonBytes); return _innerDeserializer.Deserialize(mockContext); } private class MockReceiveContext : ReceiveContext { private readonly byte[] _body; public MockReceiveContext(byte[] body) : base(new Uri("mock://"), default) { _body = body; } public override byte[] GetBody() => _body; } }
配置代码修改:
services.AddMassTransit(opt => { opt.AddConsumer(consumerType); opt.UsingAzureServiceBus((context, cfg) => { cfg.ClearSerialization(); // 注册自定义反序列化器 var rawSerializer = new RawJsonSerializer(); var customDeserializer = new CustomRawJsonDeserializer(new RawJsonDeserializer()); cfg.AddSerializer(rawSerializer); cfg.AddDeserializer(customDeserializer, isDefault: true); cfg.Host(config.GetConnectionString("ServiceBus")); cfg.Message<MyClass>(x => x.SetEntityName(config["Messaging:MembershipEndpoint"]!)); cfg.SubscriptionEndpoint<MyClass>(config["Messaging:MembershipEndpoint"]!, e => { e.ConfigureConsumer(context, consumerType); }); cfg.ConfigureEndpoints(context); }); });
内容的提问来源于stack exchange,提问作者AKA
相关产品推荐
相关产品推荐

