.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的队列中未被处理。
具体修复步骤
- 启用完成传播:链接两个块时,添加
PropagateCompletion = true,这样当TransformBlock完成时,会自动向ActionBlock发送完成信号,避免ActionBlock一直等待新消息。 - 等待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
相关产品推荐
相关产品推荐

