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

.NET Core Dataflow单元测试未等待ActionBlock执行完成求助

问题:Dataflow在单元测试中未等待ActionBlock完成

我采用Dataflow模式(包含TransformBlock与ActionBlock)实现了一个.NET Core Worker Service,Worker运行时两个块均可正常执行并返回结果,但执行XUnit单元测试时,仅TransformBlock执行,测试未等待ActionBlock完成就返回结果。尝试直接await ActionBlock,但代码无法继续执行。


Service代码

public class ServiceDataflow(IConfiguration configuration,
          ILogger<WorkerServiceDataflow> logger,
          IEService eService,
          IOptions<ENSettings> enSettings) : IServiceDataflow
{
    private enSettings _enSettings { get; } = enSettings.Value;

    public async Task<DataTable> StartProcessing(IEnumerable<ENQueue> items)
    {
        DataTable Log = new DataTable();
        Log.Columns.Add("Colum1", typeof(string));
        Log.Columns.Add("Colum2", typeof(Decimal));

        // Define dataflow blocks
        var processingBlock = new TransformBlock<ENQueue, IResult>(async email =>
        {
            // Perform asynchronous processing on the item
            try
            {
                var udTemplate = ReplacePlaceholder(email);
                string name = string.Format("{0} {1}", email.Last_Name, email.First_Name);
                await emailService.SendEmail(email.Address, email.Subject, udTemplate, name, _enSettings.IsSandBox);
                return new Result { ProcessedData = email, Status = "Success", IsProcessed = true };
            }
            catch (Exception ex)
            {
                var exceptionMsg = string.Format("ServiceParallel: Error in TransformBlock: {0} - {1} - {2}", DateTime.Now, email.Address, ex.Message);
                logger.LogError(ex, exceptionMsg);
                return new Result { ProcessedData = email, Status = "Error" , IsProcessed = false};
            }
        });

        var consumerBlock = new ActionBlock<IResult>(result =>
        {
            // Process the result (e.g., store in database, log)
            if (result.IsProcessed)
            {
                // Access the result if the task returned a value
                DataRow newRow = Log.NewRow();
                newRow["Column1"] = result.ProcessedData?.Address;
                newRow["Column2"] = result.ProcessedData?.TemplateId;
                 
                Log.Rows.Add(newRow);
                logger.LogInformation("Email Send at: {0} Recipient: {1} status{2} ", DateTimeOffset.Now, result.ProcessedData?.Address, result.Status);                   
            }
            else
            {
                logger.LogInformation("Error processing email at: {0} Recipient: {1}", DateTimeOffset.Now, result.ProcessedData?.Address);
            }
        });

        // Link the blocks
        processingBlock.LinkTo(consumerBlock);

        // Post items for processing
        foreach (var item in items)
        {
            await processingBlock.SendAsync(item);
        }

        // Signal completion of adding items
        processingBlock.Complete();

        // Wait for processing to finish using Task.WaitAll
        var completionTask = Task.WhenAll(processingBlock.Completion);

        await completionTask;
        return Log;
    }
}

单元测试代码

public class ServiceDataflowTest : ServiceDataflowFixture
{
    [Fact]
    public async Task StartProcessing_ProcessesItemsSuccessfully_ReturnsDataTable()
    {
        try
        {
            // Mock dependencies
            _mockOptions.Setup(m => m.Value).Returns(new EmailNotificationSettings { IsSandBox = true, 
                UnSubscribeURL= "https://serverurl/Unsubscribe" ,
                ProcessUser = "Test1",
                url = "https://serverurl"
            });

            // Setup for LogInformation with any message and arguments
            _mockLogger.Setup(m => m.Log(
              It.IsAny<LogLevel>(),
              It.IsAny<EventId>(),
              It.IsAny<object>(), // Capture any object passed as state
              It.IsAny<Exception>(),
              It.IsAny<Func<object, Exception, string>>() // Don't care about formatter
            )).Verifiable();

            _mockEService.Setup(m => m.SendEmail(It.IsAny<string>(), It.IsAny<string>(), It.IsAny<string>(), It.IsAny<string>(), It.IsAny<bool>()))
.Returns(Task.FromResult(new SendGrid.Response(HttpStatusCode.OK, null, null))); // Assuming this constructor exists

            // Sample data
            var sampleItem = new ENQueue
            {
                Email_Address = "xyx@example.com",
                Subject = "qwqwqw",
            };
            var items = new List<ENQueue>() { sampleItem };

            // Create the Dataflow instance
            var dataflow = new ServiceDataflow(_mockConfiguration.Object, _mockLogger.Object,
                                        _mockEService.Object, _mockOptions.Object);

            // Call the method under test
            var dataTable = await dataflow.StartProcessing(items);

            // Assertions
            Assert.NotNull(dataTable);
            Assert.Single(dataTable.Rows); // Expect one row for the processed item
            /*                Assert.Equal("test@example.com", dataTable.Rows[0]["Column1"]);
                            Assert.Equal(123, dataTable.Rows[0]["Column2"]);
                            Assert.Equal(123, dataTable.Rows[0]["Column3"]);*/

            // Verify eService.Send was called with expected arguments
            _mockEService.Verify(m => m.SendEmail(It.IsAny<string>(), It.IsAny<string>(), It.IsAny<string>(), It.IsAny<string>(), It.IsAny<bool>()), Times.Once);
            //_mockLogger.Verify(m => m.LogInformation(It.IsAny<string>(), It.IsAny<string>(), It.IsAny<string>()), Times.AtLeastOnce);

            _mockLogger.Verify();
        }
        catch (Exception ex) { }
    }
}

解决方案

核心问题分析

测试提前返回是因为只等待了TransformBlock的完成,没有等待ActionBlock处理完所有输出消息。TransformBlock的Completion仅表示它不再接受新消息、且自身已处理完所有输入,但它产出的消息可能还在ActionBlock的队列中未被处理。

具体修复步骤

  1. 启用完成传播:链接两个块时,添加PropagateCompletion = true,这样当TransformBlock完成时,会自动向ActionBlock发送完成信号,避免ActionBlock一直等待新消息。
  2. 等待ActionBlock完成:在TransformBlock完成后,必须等待ActionBlock的Completion任务,确保所有消息都被处理完毕。

修改后的关键代码片段:

// 链接块时启用完成传播
processingBlock.LinkTo(consumerBlock, new DataflowLinkOptions { PropagateCompletion = true });

// ... 其他代码 ...

// Signal completion of adding items
processingBlock.Complete();

// 先等TransformBlock处理完所有输入,再等ActionBlock处理完所有输出
await processingBlock.Completion;
await consumerBlock.Completion;

return Log;

额外优化建议

  • 移除单元测试中的空catch块,它会吞掉测试过程中的错误,导致无法排查问题。
  • 若之前单独await ActionBlock.Completion时代码卡住,就是因为没有给ActionBlock发送完成信号,启用PropagateCompletion后即可解决。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 14:05:56