如何在MassTransit中实现消息嗅探器(Sniffer)
MassTransit 消息嗅探器:如何从IntegrationEvent获取具体事件类型
问题背景
我有一套消息层级结构,想要实现一个能处理所有不同消息的消费者作为消息嗅探器,最初的代码实现如下:
public class Sniffer : IConsumer<IntegrationEvent> { readonly ILogger<Sniffer> _logger; public Sniffer(ILogger<Sniffer> logger) { _logger = logger; } public Task Consume(ConsumeContext<IntegrationEvent> context) { IntegrationEvent data = context.Message; _logger.LogInformation($"EVENT: ***{data.Code} - {data.DateTime}***"); if (data.Code == "Event1") { DataEvent1 dataEvent1 = ?????? _logger.LogInformation($"EVENT1: ***{dataEvent1.DataEvent1}***"); } else if (data.Code == "Event2") { DataEvent2 dataEvent2 = ?????? _logger.LogInformation($"EVENT2: ***{dataEvent2.DataEvent2}***"); } return Task.CompletedTask; } }
目前遇到的问题是:消费者中获取的序列化对象仅为IntegrationEvent类型,无法直接转换为DataEvent1或DataEvent2,想知道是否有办法获取具体事件对象,或者是否需要替换默认的序列化/反序列化器。
解决方案
方法1:利用ConsumeContext直接获取具体消息类型
MassTransit的ConsumeContext提供了TryGetMessage<T>方法,可以直接尝试从上下文获取指定类型的消息对象,无需手动反序列化:
public Task Consume(ConsumeContext<IntegrationEvent> context) { var data = context.Message; _logger.LogInformation($"EVENT: ***{data.Code} - {data.DateTime}***"); if (data.Code == "Event1") { if (context.TryGetMessage<DataEvent1>(out var messageWrapper)) { _logger.LogInformation($"EVENT1: ***{messageWrapper.Message.DataEvent1}***"); } } else if (data.Code == "Event2") { if (context.TryGetMessage<DataEvent2>(out var messageWrapper)) { _logger.LogInformation($"EVENT2: ***{messageWrapper.Message.DataEvent2}***"); } } return Task.CompletedTask; }
如果该方法失效,还可以读取原始消息内容手动反序列化(以System.Text.Json为例):
var rawContent = await context.GetBodyAsString(); var dataEvent1 = JsonSerializer.Deserialize<DataEvent1>(rawContent, new JsonSerializerOptions { PropertyNameCaseInsensitive = true });
方法2:配置序列化器保留类型信息
默认情况下,MassTransit的序列化器不会保留完整的类型信息,导致反序列化时只能得到基类对象。可以通过配置序列化器添加类型标识,让框架自动识别具体类型。
以System.Text.Json为例,在Bus配置中添加类型信息支持:
services.AddMassTransit(x => { x.AddConsumer<Sniffer>(); x.UsingRabbitMq((context, cfg) => { cfg.UseSystemTextJsonSerializer(options => { options.TypeInfoResolver = new DefaultJsonTypeInfoResolver { Modifiers = { typeInfo => { if (typeInfo.Type.IsAssignableTo(typeof(IntegrationEvent))) { typeInfo.PolymorphismOptions = new JsonPolymorphismOptions { TypeDiscriminatorPropertyName = "$type", DerivedTypes = { new JsonDerivedType(typeof(DataEvent1), nameof(DataEvent1)), new JsonDerivedType(typeof(DataEvent2), nameof(DataEvent2)) } }; } } } }; }); cfg.ConfigureEndpoints(context); }); });
配置完成后,context.Message会自动映射为具体的事件类型,直接用类型判断即可:
public Task Consume(ConsumeContext<IntegrationEvent> context) { var data = context.Message; _logger.LogInformation($"EVENT: ***{data.Code} - {data.DateTime}***"); if (data is DataEvent1 dataEvent1) { _logger.LogInformation($"EVENT1: ***{dataEvent1.DataEvent1}***"); } else if (data is DataEvent2 dataEvent2) { _logger.LogInformation($"EVENT2: ***{dataEvent2.DataEvent2}***"); } return Task.CompletedTask; }
注意事项
- 若使用Newtonsoft.Json序列化器,需开启
TypeNameHandling.Auto并配置类型映射,逻辑与上述一致。 - 确保
DataEvent1、DataEvent2均继承自IntegrationEvent,且Code字段的取值能准确对应消息类型。
内容的提问来源于stack exchange,提问作者Jorge
相关产品推荐
相关产品推荐

