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
相关产品推荐
相关产品推荐

