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

MassTransit消费者单元测试:Consumed断言失败及SaveChangesAsync测试

问题解决:MassTransit单元测试中Consumed断言失败及EF相关测试方案

我用MassTransit、Entity Framework和C#开发项目,消费者通过Consume方法消费MyEvent事件并插入数据到数据库。单元测试里Published<MyEvent>断言成功,但Consumed<MyEvent>始终失败。同时需要了解如何Mock消费者以及编写SaveChangesAsync的单元测试用例。

消费者原代码

public async Task Consume(ConsumeContext<MyEvent> context)
{
    private readonly MyDbContext _dbContext;  

    try
    {
        // Here logic to insert record in to new database
        var data = new MyService.TableNmae()
        {
            Id = context.Message.MyId,
            Description = "test data"
        };

        _ = _dbContext.TableName.AddAsync(data);
        _ = _dbContext.SaveChangesAsync(context.CancellationToken);
    }
    catch (Exception ex)
    {
        _logger.LogCritical($"{GetType().Name}:{nameof(Consume)} {ex}");
    }

    await Task.CompletedTask;
}

单元测试原代码

private ITestHarness _testHarness;

[SetUp]
public void Initialize()
{
    var serviceCollection = new ServiceCollection();

    serviceCollection.AddMassTransitTestHarness(busRegistrationConfigurator =>
    {
        busRegistrationConfigurator.AddConsumer<MyConsumer>();
    });

    var serviceProvider = serviceCollection.BuildServiceProvider();

    _testHarness = serviceProvider.GetRequiredService<ITestHarness>();
}

[Test]
public async Task TestMethod1()
{
    await _testHarness.Start();
    await _testHarness.Bus.Publish(new MyEvent { Code = "H"});

    Assert.That(await _testHarness.Published.Any<MyEvent>(), Is.True);
    Assert.That(await _testHarness.Consumed.Any<MyEvent>(), Is.True);
}

一、Consumed断言失败的原因及修复

1. 消费者代码的核心问题

原消费者代码存在语法和逻辑错误:

  • 缺少构造函数:消费者类必须通过构造函数注入MyDbContext和ILogger,直接在Consume方法里声明私有字段会导致无法实例化消费者。
  • 异步操作未等待:用_ =丢弃AddAsync和SaveChangesAsync的任务,没有await,会导致Consume方法提前结束,数据库操作可能还未完成甚至抛出异常被静默吞掉。
  • 异常未重新抛出:catch块只记录日志但不抛异常,MassTransit无法感知消费失败,会导致消息被静默处理,Consumed统计只记录成功消费的消息。

修复后的消费者代码:

public class MyConsumer : IConsumer<MyEvent>
{
    private readonly MyDbContext _dbContext;
    private readonly ILogger<MyConsumer> _logger;

    // 构造函数注入依赖
    public MyConsumer(MyDbContext dbContext, ILogger<MyConsumer> logger)
    {
        _dbContext = dbContext;
        _logger = logger;
    }

    public async Task Consume(ConsumeContext<MyEvent> context)
    {
        try
        {
            var data = new MyService.TableName() // 修正拼写错误TableNmae→TableName
            {
                Id = context.Message.MyId,
                Description = "test data"
            };

            await _dbContext.TableName.AddAsync(data, context.CancellationToken);
            await _dbContext.SaveChangesAsync(context.CancellationToken);
        }
        catch (Exception ex)
        {
            _logger.LogCritical(ex, $"{GetType().Name}:{nameof(Consume)}");
            // 重新抛出异常,让MassTransit处理失败逻辑(比如重试)
            throw;
        }
    }
}

2. 单元测试的依赖注入问题

原测试未注册MyDbContext,消费者无法实例化,导致消息无法被消费。可以用EF内存数据库或Mock来补充依赖:

修改后的测试代码(使用内存数据库):

private ITestHarness _testHarness;

