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

如何用MassTransit从Azure Service Bus消费FHIR类型消息并解决异常

问题解决:MassTransit消费Azure Service Bus消息时的NotSupportedException异常

问题背景

本地部署了NHS 111适配器,该适配器通过Azure Service Bus推送数据。当使用NHS 111测试套件推送数据时,测试套件会将XML payload发送给适配器,适配器将其转换为FHIR类型消息。尝试用MassTransit消费该消息时触发如下异常:

Exception on Receiver sb://[AZURESERVICEBUS]/[QUEUENAME] during ProcessMessageCallback ActiveDispatchCount(0) ErrorRequiresRecycle(True)
System.NotSupportedException: Value cannot be retrieved using the Body property.Use GetRawAmqpMessage to access the underlying Amqp Message object.
   at BinaryData Azure.Messaging.ServiceBus.Amqp.AmqpMessageExtensions.GetBody(AmqpAnnotatedMessage message)
   at BinaryData Azure.Messaging.ServiceBus.ServiceBusReceivedMessage.get_Body()
   at new MassTransit.AzureServiceBusTransport.ServiceBusReceiveContext(ServiceBusReceivedMessage message, ReceiveEndpointContext receiveEndpointContext) in /_/src/Transports/MassTransit.AzureServiceBus.Core/AzureServiceBusTransport/Contexts/ServiceBusReceiveContext.cs:line 21
   at async Task MassTransit.AzureServiceBusTransport.ServiceBusMessageReceiver.Handle(ServiceBusReceivedMessage message, CancellationToken cancellationToken, Action<ReceiveContext> contextCallback) in /_/src/Transports/MassTransit.AzureServiceBus.Core/AzureServiceBusTransport/ServiceBusMessageReceiver.cs:line 49
   at async Task Azure.Messaging.ServiceBus.ServiceBusProcessor.OnProcessMessageAsync(ProcessMessageEventArgs args)
   at async Task Azure.Messaging.ServiceBus.ReceiverManager.OnMessageHandler(EventArgs args)
   at async Task Azure.Messaging.ServiceBus.ReceiverManager.ProcessOneMessage(ServiceBusReceivedMessage triggerMessage, CancellationToken cancellationToken)

消费者实现代码:

public sealed class Nhs111ClinicalDocumentConsumer : IConsumer<FhirNhs111InboundMessageRegistration>

最初FhirNhs111InboundMessageRegistration是空类,之后创建了匹配消息结构的实体类(代码如下),但问题仍未解决。手动发送消息到队列时可以正常消费,但通过NHS 111测试套件推送时就会报错。

实体类代码:

using Helix.Nhs.Contracts.V1.Commands;
using Microservice.Contracts;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Tasks;

namespace Helix.Nhs.Contracts.V2
{
    // Root myDeserializedClass = JsonConvert.DeserializeObject<Root>(myJsonResponse);
    public sealed class Action
    {
        public List<Coding> coding { get; set; }
        public string text { get; set; }
    }

    public sealed class Agent
    {
        public string reference { get; set; }
    }

    public sealed class Author
    {
        public string reference { get; set; }
    }

    public sealed class Coding
    {
        public string system { get; set; }
        public string code { get; set; }
        public string display { get; set; }
    }

    public sealed class ConsentingParty
    {
        public string reference { get; set; }
    }

    public sealed class Context
    {
        public string reference { get; set; }
    }

    public sealed class Datum
    {
        public string meaning { get; set; }
        public Reference reference { get; set; }
    }

    public sealed class Destination
    {
        public string endpoint { get; set; }
    }

    public sealed class Encounter
    {
        public string reference { get; set; }
    }

    public sealed class Entry
    {
        public string fullUrl { get; set; }
        public Resource resource { get; set; }
        public string reference { get; set; }
        public Item item { get; set; }
    }

    public sealed class Event
    {
        public string system { get; set; }
        public string code { get; set; }
        public string display { get; set; }
    }

    public sealed class GeneralPractitioner
    {
        public string reference { get; set; }
    }

