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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 03:08:22