MassTransit消费时发消息,SendFilter无法获取ConsumeContext的解决办法
问题描述
我遇到的问题已被讨论并修复,但不知道如何应用到自身代码中。
场景概述
我需要在发送每个事件前为其头部添加Token:
- 通过MassTransit的
Send Scoped Filter可从HttpContext获取Token并添加到头部 - 消费时能通过
ConsumeContext或Consume Scoped Filter获取该Token
但在消费过程中发送另一个事件时,SendFilter中无法获取ConsumeContext,抛出MissingConsumeContext Exception。我尝试为SendFilter和ConsumeFilter创建作用域对象,但未生效。官方有组合消费与发送过滤器的相关文档,还有一个场景完全一致的示例仓库,唯一区别是我未使用请求/响应客户端发送消息。
我的代码
TokenSendFilter
public class TokenSendFilter<T> : IFilter<SendContext<T>> where T : class { private MyDependency myDependency; private readonly IWorkContext _workcontext; public TokenSendFilter(MyDependency dependency, IWorkContext workcontext) { myDependency = dependency; _workcontext = workcontext; } public Task Send(SendContext<T> context, IPipe<SendContext<T>> next) { if (!string.IsNullOrWhiteSpace(_workcontext.Token)) { myDependency.Token = _workcontext.Token; context.Headers.Set("Token", _workcontext.Token); } else if((!string.IsNullOrWhiteSpace(myDependency.Token))) { context.Headers.Set(key: "Token", myDependency.Token); } return next.Send(context); } public void Probe(ProbeContext context) { } }
TokenConsumeFilter
public class TokenConsumeFilter<T> : IFilter<ConsumeContext<T>> where T : class { private MyDependency myDependency; public TokenConsumeFilter(MyDependency dependency) { myDependency = dependency; } public Task Send(ConsumeContext<T> context, IPipe<ConsumeContext<T>> next) { if (context.Headers.TryGetHeader("Token", out object token)) { myDependency.Token = (string)token; } return next.Send(context); } public void Probe(ProbeContext context) { } }
说明:IWorkContext是获取HttpContext头部值的抽象。
注册作用域过滤器
services.AddScoped<MyDependency>(); services.AddMassTransit(x => { x.AddConsumer<MyConsumer>(); x.UsingRabbitMq((context, cfg) => { cfg.UseSendFilter(typeof(TokenSendFilter<>), context); cfg.UseConsumeFilter(typeof(TokenConsumeFilter<>), context); cfg.ConfigureEndpoints(context); }); });
发送事件代码
var address = new Uri($"exchange:{exchangeName}"); var endpoint = await _bus.GetSendEndpoint(address); await endpoint.Send(The Message);
尝试过的方法
我试过从ConsumeContext复制头部到SendContext,但在SendContext Filter中未找到ConsumeContext负载。我认为SendFilter不在ConsumeContext的作用域内,因此需要在发送另一个事件时获取ConsumeContext对象。
使用的MassTransit版本为8.0.5。
解决方案
核心问题是消费过程中发送消息时,当前的ConsumeContext未被正确传递到SendFilter的作用域中。以下是具体修复步骤:
- 调整TokenSendFilter的依赖逻辑
移除IWorkContext的直接注入,改用IConsumeContextAccessor和IHttpContextAccessor覆盖两种场景(消费中发消息、API请求中发消息),修改后的代码如下:
public class TokenSendFilter<T> : IFilter<SendContext<T>> where T : class { private readonly MyDependency _myDependency; private readonly IHttpContextAccessor _httpContextAccessor; private readonly IConsumeContextAccessor _consumeContextAccessor; public TokenSendFilter(MyDependency myDependency, IHttpContextAccessor httpContextAccessor, IConsumeContextAccessor consumeContextAccessor) { _myDependency = myDependency; _httpContextAccessor = httpContextAccessor; _consumeContextAccessor = consumeContextAccessor; } public Task Send(SendContext<T> context, IPipe<SendContext<T>> next) { string token = null; // 优先从消费上下文取Token(消费过程中发消息场景) var consumeContext = _consumeContextAccessor.ConsumeContext; if (consumeContext != null && consumeContext.Headers.TryGetHeader("Token", out object consumeToken)) { token = consumeToken as string; } // 其次从HttpContext取(API请求发消息场景) else if (_httpContextAccessor.HttpContext != null && !string.IsNullOrWhiteSpace(_httpContextAccessor.HttpContext.Request.Headers["Token"])) { token = _httpContextAccessor.HttpContext.Request.Headers["Token"]; } // 最后从作用域依赖取 else if (!string.IsNullOrWhiteSpace(_myDependency.Token)) { token = _myDependency.Token; } if (!string.IsNullOrWhiteSpace(token)) { _myDependency.Token = token; context.Headers.Set("Token", token); } return next.Send(context); } public void Probe(ProbeContext context) { } }
- 注册必要的基础服务
在DI容器中添加上下文访问器的注册:
services.AddHttpContextAccessor(); services.AddScoped<IConsumeContextAccessor, ConsumeContextAccessor>(); services.AddScoped<MyDependency>();
- 调整过滤器注册顺序
确保ConsumeFilter先于SendFilter注册,保证作用域内的Token先被消费过滤器初始化:
services.AddMassTransit(x => { x.AddConsumer<MyConsumer>(); x.UsingRabbitMq((context, cfg) => { cfg.UseConsumeFilter(typeof(TokenConsumeFilter<>), context); cfg.UseSendFilter(typeof(TokenSendFilter<>), context); cfg.ConfigureEndpoints(context); }); });
- 消费中发送消息的正确方式
在消费者中优先使用ConsumeContext.Send方法发送消息,而非直接从IBus获取端点,这样能自动传递当前消费上下文:
public class MyConsumer : IConsumer<YourMessage> { public async Task Consume(ConsumeContext<YourMessage> context) { // 直接通过ConsumeContext发送,自动传递上下文 await context.Send(new Uri($"exchange:{exchangeName}"), new AnotherMessage()); } }
关键原理
IConsumeContextAccessor会在消费过程中自动注入当前的ConsumeContext,SendFilter可通过它获取消费上下文的Token- 优先级顺序覆盖了两种核心场景,确保Token传递不中断
- 使用
ConsumeContext.Send能避免手动获取端点导致的上下文丢失问题
内容的提问来源于stack exchange,提问作者Haydar Can Kubilay Gümüş
相关产品推荐
相关产品推荐