[SetUp]
public void Initialize()
{
    var serviceCollection = new ServiceCollection();

    // 注册内存数据库的DbContext
    serviceCollection.AddDbContext<MyDbContext>(options =>
    {
        options.UseInMemoryDatabase(Guid.NewGuid().ToString());
    });

    serviceCollection.AddMassTransitTestHarness(busRegistrationConfigurator =>
    {
        busRegistrationConfigurator.AddConsumer<MyConsumer>();
    });

    var serviceProvider = serviceCollection.BuildServiceProvider();

    _testHarness = serviceProvider.GetRequiredService<ITestHarness>();
}

[Test]
public async Task TestMethod1()
{
    await _testHarness.Start();
    try
    {
        await _testHarness.Bus.Publish(new MyEvent { MyId = Guid.NewGuid(), Code = "H" });

        // 等待消息被消费,增加超时判断避免异步延迟
        Assert.That(await _testHarness.Published.Any<MyEvent>(), Is.True);
        Assert.That(await _testHarness.Consumed.Any<MyEvent>(x => x.Context.Message.Code == "H"), Is.True, "消息未被消费者处理");
    }
    finally
    {
        await _testHarness.Stop();
    }
}

二、Mock消费者及测试SaveChangesAsync

1. Mock消费者依赖(DbContext)

如果不想用内存数据库,可用Moq Mock DbContext和DbSet,验证数据库操作是否被正确调用:

[Test]
public async Task Consumer_Should_Add_Data_And_Save_Changes()
{
    // Mock DbSet
    var mockSet = new Mock<DbSet<MyService.TableName>>();
    // Mock DbContext
    var mockContext = new Mock<MyDbContext>();
    mockContext.Setup(c => c.TableName).Returns(mockSet.Object);

    // Mock Logger
    var mockLogger = new Mock<ILogger<MyConsumer>>();

    // 创建消费者实例
    var consumer = new MyConsumer(mockContext.Object, mockLogger.Object);

    // 构造测试用的ConsumeContext
    var myEvent = new MyEvent { MyId = Guid.NewGuid() };
    var context = new Mock<ConsumeContext<MyEvent>>();
    context.Setup(c => c.Message).Returns(myEvent);
    context.Setup(c => c.CancellationToken).Returns(CancellationToken.None);

    // 调用Consume方法
    await consumer.Consume(context.Object);

    // 验证AddAsync是否被调用,参数是否匹配
    mockSet.Verify(s => s.AddAsync(
        It.Is<MyService.TableName>(d => d.Id == myEvent.MyId && d.Description == "test data"),
        It.IsAny<CancellationToken>()), Times.Once);

    // 验证SaveChangesAsync是否被调用
    mockContext.Verify(c => c.SaveChangesAsync(It.IsAny<CancellationToken>()), Times.Once);
}

2. 测试SaveChangesAsync异常场景

模拟数据库保存失败,验证消费者是否正确记录日志并抛出异常:

[Test]
public async Task Consumer_Should_Log_And_Rethrow_Exception_When_Save_Fails()
{
    var mockSet = new Mock<DbSet<MyService.TableName>>();
    var mockContext = new Mock<MyDbContext>();
    mockContext.Setup(c => c.TableName).Returns(mockSet.Object);
    
    // 模拟SaveChangesAsync抛出异常
    var testException = new InvalidOperationException("Test save failure");
    mockContext.Setup(c => c.SaveChangesAsync(It.IsAny<CancellationToken>()))
        .ThrowsAsync(testException);

    var mockLogger = new Mock<ILogger<MyConsumer>>();
    var consumer = new MyConsumer(mockContext.Object, mockLogger.Object);

    var context = new Mock<ConsumeContext<MyEvent>>();
    context.Setup(c => c.Message).Returns(new MyEvent { MyId = Guid.NewGuid() });

    // 验证是否抛出异常
    var exception = Assert.ThrowsAsync<InvalidOperationException>(() => consumer.Consume(context.Object));
    Assert.That(exception, Is.EqualTo(testException));

    // 验证日志是否被记录
    mockLogger.Verify(
        x => x.Log(
            LogLevel.Critical,
            It.IsAny<EventId>(),
            It.Is<It.IsAnyType>((v, t) => v.ToString().Contains(nameof(MyConsumer))),
            testException,
            It.IsAny<Func<It.IsAnyType, Exception, string>>()),
        Times.Once);
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 04:55:23