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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 05:37:05