MassTransit含IBus依赖的消费者单元测试问题排查
MassTransit消费者单元测试问题及优化
消费者代码
public sealed class QuotationMessageConsumer : IConsumer<SendQuotationMessage> { private readonly ILogger<QuotationMessageConsumer> _logger; private readonly IBus _bus; public QuotationMessageConsumer(ILogger<QuotationMessageConsumer> logger, IBus bus) { _logger = logger; _bus = bus; } public async Task Consume(ConsumeContext<SendQuotationMessage> context) { _logger?.LogInformation("Sending message to topic..."); var request = context.Message; if (request?.Quote is null || request?.State is null) { _logger?.LogError("Invalid request"); await context.RespondAsync(new BaseResult { ErrorMessage = "Invalid request", Success = false }); return; } try { var quotationMessage = BuildQuotationMessage(request.Quote, request.State); await _bus.Send(quotationMessage); _logger?.LogInformation("Message sent"); await context.RespondAsync(new BaseResult { Success = true }); } catch (Exception ex) { _logger?.LogError(ex, "An error occured while sending the message: {error}", ex.Message); await context.RespondAsync(new BaseResult { Success = false, ErrorMessage = ex.Message }); } } }
单元测试代码
[Fact] public async Task Consume_ValidRequest_SendsMessageToBus() { // Arrange var provider = new ServiceCollection() .AddMassTransitTestHarness(cfg => { cfg.AddConsumer<QuotationMessageConsumer>(); }) .BuildServiceProvider(true); var harness = provider.GetRequiredService<ITestHarness>(); await harness.Start(); try { var bus = provider.GetRequiredService<IBus>(); var client = bus.CreateRequestClient<SendQuotationMessage>(); // Act await client.GetResponse<BaseResult>(BuildEmptySendQuotationMessage()); // Assert Assert.True(await harness.Consumed.Any<SendQuotationMessage>()); var consumerHarness = provider.GetRequiredService<IConsumerTestHarness<QuotationMessageConsumer>>(); Assert.True(await consumerHarness.Consumed.Any<SendQuotationMessage>()); } finally { await harness.Stop(); await provider.DisposeAsync(); } }
遇到的问题
- 消费者调用
IBus.Send时抛出异常:'A convention for the message type QuotationMessage.QuotationMessage was not found',需要了解如何MockIBus; - 最初不知道如何测试
BaseResult,后续了解可通过GetResponse的结果进行断言。
后续优化配置
builder.Services.AddMassTransit(x => { x.UsingAzureServiceBus((context, cfg) => { cfg.ClearSerialization(); cfg.UseJsonSerializer(); var busConfig = context.GetRequiredService<IOptions<BusConfig>>().Value; cfg.Host(busConfig.ConnectionString); cfg.Message<Message>(x => { x.SetEntityName(busConfig.FaiMotorTopic); }); }); });
新增消息服务
public class MessageService : IMessageService { private readonly ILogger<MessageService> _logger; private readonly IPublishEndpoint _publishEndpoint; /// <summary> /// 构造函数 /// </summary> /// <param name="logger"></param> /// <param name="publishEndpoint"></param> public MessageService(ILogger<MessageService> logger, IPublishEndpoint publishEndpoint) { _logger = logger; _publishEndpoint = publishEndpoint; } /// <inheritdoc/> public async Task SendMessageAsync<T>(T message, CancellationToken cancellationToken = default) where T : class { _logger.LogInformation("Sending message..."); if (message == null) throw new ArgumentNullException(nameof(message)); try { await _publishEndpoint.Publish(message, cancellationToken); _logger.LogInformation("Message sent"); } catch (Exception ex) { _logger.LogError(ex, "An error occured while sending the message: {error}", ex.Message); throw; } } }
问题解决方案
1. 解决IBus调用异常及Mock方法
使用MassTransit测试工具时,无需手动Mock IBus,只需在测试Harness中配置消息类型约定,或改用IPublishEndpoint(如后续优化的MessageService)。若继续使用IBus.Send,需在测试配置中添加消息接收端点:
.AddMassTransitTestHarness(cfg => { cfg.AddConsumer<QuotationMessageConsumer>(); // 添加QuotationMessage的请求客户端,让测试Harness识别消息类型 cfg.AddRequestClient<QuotationMessage>(); cfg.UsingInMemory((context, cfg) => { cfg.ConfigureEndpoints(context); // 也可显式配置消息实体名称 cfg.Message<QuotationMessage>(x => x.SetEntityName("quotation-topic")); }); })
若要手动Mock IBus,可使用Moq:
var mockBus = new Mock<IBus>(); mockBus.Setup(b => b.Send(It.IsAny<QuotationMessage>(), default)).Returns(Task.CompletedTask); // 将mockBus.Object注入消费者实例 var consumer = new QuotationMessageConsumer(mockLogger.Object, mockBus.Object);
2. 测试BaseResult的方法
通过GetResponse获取返回结果后直接断言:
// Act var response = await client.GetResponse<BaseResult>(validRequest); // 断言成功场景 Assert.True(response.Message.Success); // 断言错误场景(如无效请求) var errorResponse = await client.GetResponse<BaseResult>(invalidRequest); Assert.False(errorResponse.Message.Success); Assert.Equal("Invalid request", errorResponse.Message.ErrorMessage);
内容的提问来源于stack exchange,提问作者user3075478
相关产品推荐
相关产品推荐

