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

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的作用域中。以下是具体修复步骤:

  1. 调整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)
    {
    }
}
  1. 注册必要的基础服务
    在DI容器中添加上下文访问器的注册:
services.AddHttpContextAccessor();
services.AddScoped<IConsumeContextAccessor, ConsumeContextAccessor>();
services.AddScoped<MyDependency>();
  1. 调整过滤器注册顺序
    确保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);
    });
});
  1. 消费中发送消息的正确方式
    在消费者中优先使用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üş

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 10:14:57