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

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',需要了解如何Mock IBus;
  • 最初不知道如何测试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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 13:43:28