    public sealed class Identifier
    {
        public Type type { get; set; }
        public string value { get; set; }
    }

    public sealed class Individual
    {
        public string reference { get; set; }
    }

    public sealed class Item
    {
        public string reference { get; set; }
    }

    public sealed class Location
    {
        public Location location { get; set; }
        public string status { get; set; }
        public string reference { get; set; }
    }

    public sealed class ManagingOrganization
    {
        public string reference { get; set; }
    }

    public sealed class Meta
    {
        public DateTime lastUpdated { get; set; }
    }

    public sealed class OccurrencePeriod
    {
        public DateTime start { get; set; }
    }

    public sealed class OnBehalfOf
    {
        public string reference { get; set; }
    }

    public sealed class OrderedBy
    {
        public List<Coding> coding { get; set; }
    }

    public sealed class Participant
    {
        public List<Type> type { get; set; }
        public Individual individual { get; set; }
    }

    public sealed class Patient
    {
        public string reference { get; set; }
    }

    public sealed class Period
    {
        public DateTime start { get; set; }
        public DateTime end { get; set; }
    }

    public sealed class Practitioner
    {
        public string reference { get; set; }
    }

    public sealed class ProvidedBy
    {
        public string reference { get; set; }
    }

    public sealed class Reason
    {
        public List<Coding> coding { get; set; }
    }

    public sealed class ReasonReference
    {
        public string reference { get; set; }
    }

    public sealed class Recipient
    {
        public string reference { get; set; }
    }

    public sealed class Reference
    {
        public string reference { get; set; }
    }

    public sealed class Relationship
    {
        public List<Coding> coding { get; set; }
    }

    public sealed class Requester
    {
        public Agent agent { get; set; }
        public OnBehalfOf onBehalfOf { get; set; }
    }

    public sealed class Resource
    {
        public string resourceType { get; set; }
        public string id { get; set; }
        public Event @event { get; set; }
        public List<Destination> destination { get; set; }
        public DateTime timestamp { get; set; }
        public Source source { get; set; }
        public Reason reason { get; set; }
        public object identifier { get; set; }
        public string status { get; set; }
        public object type { get; set; }
        public Subject subject { get; set; }
        public List<Participant> participant { get; set; }
        public Period period { get; set; }
        public List<Location> location { get; set; }
        public ServiceProvider serviceProvider { get; set; }
        public bool? active { get; set; }
        public object name { get; set; }
        public ManagingOrganization managingOrganization { get; set; }
        public object address { get; set; }
        public List<Telecom> telecom { get; set; }
        public string gender { get; set; }
        public string birthDate { get; set; }
        public List<GeneralPractitioner> generalPractitioner { get; set; }
        public ProvidedBy providedBy { get; set; }
        public string intent { get; set; }
        public string priority { get; set; }
        public Context context { get; set; }
        public OccurrencePeriod occurrencePeriod { get; set; }
        public DateTime? authoredOn { get; set; }
        public Requester requester { get; set; }
        public List<Recipient> recipient { get; set; }
        public List<ReasonReference> reasonReference { get; set; }
        public List<SupportingInfo> supportingInfo { get; set; }
        public bool? doNotPerform { get; set; }
        public object code { get; set; }
        public Encounter encounter { get; set; }
        public DateTime? date { get; set; }
        public List<Author> author { get; set; }
        public string title { get; set; }
        public string confidentiality { get; set; }
        public List<Section> section { get; set; }
        public Patient patient { get; set; }
        public List<ConsentingParty> consentingParty { get; set; }
        public List<Action> action { get; set; }
        public object organization { get; set; }
        public string policyRule { get; set; }
        public List<Datum> data { get; set; }
        public string clinicalStatus { get; set; }
        public string verificationStatus { get; set; }
        public Practitioner practitioner { get; set; }
        public Relationship relationship { get; set; }
        public string model { get; set; }
        public string version { get; set; }
        public string mode { get; set; }
        public OrderedBy orderedBy { get; set; }
        public List<Entry> entry { get; set; }
    }

