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

在Azure函数中用NSubstitute模拟EventHubAsyncCollector遇异常问题

问题描述

我有一个基于CosmosDBTrigger的Azure函数,调用带partitionKey参数的AddAsync方法向EventHub发送数据。使用NSubstitute模拟IAsyncCollector<EventData>进行单元测试时,触发了EventHubWebJobsExtensions的异常——该扩展方法要求实例必须是EventHubAsyncCollector,而非NSubstitute代理对象。我已实现支持无参AddAsync的MockAsyncCollector,但无法解决带参AddAsync的测试问题,当前测试代码会抛出异常,寻求解决方案。

原函数代码

[FunctionName("TestFunction")]
public async Task Run([CosmosDBTrigger(
    databaseName: "Test",
    containerName: "test",
    CreateLeaseContainerIfNotExists = true,
    MaxItemsPerInvocation = 200,
    FeedPollDelay = 10000,
    Connection = "TestConnection",            
    LeaseContainerName = "testLeases")]IReadOnlyList<Document> inputs,
    [EventHub("test", Connection = "TestConnectionString")] IAsyncCollector<EventData> outputEvents,
    ILogger log)
{
    if (inputs != null && inputs.Count > 0)
    {
        try
        {
            await inputs.ForEachAsync(dop: 10, input =>
            {
                var id = input.GetPropertyValue<string>("id");
                var eventData = new EventData(Encoding.UTF8.GetBytes(input.ToString()));
                return outputEvents.AddAsync(eventData, id);
            });
        }
        catch (Exception exception)
        {
            _logger.LogError(exception, $"Error occurred while sending reports feed to datalake: {exception.Message}");
        }
    }
}

异常来源的扩展方法

public static Task AddAsync(this IAsyncCollector<EventData> instance, EventData eventData, string partitionKey, CancellationToken cancellationToken = default(CancellationToken)) =>
    instance switch
    {
        EventHubAsyncCollector ehCollector => ehCollector.AddAsync(eventData, partitionKey, cancellationToken),
        _ => throw new InvalidOperationException("Adding with a partition key is only available when using the Event Hubs extension package.")
    };

已实现的MockAsyncCollector

public class MockAsyncCollector<T> : IAsyncCollector<T>
{
    public readonly List<T> Items = new List<T>();

    public Task AddAsync(T item, CancellationToken cancellationToken = default)
    {
        Items.Add(item);
        return Task.FromResult(true);
    }

    public Task FlushAsync(CancellationToken cancellationToken = default)
    {
        return Task.FromResult(true);
    }
}

当前测试代码

public class Tests
{
    private readonly MockLogger<Test> _logger;
    private readonly Test _processor;

    public Tests()
    {
        _logger = Substitute.For<MockLogger<Test>>();
        _processor = new Test(_logger);
    }
    
    [Fact]
    public async Task Run_GivenProperMessage_ShouldNotThrowAndCallAddAsync()
    {
        // Arrange
        var input = new List<Document>();
        var inputDocument = new Document();
        inputDocument.SetPropertyValue("id", Guid.NewGuid().ToString());
        input.Add(inputDocument);

        var output = Substitute.For<IAsyncCollector<EventData>>();
        
        // Act
        var act = async () => await _processor.Run(input, output, _logger);

        // Assert
        await act.Should().NotThrowAsync();
        await output.Received(1).AddAsync(Arg.Any<EventData>(), Arg.Any<string>());
    }
}
解决方案

方案1:扩展MockAsyncCollector,适配带partitionKey的AddAsync

在测试项目中添加针对MockAsyncCollector<EventData>的扩展方法,让调用带partitionKey的AddAsync时绕过系统扩展的类型检查:

public static class MockAsyncCollectorExtensions
{
    // 可选:给MockAsyncCollector添加存储partitionKey的集合,用于验证
    public static List<string> PartitionKeys = new List<string>();

    public static Task AddAsync(this MockAsyncCollector<EventData> instance, EventData eventData, string partitionKey, CancellationToken cancellationToken = default)
    {
        PartitionKeys.Add(partitionKey);
        return instance.AddAsync(eventData, cancellationToken);
    }
}

修改测试代码,使用自定义Mock替代NSubstitute代理:

