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

ASP.NET Core Web API中RabbitMQ MessageProducer单元测试问题

RabbitMQ发布确认生产者单元测试问题解决

问题背景

我按照RabbitMQ官方文档实现了带发布确认的MessageProducer,使用XUnit和NSubstitute做单元测试时遇到两个问题:

  • 无法验证IModel的BasicAcks/BasicNacks回调函数是否正确触发;
  • 测试中_channel.NextPublishSeqNo始终不会递增。

服务方法代码

public async Task SendMessagesWithConfirmAsync<T>(IEnumerable<T> messages, string queueName, string routingKey)
{
    _channel.QueueDeclare(queueName, true, false);

    _channel.ConfirmSelect();

    // 注册确认回调
    _channel.BasicAcks += (sender, ea) => CleanOutstandingConfirms(ea.DeliveryTag, ea.Multiple);

    _channel.BasicNacks += (sender, ea) =>
        {
            _outstandingConfirms.TryGetValue(ea.DeliveryTag, out var body);

            Console.WriteLine(
                $"Message with body {body} has been nack-ed. Sequence number: {ea.DeliveryTag}, multiple: {ea.Multiple}"
            );

            CleanOutstandingConfirms(ea.DeliveryTag, ea.Multiple);
    };

    foreach (var message in messages)
    {
        var body = JsonSerializer.Serialize(message);
        _outstandingConfirms.TryAdd(_channel.NextPublishSeqNo, body);
        _channel.BasicPublish(queueName, routingKey, null, Encoding.UTF8.GetBytes(body));
    }

    await Task.CompletedTask;
}

单元测试代码

[Theory]
[InlineData("Test 1", "Test 2", "Test 3")]
public async void SendMessageWithConfirm_MultipleMessages_ShouldPublishMessagesAndWaitForConfirmOrDie(
    params string[] messages)
{
    // Arrange
    var messageProducer = new RabbitMqMessageProducer(_connectionFactory);

    // Act
    await messageProducer.SendMessagesWithConfirmAsync(messages, "invitations", "invitation");

    // Assert
    _channel.Received(messages.Length).BasicPublish(Arg.Any<string>(), Arg.Any<string>(),
        null, Arg.Any<ReadOnlyMemory<byte>>());

    // 无法验证回调事件的代码位置
    // ...
}

解决方案

1. 修复NextPublishSeqNo不递增的问题

NSubstitute的默认替身不会自动模拟NextPublishSeqNo的递增逻辑,需要手动设置属性的getter并拦截BasicPublish调用实现序号递增:

// Arrange阶段添加
var sequenceNumber = 0;
// 设置NextPublishSeqNo的返回值为当前计数器
_channel.NextPublishSeqNo.Returns(x => sequenceNumber);
// 每次调用BasicPublish后,计数器自增
_channel.When(x => x.BasicPublish(Arg.Any<string>(), Arg.Any<string>(), Arg.Any<bool>(), Arg.Any<IBasicProperties>(), Arg.Any<ReadOnlyMemory<byte>>()))
        .Do(_ => sequenceNumber++);

2. 验证回调函数的触发逻辑

要测试回调,需要手动触发BasicAcks/BasicNacks事件,同时需要让_outstandingConfirms的状态可被测试访问(比如将其改为内部属性,通过[InternalsVisibleTo("YourTestProjectName")]让测试项目可见,或者提供公开方法获取其状态)。

修改后的测试示例:

[Theory]
[InlineData("Test 1", "Test 2", "Test 3")]
public async void SendMessageWithConfirm_MultipleMessages_ShouldHandleAcks()
{
    // Arrange
    var sequenceNumber = 0;
    _channel.NextPublishSeqNo.Returns(x => sequenceNumber);
    _channel.When(x => x.BasicPublish(Arg.Any<string>(), Arg.Any<string>(), null, Arg.Any<ReadOnlyMemory<byte>>()))
            .Do(_ => sequenceNumber++);

    var messageProducer = new RabbitMqMessageProducer(_connectionFactory);

    // Act
    await messageProducer.SendMessagesWithConfirmAsync(messages, "invitations", "invitation");
    // 手动触发BasicAcks事件,模拟RabbitMQ的确认响应
    _channel.BasicAcks += Raise.Event(null, new BasicAckEventArgs { DeliveryTag = 3, Multiple = true });

    // Assert
    _channel.Received(messages.Length).BasicPublish(Arg.Any<string>(), Arg.Any<string>(),
        null, Arg.Any<ReadOnlyMemory<byte>>());
    // 验证未确认消息集合被清空
    Assert.Empty(messageProducer.OutstandingConfirms);
}

额外优化建议

  • 不要在每次调用SendMessagesWithConfirmAsync时重复注册BasicAcks/BasicNacks回调,会导致同一事件触发多次处理逻辑。建议在Channel初始化阶段(比如生产者构造函数内)只注册一次回调。
  • 方法中没有实际异步操作,可移除await Task.CompletedTask,直接返回Task.CompletedTask,或者将方法改为同步方法。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 21:39:36