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

MassTransit集成测试问题:TestHarness未消费消息

MassTransit集成测试消费消息失败排查求助

基于Web Application Factory搭建MassTransit集成测试,预期实现消费消息后在数据库创建记录的逻辑。已确认TestHarness成功发布消息,但始终无法消费消息,IsMessageConsumedByService函数持续返回false,调试时消费者测试工具显示消费消息数为0。相关代码如下,恳请技术帮助。

Web Application Factory代码

public class CustomWebApplicationFactory
    : WebApplicationFactory<Program>, IAsyncLifetime
{
    private readonly MsSqlContainer _mssqlContainer = new MsSqlBuilder().Build();
    private readonly RabbitMqContainer _rabbitMqContainer = new RabbitMqBuilder().Build();
    protected override void ConfigureWebHost(IWebHostBuilder builder)
    {
        builder.ConfigureTestServices(services =>
        {
            services.RemoveAll<TenantStoreDbContext>();
            services.RemoveAll<DbConnection>();
            services.RemoveAll(typeof(DbContextOptions<TenantStoreDbContext>));

            services.AddDbContext<TenantStoreDbContext>((container, options) =>
            {
                options.UseSqlServer(_mssqlContainer.GetConnectionString());
            });

            services.AddMassTransitTestHarness(x =>
            {
                x.AddConsumer<StripeSubscriptionCreatedConsumer>();
                // ToDo: Test with RabbitMq container after in memory works
            });
        });

        builder.UseEnvironment("Local");
    }
    public async Task InitializeAsync() 
    {
        await _mssqlContainer.StartAsync();
        await _rabbitMqContainer.StartAsync();
    }

    public new async Task DisposeAsync() 
    {
        await _mssqlContainer.DisposeAsync();
        await _rabbitMqContainer.DisposeAsync();
    } 
}

测试类代码

public class NewTenantTests : IClassFixture<CustomWebApplicationFactory>
    {
        private readonly CustomWebApplicationFactory _factory;

        public NewTenantTests(CustomWebApplicationFactory factory)
        {
            _factory = factory;
        }

        [Fact]
        public async void When_stripe_created_event_is_recieved_a_new_tenant_is_created()
        {
            // Arrange
            var serviceProvider = _factory.Services.GetRequiredService<IServiceProvider>();
            var testHarness = await StartTestHarness(serviceProvider);
            var subscriptionEvent = BuildStripeSubscriptionCreated();

            // Act
            await PublishMessage(subscriptionEvent, testHarness);

            // Assert
            (await IsMessageConsumedByService(testHarness, subscriptionEvent)).ShouldBe(true);
            (await IsMessageConsumedByConsumer(serviceProvider, subscriptionEvent)).ShouldBe(true);
            (await AssertNewTenantWasCreated(serviceProvider, subscriptionEvent.TenantIdentifier)).ShouldBe(true);

        }

        private static async Task<ITestHarness> StartTestHarness(IServiceProvider provider)
        {
            var harness = provider.GetTestHarness();
            await harness.Start();

            return harness;
        }

        private static StripeSubscriptionCreatedEvent BuildStripeSubscriptionCreated()
        {
            var subscriptionEvent = new StripeSubscriptionCreatedEvent()
            {
                SubscriptionId = "abc",
                CustomerId = "abc",
                Description = "description",
                CurrentPeriodStart = DateTime.Now,
                CurrentPeriodEnd = DateTime.Now.AddMonths(1),
                TenantIdentifier = Guid.NewGuid().ToString()
            };

            return subscriptionEvent;
        }

        private static async Task PublishMessage(StripeSubscriptionCreatedEvent subscriptionCreated, ITestHarness harness)
        {
            await harness.Bus.Publish(subscriptionCreated);
        }

        private static async Task<bool> IsMessageConsumedByConsumer(IServiceProvider serviceProvider,
            StripeSubscriptionCreatedEvent subscriptionCreated)
        {
            return await serviceProvider.
                GetRequiredService<IConsumerTestHarness<StripeSubscriptionCreatedConsumer>>().
                Consumed.Any<StripeSubscriptionCreatedEvent>(x => x.Context.Message.EventId == subscriptionCreated.EventId);
        }

        private static async Task<bool> IsMessageConsumedByService(ITestHarness harness,
            StripeSubscriptionCreatedEvent subscriptionCreated)
        {
            var consumerHarnes = harness.GetConsumerHarness<StripeSubscriptionCreatedConsumer>();
            return await consumerHarnes.Consumed.Any<StripeSubscriptionCreatedEvent>(x => x.Context.Message.EventId == subscriptionCreated.EventId);
        }

        private static async Task<bool> AssertNewTenantWasCreated(IServiceProvider serviceProvider, string tenantIdentifier)
        {
            var repository = serviceProvider.GetRequiredService<ITenantStoreRepository>();
            var identifier = TenantIdentifier.Create(tenantIdentifier).Value;
            var tenant = await repository.GetTenantByIdentifierAsync(identifier);

            return tenant != null;
        }
    }

排查与解决方案

  • 配置消息与消费者的路由:当前AddMassTransitTestHarness仅添加了消费者,但未配置消息路由。测试使用内存传输时,需显式配置端点:
    services.AddMassTransitTestHarness(x =>
    {
        x.AddConsumer<StripeSubscriptionCreatedConsumer>();
        x.UsingInMemory((context, cfg) =>
        {
            cfg.ConfigureEndpoints(context);
        });
    });
    
  • 修复测试方法返回类型:将测试方法的async void改为async Task,否则测试框架无法正确等待异步操作完成:
    [Fact]
    public async Task When_stripe_created_event_is_recieved_a_new_tenant_is_created()
    {
        // 原有逻辑不变
    }
    
  • 修正消息匹配条件:BuildStripeSubscriptionCreated方法未给EventId赋值,导致断言时无法匹配。要么给EventId设置唯一值,要么改用已赋值的TenantIdentifier作为判断条件:
    // 修改断言逻辑
    return await consumerHarnes.Consumed.Any<StripeSubscriptionCreatedEvent>(x => 
        x.Context.Message.TenantIdentifier == subscriptionCreated.TenantIdentifier);
    
  • 等待消费完成:发布消息后,添加等待逻辑确保消费者完成处理:
    await PublishMessage(subscriptionEvent, testHarness);
    // 等待指定消费者处理消息
    await testHarness.WaitForConsumer<StripeSubscriptionCreatedConsumer>();
    
  • 检查消费者依赖注入:确认StripeSubscriptionCreatedConsumer的所有依赖(如ITenantStoreRepository、TenantStoreDbContext)在测试服务中已正确注册,避免因依赖缺失导致消费者初始化失败。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 20:05:22