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

如何用Mock测试C#中RabbitMQClient的Received事件处理器?

测试DataOutConsumer中consumer.Received事件处理器的方案

我需要为DataOutConsumer类Start方法内的consumer.Received事件处理器编写Mock测试,但不清楚如何用Mock数据测试该代码块,也不知道如何触发该事件。我考虑过把这段代码提取到其他方法或类中,但不知道怎么关联。相关代码如下:

public class DataOutConsumer : IDataOutConsumer
{
private readonly IRabbitMqConnectionDetails _mqConnectionDetails;
private readonly IRabbitMqConnectionFactory _rabbitMqConnectionFactory;

private readonly IDictionary<string, IMessageFormatter> _formatters;
private readonly ILogger _logger;
private readonly IMessageRepository _repository;

private IConnection _connection;
private IModel _channel;

public DataOutConsumer(
    IRabbitMqConnectionFactory rabbitMqConnectionFactory,
    IRabbitMqConnectionDetails mqConnectionDetails,
    IMessageRepository repository,
    IEnumerable<IMessageFormatter> formatters,
    ILogger logger)
{
    _rabbitMqConnectionFactory = rabbitMqConnectionFactory;
    _mqConnectionDetails = mqConnectionDetails;
    _repository = repository;
    _formatters = formatters.Select(f => new { f.Source, f }).ToDictionary(fs => fs.Source, fs => fs.f);
    _logger = logger;
}

/// <summary>
/// Start the consumer that will process messages on the DataOut queue
/// </summary>
public void Start()
{
    _connection = _rabbitMqConnectionFactory.CreateConnection();

    var channel = _connection.CreateModel(); 
    channel.BasicQos(0, 1, false);

    var consumer = new EventingBasicConsumer(channel);
    consumer.Received += (model, ea) =>
    {
        try
        {
            var properties = ea.BasicProperties;

            var headers = properties.Headers;

            var endpoint = headers.ReadHeaderValue("NServiceBus.OriginatingEndpoint");
            var formatter = _formatters[endpoint];
            var formattedMessage = formatter.Format(Encoding.UTF8.GetString(ea.Body.ToArray()));

            _logger.Information($"PreAddMessage: Originating Endpoint: {endpoint}, Enclosed Message Types: {headers.ReadHeaderValue("NServiceBus.EnclosedMessageTypes")}, Message Endpoint: DataOut, Payload: {formattedMessage.Payload}");

            var task = _repository.AddMessage(
                formattedMessage.Id,
                "DataOut",
                endpoint,
                headers.ReadHeaderValue("NServiceBus.EnclosedMessageTypes"),
                formattedMessage.Payload,
                formattedMessage.DateQueued);

            task.GetAwaiter().GetResult();

            // Acknowledge that the message has been processed
            channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
        }
        catch (Exception e)
        {
            // Mark as Nack'ed
            // These seem to be sent to the error queue at least, also logs as an error
            consumer.Model.BasicNack(ea.DeliveryTag, false, false);
            _logger.Error(e.Message, e);
        }
    };
    channel.BasicConsume(queue: _mqConnectionDetails.Endpoint,
        autoAck: false,
        consumer: consumer);
}

/// <summary>
/// delete the channel and connection, don't believe this ever needs calling
/// </summary>
public void Stop()
{
    try
    {
        _channel.Close();
        _connection.Close();
    }
    catch (Exception e)
    {
        _logger.Error(e.Message, e);
    }
}
}

一、优先方案:重构代码提取事件逻辑

把Received事件内的匿名委托逻辑提取为独立的实例方法,这是最简洁易测的方式,避免直接操作事件触发的复杂操作。

重构后的DataOutConsumer代码

public class DataOutConsumer : IDataOutConsumer
{
    // 原有字段和构造函数不变...