[Fact]
public async Task Run_GivenProperMessage_ShouldNotThrowAndCallAddAsync()
{
    // Arrange
    var input = new List<Document>();
    var testId = Guid.NewGuid().ToString();
    var inputDocument = new Document();
    inputDocument.SetPropertyValue("id", testId);
    input.Add(inputDocument);

    var output = new MockAsyncCollector<EventData>();
    
    // Act
    await _processor.Run(input, output, _logger);

    // Assert
    output.Items.Should().HaveCount(1);
    MockAsyncCollectorExtensions.PartitionKeys.Should().Contain(testId);
}

方案2:抽象发送逻辑,解耦依赖

定义封装EventHub发送逻辑的接口,让函数依赖该接口而非直接依赖IAsyncCollector<EventData>,降低测试复杂度:

// 定义发送接口
public interface IEventHubSender
{
    Task SendEventAsync(EventData eventData, string partitionKey);
}

// 生产环境实现
public class EventHubSender : IEventHubSender
{
    private readonly IAsyncCollector<EventData> _collector;

    public EventHubSender(IAsyncCollector<EventData> collector)
    {
        _collector = collector;
    }

    public async Task SendEventAsync(EventData eventData, string partitionKey)
    {
        await _collector.AddAsync(eventData, partitionKey);
    }
}

修改函数代码,注入IEventHubSender(建议通过构造函数注入,更符合依赖倒置原则):

[FunctionName("TestFunction")]
public async Task Run([CosmosDBTrigger(
    // 触发器参数保持不变
    )]IReadOnlyList<Document> inputs,
    [EventHub("test", Connection = "TestConnectionString")] IAsyncCollector<EventData> outputEvents,
    ILogger log)
{
    var sender = new EventHubSender(outputEvents);
    if (inputs != null && inputs.Count > 0)
    {
        try
        {
            await inputs.ForEachAsync(dop: 10, input =>
            {
                var id = input.GetPropertyValue<string>("id");
                var eventData = new EventData(Encoding.UTF8.GetBytes(input.ToString()));
                return sender.SendEventAsync(eventData, id);
            });
        }
        catch (Exception exception)
        {
            log.LogError(exception, $"Error occurred while sending reports feed to datalake: {exception.Message}");
        }
    }
}

测试时直接模拟IEventHubSender:

[Fact]
public async Task Run_GivenProperMessage_ShouldNotThrowAndCallSendEventAsync()
{
    // Arrange
    var input = new List<Document>();
    var testId = Guid.NewGuid().ToString();
    var inputDocument = new Document();
    inputDocument.SetPropertyValue("id", testId);
    input.Add(inputDocument);

    var mockSender = Substitute.For<IEventHubSender>();
    var processor = new Test(_logger, mockSender);
    
    // Act
    await processor.Run(input, null, _logger);

    // Assert
    await mockSender.Received(1).SendEventAsync(Arg.Any<EventData>(), testId);
}

方案3:模拟EventHubAsyncCollector(需处理内部类型访问)

若必须直接模拟IAsyncCollector<EventData>,可通过InternalsVisibleTo让测试项目访问内部类型EventHubAsyncCollector:

  1. 在函数项目的AssemblyInfo.cs中添加:
[assembly: InternalsVisibleTo("YourTestProjectName")]
  1. 测试中模拟EventHubAsyncCollector(需处理构造依赖):
[Fact]
public async Task Run_GivenProperMessage_ShouldNotThrowAndCallAddAsync()
{
    // Arrange
    var input = new List<Document>();
    var inputDocument = new Document();
    inputDocument.SetPropertyValue("id", Guid.NewGuid().ToString());
    input.Add(inputDocument);

    var mockCollector = Substitute.For<EventHubAsyncCollector>(
        Substitute.For<EventHubClient>(),
        Substitute.For<EventHubOptions>(),
        Substitute.For<ILogger<EventHubAsyncCollector>>()
    );
    mockCollector.AddAsync(Arg.Any<EventData>(), Arg.Any<string>(), Arg.Any<CancellationToken>()).Returns(Task.CompletedTask);
    
    // Act
    var act = async () => await _processor.Run(input, mockCollector, _logger);

    // Assert
    await act.Should().NotThrowAsync();
    await mockCollector.Received(1).AddAsync(Arg.Any<EventData>(), Arg.Any<string>(), Arg.Any<CancellationToken>());
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 17:36:08