如何阻止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,提问作者Никита Викулов
相关产品推荐
相关产品推荐

