为何我的ActionBlock未手动设置就进入Completed状态?
问题描述
运行约15分钟后,ParallelWorker.SendAsync(rd)返回false,且ParallelWorker.Complete为true,触发异常。已确认ProcessSetupsAsync内部同步代码未抛出异常(catch块从未执行)。
场景:通过静态ActionBlock处理大量生成的工作项,更换数据集可复现该问题。
原始代码
public static class Setups { private struct RunData { internal MyClass1 Setup; internal MyClass2 Positions; internal string Set; } private static readonly ActionBlock<RunData> ParallelWorker = new(d => ProcessSetupsAsync(d.Setup, d.Positions, d.Set), new ExecutionDataflowBlockOptions { BoundedCapacity = Environment.ProcessorCount * 10, MaxDegreeOfParallelism = Environment.ProcessorCount, SingleProducerConstrained = false }); public static async Task GetSetups(ItemType[] bar1Filter, ItemType[] bar2Filter, bool sameSet) { for (/* do some work*/) { foreach (RunData rd in MyMethod1(/*variables*/).Select(final => new RunData { Positions = positions, Set = set, Setup = final })) { if (!await ParallelWorker.SendAsync(rd).ConfigureAwait(false)) { // Fails after running for about 15 minutes // with ParallelWorker.Complete is true throw new Exception("xxxxxxxx"); } } } ParallelWorker.Complete(); await ParallelWorker.Completion.ConfigureAwait(false); } private static Task ProcessSetupsAsync (MyClass1 setup, MyClass2 positions, string set) { try { List<MyClass1> setups = new(); /* Do some work */ return setups.Count > 0 ? LocalAccess.AddSetupAsync(setups) : Task.CompletedTask; } catch (Exception ex) { Console.WriteLine(ex); //<-Never gets hit. throw; } } }
问题分析与解决方案
核心原因
- 异步任务异常未被捕获:
ProcessSetupsAsync的try/catch仅覆盖同步代码逻辑,LocalAccess.AddSetupAsync返回的异步任务若抛出异常,属于异步异常,不会被当前捕获逻辑处理。此时ActionBlock会因委托任务故障进入完成状态,拒绝接收新消息。 - 静态实例被意外终止:
ParallelWorker是静态字段,若应用内其他代码路径调用了ParallelWorker.Complete(),会导致块提前进入完成状态。 - 故障状态未被感知:
ActionBlock进入故障状态后,Completion任务会包含异常,但原代码未捕获这部分异常,无法定位具体错误根源。
解决步骤
- 修复异步异常捕获逻辑:将
ProcessSetupsAsync改为async Task,通过await捕获异步任务的异常:private static async Task ProcessSetupsAsync(MyClass1 setup, MyClass2 positions, string set) { try { List<MyClass1> setups = new(); /* Do some work */ if (setups.Count > 0) { await LocalAccess.AddSetupAsync(setups).ConfigureAwait(false); } } catch (Exception ex) { Console.WriteLine(ex); throw; } } - 避免静态实例的状态污染:将
ParallelWorker改为非静态,每次调用GetSetups时创建新实例,防止跨调用的状态干扰:public static async Task GetSetups(ItemType[] bar1Filter, ItemType[] bar2Filter, bool sameSet) { var parallelWorker = new ActionBlock<RunData>(d => ProcessSetupsAsync(d.Setup, d.Positions, d.Set), new ExecutionDataflowBlockOptions { BoundedCapacity = Environment.ProcessorCount * 10, MaxDegreeOfParallelism = Environment.ProcessorCount, SingleProducerConstrained = false }); for (/* do some work*/) { foreach (RunData rd in MyMethod1(/*variables*/).Select(final => new RunData { Positions = positions, Set = set, Setup = final })) { if (!await parallelWorker.SendAsync(rd).ConfigureAwait(false)) { throw new Exception("xxxxxxxx"); } } } parallelWorker.Complete(); await parallelWorker.Completion.ConfigureAwait(false); } - 捕获
ActionBlock的故障异常:在等待Completion时捕获AggregateException,定位具体错误:try { await parallelWorker.Completion.ConfigureAwait(false); } catch (AggregateException ex) { foreach (var innerEx in ex.InnerExceptions) { Console.WriteLine($"ActionBlock故障原因:{innerEx.Message}\n{innerEx.StackTrace}"); } throw; }
内容的提问来源于stack exchange,提问作者Iaman Swtrse
相关产品推荐
相关产品推荐