    // 新增独立方法承载事件逻辑
    public void ProcessReceivedMessage(object model, BasicDeliverEventArgs ea, IModel channel)
    {
        try
        {
            var properties = ea.BasicProperties;
            var headers = properties.Headers;

            var endpoint = headers.ReadHeaderValue("NServiceBus.OriginatingEndpoint");
            var formatter = _formatters[endpoint];
            var formattedMessage = formatter.Format(Encoding.UTF8.GetString(ea.Body.ToArray()));

            _logger.Information($"PreAddMessage: Originating Endpoint: {endpoint}, Enclosed Message Types: {headers.ReadHeaderValue("NServiceBus.EnclosedMessageTypes")}, Message Endpoint: DataOut, Payload: {formattedMessage.Payload}");

            var task = _repository.AddMessage(
                formattedMessage.Id,
                "DataOut",
                endpoint,
                headers.ReadHeaderValue("NServiceBus.EnclosedMessageTypes"),
                formattedMessage.Payload,
                formattedMessage.DateQueued);

            task.GetAwaiter().GetResult();

            channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
        }
        catch (Exception e)
        {
            channel.BasicNack(ea.DeliveryTag, false, false);
            _logger.Error(e.Message, e);
        }
    }

    public void Start()
    {
        _connection = _rabbitMqConnectionFactory.CreateConnection();
        var channel = _connection.CreateModel(); 
        channel.BasicQos(0, 1, false);

        var consumer = new EventingBasicConsumer(channel);
        // 绑定事件到新增的方法
        consumer.Received += (model, ea) => ProcessReceivedMessage(model, ea, channel);
        
        channel.BasicConsume(queue: _mqConnectionDetails.Endpoint,
            autoAck: false,
            consumer: consumer);
    }

    // Stop方法不变...
}

二、编写Mock测试(基于xUnit + Moq)

下面是覆盖正常流程和异常流程的测试示例:

1. 正常流程测试

验证消息处理成功时,正确调用BasicAck和AddMessage:

using Xunit;
using Moq;
using RabbitMQ.Client;
using System.Text;
using System.Collections.Generic;

public class DataOutConsumerTests
{
    [Fact]
    public void ProcessReceivedMessage_ValidMessage_CallsAckAndAddMessage()
    {
        // 1. 准备Mock依赖
        var mockConnectionFactory = new Mock<IRabbitMqConnectionFactory>();
        var mockConnectionDetails = new Mock<IRabbitMqConnectionDetails>();
        var mockRepository = new Mock<IMessageRepository>();
        var mockLogger = new Mock<ILogger>();

        // Mock消息格式化器
        var mockFormatter = new Mock<IMessageFormatter>();
        var testFormattedMsg = new FormattedMessage
        {
            Id = "test-msg-id",
            Payload = "test-payload",
            DateQueued = DateTime.UtcNow
        };
        mockFormatter.Setup(f => f.Format(It.IsAny<string>())).Returns(testFormattedMsg);
        mockFormatter.Setup(f => f.Source).Returns("TestOriginEndpoint");

        // 2. 创建被测实例
        var consumer = new DataOutConsumer(
            mockConnectionFactory.Object,
            mockConnectionDetails.Object,
            mockRepository.Object,
            new List<IMessageFormatter> { mockFormatter.Object },
            mockLogger.Object);

        // 3. 构造测试用的消息参数
        var mockChannel = new Mock<IModel>();
        var mockProperties = new Mock<IBasicProperties>();
        var headers = new Dictionary<string, object>
        {
            { "NServiceBus.OriginatingEndpoint", "TestOriginEndpoint" },
            { "NServiceBus.EnclosedMessageTypes", "TestMessageType" }
        };
        mockProperties.Setup(p => p.Headers).Returns(headers);

        var ea = new BasicDeliverEventArgs
        {
            BasicProperties = mockProperties.Object,
            Body = Encoding.UTF8.GetBytes("raw-test-body"),
            DeliveryTag = 12345
        };

        // 4. 调用被测方法
        consumer.ProcessReceivedMessage(null, ea, mockChannel.Object);

        // 5. 验证行为
        mockRepository.Verify(r => r.AddMessage(
            testFormattedMsg.Id,
            "DataOut",
            "TestOriginEndpoint",
            "TestMessageType",
            testFormattedMsg.Payload,
            testFormattedMsg.DateQueued), Times.Once);

        mockChannel.Verify(c => c.BasicAck(12345, false), Times.Once);
    }
}

