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

MassTransit测试中请求超时问题排查

问题

自研API可正常运行,但测试环节出现异常:测试方法通过MediatR发送命令后,命令能抵达Handler,但Handler内通过MassTransit发送的请求始终无法获取响应。


测试类代码

[TestClass]
public class UnbindVisilibityTests : InfrastructureDataTestBase
{
    private readonly InfrastructureDataTestDal _infrastructureDal = new InfrastructureDataTestDal();
    private ITestHarness _harness;
    private ServiceProvider _provider;

    UnbindVisibilityCommand command;

    int idVis;

    [TestInitialize]
    public void Init()
    {
        _provider = new ServiceCollection()
            .AddMassTransitTestHarness(cfg =>
            {
                cfg.AddConsumer<UnbindUsersFromVisibilityFilterRequestConsumer>();
                cfg.AddHandler<UnbindVisibilityCommand>(context => context.RespondAsync(true));
            })
            .BuildServiceProvider(true);

        _harness = _provider.GetRequiredService<ITestHarness>();

        idVis = _infrastructureDal.InsertVisibility("UnbindVisibilityCommandTest", "TestAppl");

        command = new UnbindVisibilityCommand
        {
            Id = idVis,
        };
    }

    [TestMethod]
    public async Task Unbind_OK()
    {
        await _harness.Start();

        var result = Mediator.Send(command).Result;

        Assert.IsTrue(result is SuccessCommandResult);

        await _harness.Stop();
    }

    [TestCleanup]
    public void CleanUpTest()
    {
        _infrastructureDal.DeleteVisibility(new int[] { idVis });
    }
}

Handler代码

public class UnbindVisibilityCommandHandler : IRequestHandler<UnbindVisibilityCommand, CommandResult>
{
    private readonly IClientFactory _clientFactory;

    public UnbindVisibilityCommandHandler(IClientFactory clientFactory)
    {
        _clientFactory = clientFactory;
    }

    public async Task<CommandResult> Handle(UnbindVisibilityCommand request, CancellationToken cancellationToken)
    {
        var client = _clientFactory.CreateRequestClient<IUnbindUsersFromVisibilityFilterRequest>();

        var response = await client.GetResponse<IUnbindUsersFromVisibilityFilterResult>(
        new
        {
            Id = request.Id
        });

        return new SuccessCommandResult();
    }
}

消费者代码

public class UnbindUsersFromVisibilityFilterRequestConsumer : IConsumer<IUnbindUsersFromVisibilityFilterRequest>
{
    private readonly IApplicativeUsersCommandDal _dal;

    public UnbindUsersFromVisibilityFilterRequestConsumer(IApplicativeUsersCommandDal dal)
    {
        _dal = dal;
    }

    public async Task Consume(ConsumeContext<IUnbindUsersFromVisibilityFilterRequest> context)
    {
        await _dal.UnbindUsersFromVisibilityAsync(context.Message.Id);

        await context.RespondAsync<IUnbindUsersFromVisibilityFilterResult>(new { });
    }
}

测试基类(InfrastructureDataTestBase)代码

public abstract class InfrastructureDataTestBase : TestBase
{
    internal static IMediator Mediator;
    protected IServiceProvider serviceProvider;
    static bool _isInitialized;
    public IConfiguration configuration;

    public InfrastructureDataTestBase()
    {
        var builder = new ConfigurationBuilder()
              .SetBasePath(Directory.GetCurrentDirectory() + ".\\InfrastructureData.Tests")
              .AddJsonFile("appsettings.json", optional: false, reloadOnChange: true)
              .AddJsonFile("TokenVerificationSettings.json", optional: false, reloadOnChange: true)
              .AddEnvironmentVariables();

        configuration = builder.Build();

        if (_isInitialized == false)
        {
            serviceProvider = GetServiceProvider(configuration);

            Mediator = serviceProvider.GetService<IMediator>();
        }
    }