    public sealed class FhirNhs111InboundMessageRegistration
    {
        public string resourceType { get; set; }
        public Meta meta { get; set; }
        public Identifier identifier { get; set; }
        public string type { get; set; }
        public List<Entry> entry { get; set; }
    }

    public sealed class Section
    {
        public List<Section> section { get; set; }
        public string title { get; set; }
        public List<Entry> entry { get; set; }
        public Text text { get; set; }
    }

    public sealed class ServiceProvider
    {
        public string reference { get; set; }
    }

    public sealed class Source
    {
        public string name { get; set; }
        public string endpoint { get; set; }
        public string reference { get; set; }
    }

    public sealed class Subject
    {
        public string reference { get; set; }
    }

    public sealed class SupportingInfo
    {
        public string reference { get; set; }
    }

    public sealed class Telecom
    {
        public string system { get; set; }
        public string value { get; set; }
        public string use { get; set; }
    }

    public sealed class Text
    {
        public string status { get; set; }
        public string div { get; set; }
    }

    public sealed class Type
    {
        public string text { get; set; }
        public List<Coding> coding { get; set; }
    }
}

解决方案

这个异常的核心原因是NHS 111适配器推送的消息采用了AMQP的复杂消息格式,而非MassTransit默认期望的简单二进制/JSON格式,导致MassTransit尝试通过ServiceBusReceivedMessage.Body读取消息时失败。

可以通过以下两种方式解决:

方式一:自定义消息接收上下文(推荐)

在MassTransit配置中,自定义处理AMQP消息的逻辑,直接获取原始AMQP消息并解析:

services.AddMassTransit(x =>
{
    x.AddConsumer<Nhs111ClinicalDocumentConsumer>();

    x.UsingAzureServiceBus((context, cfg) =>
    {
        cfg.Host("sb://[AZURESERVICEBUS]/");

        cfg.ReceiveEndpoint("[QUEUENAME]", e =>
        {
            // 覆盖默认的消息接收逻辑,处理原始AMQP消息
            e.MessageDeserializer = new CustomAmqpMessageDeserializer();
            e.Consumer<Nhs111ClinicalDocumentConsumer>(context);
        });
    });
});

// 自定义AMQP消息反序列化器
public class CustomAmqpMessageDeserializer : IMessageDeserializer
{
    public ContentType ContentType => new ContentType("application/fhir+json");

    public MessageContext Deserialize(ReceiveContext receiveContext)
    {
        var serviceBusReceiveContext = receiveContext as ServiceBusReceiveContext;
        if (serviceBusReceiveContext == null)
            throw new InvalidOperationException("Not a Service Bus receive context");

        // 获取原始AMQP消息
        var amqpMessage = serviceBusReceiveContext.Message.GetRawAmqpMessage();
        
        // 从AMQP消息中提取FHIR JSON数据(根据适配器实际的消息结构调整)
        var bodyStream = amqpMessage.Body.Value.GetStream();
        using var reader = new StreamReader(bodyStream);
        var fhirJson = reader.ReadToEnd();

        // 反序列化为目标类型
        var message = JsonConvert.DeserializeObject<FhirNhs111InboundMessageRegistration>(fhirJson);

        // 创建消息上下文
        var messageContext = new MessageContext(receiveContext, message);
        return messageContext;
    }
}

方式二:修改适配器的消息发送格式

如果可以修改NHS 111适配器的代码,调整其发送Azure Service Bus消息的方式,将FHIR消息以简单二进制JSON格式发送,而非复杂AMQP结构。例如,确保消息直接设置Body为JSON字节数组,而非使用AMQP的高级消息属性。

验证点

  • 确认适配器推送的消息内容格式:通过Azure Service Bus Explorer查看测试套件推送的消息,对比手动发送的消息的Body结构差异
  • 检查消息的Content-Type头:确保适配器发送的消息Content-Type与MassTransit配置的反序列化器匹配

内容的提问来源于stack exchange,提问作者Apexa Shah

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 04:07:01