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

如何阻止MassTransit的IConsumeMessageObserver拦截RawJson序列化消息?

解决IConsumeMessageObserver不拦截RawJson序列化消息的方案

下面提供几种实用的实现方式,可根据你使用的消息框架(比如MassTransit等)选择适配的方案:

方案1:通过消费者定义元数据标记过滤

在TestRabbitResponseConsumerDefinition中添加自定义元数据标记,然后在观察者的拦截方法里检查该标记,若存在则跳过处理:

步骤1:修改消费者定义添加元数据

public class TestRabbitResponseConsumerDefinition : ConsumerDefinition<TestRabbitResponseConsumer>
{
    public TestRabbitResponseConsumerDefinition()
    {
        // 原有配置
        DisableConsumeTopology = true;
        ClearSerialization();
        UseSerializer<RawJsonSerializer>();
        
        // 添加自定义元数据,标记该消费者使用RawJson序列化
        Metadata.Add("IsRawJsonConsumer", "true");
    }
}

步骤2:在观察者中添加过滤逻辑

public class YourConsumeObserver : IConsumeMessageObserver
{
    public Task PreConsume<T>(ConsumeContext<T> context) where T : class
    {
        // 检查当前消费者是否标记为RawJson类型
        if (context.ConsumerDefinition?.Metadata.TryGetValue("IsRawJsonConsumer", out var flag) == true 
            && flag == "true")
        {
            // 跳过拦截逻辑
            return Task.CompletedTask;
        }

        // 你的原有拦截逻辑
        // ...
        return Task.CompletedTask;
    }

    public Task PostConsume<T>(ConsumeContext<T> context) where T : class
    {
        // 同PreConsume的判断逻辑
        if (context.ConsumerDefinition?.Metadata.TryGetValue("IsRawJsonConsumer", out var flag) == true 
            && flag == "true")
        {
            return Task.CompletedTask;
        }

        // 你的原有拦截逻辑
        // ...
        return Task.CompletedTask;
    }
}

方案2:直接检查消息使用的序列化器类型

如果你的消息框架允许从ConsumeContext中获取当前使用的序列化器实例,可以直接判断序列化器类型来跳过拦截:

public class YourConsumeObserver : IConsumeMessageObserver
{
    public Task PreConsume<T>(ConsumeContext<T> context) where T : class
    {
        // 根据实际框架API获取序列化器,这里以MassTransit为例
        var serializer = context.Serializer;
        if (serializer is RawJsonSerializer || serializer.GetType().FullName.Contains("SystemTextJsonRawSerializerContext"))
        {
            return Task.CompletedTask;
        }

        // 原有拦截逻辑
        // ...
        return Task.CompletedTask;
    }

    // PostConsume方法同理处理
}

方案3:给RawJson序列化的消息添加自定义Header,观察者通过Header过滤

在消费者接收消息时(或者自定义RawJson序列化器中)添加专属Header,观察者检查到该Header则跳过处理:

步骤1:添加自定义Header

可以在消费者的Consume方法开头添加,或者在自定义序列化器的反序列化方法中添加:

public class TestRabbitResponseConsumer : IConsumer<Test>
{
    public async Task Consume(ConsumeContext<Test> context)
    {
        // 添加标记Header
        context.Headers.Set("X-RawJsonSerialized", true);
        
        // 你的原有消费逻辑
        // ...
    }
}

步骤2:观察者中检查Header

public class YourConsumeObserver : IConsumeMessageObserver
{
    public Task PreConsume<T>(ConsumeContext<T> context) where T : class
    {
        if (context.Headers.TryGetHeader("X-RawJsonSerialized", out var value) 
            && value is bool isRaw && isRaw)
        {
            return Task.CompletedTask;
        }

        // 原有拦截逻辑
        // ...
        return Task.CompletedTask;
    }

    // PostConsume方法同理处理
}

方案4:创建过滤基类,统一处理跳过逻辑

如果有多个观察者需要过滤这类消息,可以写一个基类封装判断逻辑,其他观察者继承即可:

public abstract class RawJsonExcludedObserverBase : IConsumeMessageObserver
{
    protected virtual bool ShouldSkip(ConsumeContext context)
    {
        // 这里可以选用上面任意一种判断逻辑,比如元数据判断
        return context.ConsumerDefinition?.Metadata.TryGetValue("IsRawJsonConsumer", out var flag) == true 
               && flag == "true";
    }

    public virtual Task PreConsume<T>(ConsumeContext<T> context) where T : class
    {
        if (ShouldSkip(context)) return Task.CompletedTask;
        return PreConsumeCore(context);
    }

    public virtual Task PostConsume<T>(ConsumeContext<T> context) where T : class
    {
        if (ShouldSkip(context)) return Task.CompletedTask;
        return PostConsumeCore(context);
    }

    // 子类实现具体的拦截逻辑
    protected abstract Task PreConsumeCore<T>(ConsumeContext<T> context) where T : class;
    protected abstract Task PostConsumeCore<T>(ConsumeContext<T> context) where T : class;
}

然后你的观察者继承这个基类:

public class YourConsumeObserver : RawJsonExcludedObserverBase
{
    protected override Task PreConsumeCore<T>(ConsumeContext<T> context) where T : class
    {
        // 你的原有拦截逻辑
        // ...
        return Task.CompletedTask;
    }

    protected override Task PostConsumeCore<T>(ConsumeContext<T> context) where T : class
    {
        // 你的原有拦截逻辑
        // ...
        return Task.CompletedTask;
    }
}

内容的提问来源于stack exchange,提问作者Никита Викулов

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 09:50:27