2. 异常流程测试

验证处理抛出异常时,正确调用BasicNack并记录错误日志:

[Fact]
public void ProcessReceivedMessage_ThrowsException_CallsNackAndLogsError()
{
    // 1. 准备Mock依赖
    var mockConnectionFactory = new Mock<IRabbitMqConnectionFactory>();
    var mockConnectionDetails = new Mock<IRabbitMqConnectionDetails>();
    var mockRepository = new Mock<IMessageRepository>();
    var mockLogger = new Mock<ILogger>();

    var mockFormatter = new Mock<IMessageFormatter>();
    mockFormatter.Setup(f => f.Source).Returns("TestOriginEndpoint");
    // 让格式化方法抛出异常
    mockFormatter.Setup(f => f.Format(It.IsAny<string>())).Throws(new InvalidOperationException("Mock processing error"));

    // 2. 创建被测实例
    var consumer = new DataOutConsumer(
        mockConnectionFactory.Object,
        mockConnectionDetails.Object,
        mockRepository.Object,
        new List<IMessageFormatter> { mockFormatter.Object },
        mockLogger.Object);

    // 3. 构造测试参数
    var mockChannel = new Mock<IModel>();
    var mockProperties = new Mock<IBasicProperties>();
    var headers = new Dictionary<string, object>
    {
        { "NServiceBus.OriginatingEndpoint", "TestOriginEndpoint" }
    };
    mockProperties.Setup(p => p.Headers).Returns(headers);

    var ea = new BasicDeliverEventArgs
    {
        BasicProperties = mockProperties.Object,
        Body = Encoding.UTF8.GetBytes("raw-test-body"),
        DeliveryTag = 12345
    };

    // 4. 调用被测方法
    consumer.ProcessReceivedMessage(null, ea, mockChannel.Object);

    // 5. 验证行为
    mockChannel.Verify(c => c.BasicNack(12345, false, false), Times.Once);
    mockLogger.Verify(l => l.Error("Mock processing error", It.IsAny<InvalidOperationException>()), Times.Once);
}

三、不重构的备选方案:直接触发事件

如果暂时不想修改业务代码,可通过反射获取Start方法中创建的EventingBasicConsumer实例,手动触发Received事件。但这种方式耦合度高,维护成本大,仅作参考:

[Fact]
public void Start_ReceivedEventTriggered_ProcessesMessage()
{
    // 准备Mock连接和通道
    var mockConnection = new Mock<IConnection>();
    var mockChannel = new Mock<IModel>();
    var mockConnectionFactory = new Mock<IRabbitMqConnectionFactory>();
    mockConnectionFactory.Setup(f => f.CreateConnection()).Returns(mockConnection.Object);
    mockConnection.Setup(c => c.CreateModel()).Returns(mockChannel.Object);

    // 其他Mock依赖准备...
    var mockConnectionDetails = new Mock<IRabbitMqConnectionDetails>();
    var mockRepository = new Mock<IMessageRepository>();
    var mockLogger = new Mock<ILogger>();
    var mockFormatter = new Mock<IMessageFormatter>();
    mockFormatter.Setup(f => f.Source).Returns("TestOriginEndpoint");
    var testFormattedMsg = new FormattedMessage
    {
        Id = "test-id",
        Payload = "test-payload",
        DateQueued = DateTime.UtcNow
    };
    mockFormatter.Setup(f => f.Format(It.IsAny<string>())).Returns(testFormattedMsg);

    // 创建实例并启动
    var consumer = new DataOutConsumer(
        mockConnectionFactory.Object,
        mockConnectionDetails.Object,
        mockRepository.Object,
        new List<IMessageFormatter> { mockFormatter.Object },
        mockLogger.Object);
    consumer.Start();

    // 通过反射获取局部变量consumer(需要借助第三方库如Mono.Cecil,或修改代码将consumer改为私有字段)
    // 此步骤复杂度高,不推荐,仅作思路展示
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 12:18:26