    public IServiceProvider GetServiceProvider(IConfiguration config)
    {
        ServiceCollection services = new ServiceCollection();

        services.AddControllers();

        DbProviderFactories.RegisterFactory("System.Data.SqlClient", SqlClientFactory.Instance);

        services.AddSwaggerApiVersioning<ConfigureSwaggerOptions>();

        var connProvider = ConnectionStringProviderFactory.CreateInstance(config.GetValue<string>("ConnectionStringFullFilePath"), config.GetValue<string>("RsaKeyFullFilePath"));

        services.AddSingleton(connProvider);

        services.AddLogging();

        services.AddSingleton<IConfiguration>(config);

        services.AddPagingProvider();

        services.AddHealthChecks();

        services.AddOptions(config);

        services.AddAccessTokenValidation();

        services.AddIdentityProvider(configuration.GetValue<string>("IDP.BaseUri"));

        services.AddApplicationLayer();

        services.AddDataInfrastructureLayer(configuration, connProvider);

        ConfigureServices(services, configuration);

        MessageBrokerOptions messageBrokerOptions = new();

        configuration.GetSection(MessageBrokerOptions.SectionName).Bind(messageBrokerOptions);

        services.AddMassTransit(x =>
        {                
            x.AddConsumer<UnbindUsersFromVisibilityFilterRequestConsumer>();

            x.UsingRabbitMq((context, cfg) =>
            {
                cfg.Host(messageBrokerOptions.HostName, (ushort)messageBrokerOptions.Port, messageBrokerOptions.VirtualHost, h =>
                {
                    var rsaFile = config.GetValue<string>("RsaKeyFullFilePath");

                    var cryptoService = string.IsNullOrWhiteSpace(rsaFile)
                        ? Digitronica.IT.Infrastructure.Security.Cryptography.CryptoService.CreateDefault()
                        : Digitronica.IT.Infrastructure.Security.Cryptography.CryptoService.CreateFromXmlFile(rsaFile);

                    h.Username(messageBrokerOptions.Username);
                    h.Password(cryptoService.Decrypt(messageBrokerOptions.Password));
                });

                cfg.ConfigureEndpoints(context);
            });

            x.AddRequestClient<IUnbindUsersFromVisibilityFilterRequest>();
        });

        _isInitialized = true;

        return services.BuildServiceProvider();
    }
}

TestBase类代码

public abstract class TestBase
{
    public void ConfigureServices(IServiceCollection services, IConfiguration config)
    {
        services.AddControllers();

        DbProviderFactories.RegisterFactory("System.Data.SqlClient", SqlClientFactory.Instance);

        services.AddSwaggerApiVersioning<ConfigureSwaggerOptions>();

        services.AddConnectionStringProvider(config.GetValue<string>("ConnectionStringFullFilePath"), config.GetValue<string>("RsaKeyFullFilePath"));

        services.AddLogging();

        services.AddSingleton<IConfiguration>(config);

        services.AddPagingProvider();

        services.AddHealthChecks();

        services.AddOptions(config);

        services.AddAccessTokenValidation();
    }
}

问题原因

测试类中单独创建了一个包含MassTransit测试Harness的ServiceProvider,但测试用的Mediator是从基类InfrastructureDataTestBase的独立容器中获取的,两个容器完全隔离:

  • 基类容器配置了真实的RabbitMQ连接,Handler中的IClientFactory来自这个容器,请求会发送到真实RabbitMQ实例
  • 测试Harness的消费者仅注册在测试类自己创建的容器中,无法处理真实RabbitMQ中的请求

最终导致Handler发送的MassTransit请求没有对应的消费者处理,无法收到响应。


解决方案

方案1:复用基类容器集成测试Harness

删除测试类中单独创建的ServiceProvider,在基类容器基础上配置MassTransit测试Harness,确保所有组件在同一容器中:

[TestClass]
public class UnbindVisilibityTests : InfrastructureDataTestBase
{
    private readonly InfrastructureDataTestDal _infrastructureDal = new InfrastructureDataTestDal();
    private ITestHarness _harness;

    UnbindVisibilityCommand command;
    int idVis;

    [TestInitialize]
    public void Init()
    {
        _harness = serviceProvider.GetRequiredService<ITestHarness>();

        idVis = _infrastructureDal.InsertVisibility("UnbindVisibilityCommandTest", "TestAppl");

        command = new UnbindVisibilityCommand
        {
            Id = idVis,
        };
    }

    [TestMethod]
    public async Task Unbind_OK()
    {
        await _harness.Start();

        // 使用await替代.Result避免死锁
        var result = await Mediator.Send(command);

        Assert.IsTrue(result is SuccessCommandResult);

        await _harness.Stop();
    }

    [TestCleanup]
    public void CleanUpTest()
    {
        _infrastructureDal.DeleteVisibility(new int[] { idVis });
    }
}

方案2:基类中切换测试环境配置

在基类的GetServiceProvider方法中,替换真实RabbitMQ配置为测试Harness:

public IServiceProvider GetServiceProvider(IConfiguration config)
{
    ServiceCollection services = new ServiceCollection();

    // ... 其他服务配置保持不变 ...

    // 替换为测试Harness
    services.AddMassTransitTestHarness(x =>
    {
        x.AddConsumer<UnbindUsersFromVisibilityFilterRequestConsumer>();
        x.AddRequestClient<IUnbindUsersFromVisibilityFilterRequest>();
    });

    _isInitialized = true;

    return services.BuildServiceProvider();
}

额外优化

  • 测试方法中使用await Mediator.Send(command)替代.Result,避免异步操作死锁
  • 确保消费者依赖的IApplicativeUsersCommandDal在测试容器中已注册(可使用Mock替代真实实现)

内容的提问来源于stack exchange,提问作者Francesco Vargas

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 02:02:01