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

