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

使用MassTransit接收Azure Service Bus无体消息时遇序列化异常

问题分析与解决方案

错误原因

你遇到的序列化错误,核心是MassTransit默认会尝试将消息体反序列化为指定类型,但无消息体时JSON反序列化器无法解析空输入。另外,直接消费ServiceBusReceivedMessage(Azure Service Bus SDK的底层类型)并非MassTransit的推荐用法,它更适合处理封装后的消息契约。


解决方案

方案一:使用非泛型消费者接收原始消息(推荐)

通过非泛型IConsumer直接获取上下文里的原始Azure Service Bus消息,跳过消息体反序列化步骤:

修改消费者代码

public class MessageConsumer : IConsumer
{
    private readonly ILogger<MessageConsumer> _logger;

    public MessageConsumer(ILogger<MessageConsumer> logger)
    {
        _logger = logger;
    }

    public Task Consume(ConsumeContext context)
    {
        // 获取原始Azure Service Bus消息
        var serviceBusMessage = context.GetPayload<ServiceBusReceivedMessage>();
        
        _logger.LogInformation("Message id: {MessageId}", serviceBusMessage.MessageId);
        
        // 遍历处理消息头
        foreach (var header in serviceBusMessage.ApplicationProperties)
        {
            _logger.LogInformation("Header {Key}: {Value}", header.Key, header.Value);
        }
        
        return Task.CompletedTask;
    }
}

修改MassTransit配置

在订阅端点中配置原始JSON序列化器,兼容空消息体:

serviceCollection.AddMassTransit(x =>
{
    x.AddConsumers(typeof(MessageConsumer).Assembly);
    x.UsingAzureServiceBus((context, cfg) =>
    {
        cfg.Host(configuration.GetValue<string>(connectionString));
        
        cfg.SubscriptionEndpoint("subscriptionName", "topicPath", endpointCfg =>
        {
            // 启用原始JSON序列化器,支持空消息体
            endpointCfg.UseRawJsonSerializer();
            
            endpointCfg.Consumer<MessageConsumer>(context);
        });

        cfg.ConfigureEndpoints(context);
    });
});

方案二:自定义序列化器处理空消息体

如果需要保留泛型消费者,可以自定义序列化器,当检测到空消息体时返回默认实例:

自定义序列化器

public class EmptyBodySerializer : IMessageSerializer
{
    private readonly IMessageSerializer _innerSerializer;

    public EmptyBodySerializer(IMessageSerializer innerSerializer)
    {
        _innerSerializer = innerSerializer;
    }

    public ContentType ContentType => _innerSerializer.ContentType;

    public MessageBody Serialize<T>(T message) where T : class
    {
        return _innerSerializer.Serialize(message);
    }

    public T Deserialize<T>(MessageBody body) where T : class
    {
        return body.Length == 0 ? default : _innerSerializer.Deserialize<T>(body);
    }

    public object Deserialize(MessageBody body, Type type)
    {
        return body.Length == 0 ? Activator.CreateInstance(type) : _innerSerializer.Deserialize(body, type);
    }
}

配置中替换序列化器

endpointCfg.UseSerializer(context => new EmptyBodySerializer(context.GetRequiredService<IMessageSerializer>()));

内容的提问来源于stack exchange,提问作者ChinnMan Spam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 11:46:32