如何用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
相关产品推荐
相关产品推荐

