尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

深入理解 .NET 任务并行库 ContinueWhenAll:多任务合流与延续机制

深入理解 .NET 任务并行库 ContinueWhenAll:多任务合流与延续机制 搞异步编程这么多年我一直在跟 .NET 的任务并行库TPL打交道。最初接触 TPL 的时候处理并行任务之间的先后顺序最让我头疼尤其是“一批任务全部跑完再做下一件事”这种常见的场景。当时我用得最多的就是Task.Factory.ContinueWhenAll。这个 API 虽然看起来不起眼但它代表了 TPL 延续模型里很核心的一环把一组任务的结果汇合点变成一个全新的任务。很多老项目里大量使用了这种写法你如果看不懂它读代码会非常吃力。这篇文章我就把这个 API 从头到尾拆开讲清楚包括方法签名、调度细节、实际案例还有我踩过的一些坑。我先把结论放在前面如果是新写的 async/await 代码绝大多数情况下用Task.WhenAll会更好代码也更简洁但如果你在做延续式任务链、处理回调式 TPL 管道或者需要维护旧代码Task.Factory.ContinueWhenAll仍然是绕不开的基础。理解它不是为了让你把它当成首选而是为了让你在遇到这类代码时能准确判断它的行为不至于被“延续不执行”“异常没抛出来”之类的问题卡住。1. 从任务到延续ContinueWhenAll 到底解决了什么问题1.1 TPL 的“任务接力”模型TPL 的核心抽象是Task表示一个未来会完成的操作。这个抽象本身不难理解难的是把多个任务串起来。你可能有三个下载任务、两个计算任务它们各自独立启动但你希望它们在全部结束之后再进入下一步。最朴素的做法是逐个Wait()然后按顺序处理结果。这当然可行但问题是当你调用Wait()的时候当前线程就阻塞了。在异步或并行场景里阻塞意味着浪费线程资源也可能造成线程饥饿。TPL 给出的方案是“延续”。所谓延续就是给一个任务挂上后续动作前一个任务完成时后续动作会自动启动。ContinueWith是单任务延续的入口而ContinueWhenAll是“多任务合流”的入口。简单说Task.Factory.ContinueWhenAll(tasks, action)做的事情是等待tasks数组里的所有任务都处于完成状态然后启动新的延续任务在延续任务里执行action。这个设计思路有点像接力赛你把接力棒分别交给几个选手等所有人都到达终点下一棒的选手才开始跑。各位选手到达的先后顺序不重要重要的是“全部到达”这个聚合条件。1.2 ContinueWhenAll、ContinueWith、ContinueWhenAny 的定位区别很多初学者把ContinueWith和ContinueWhenAll搞混。它们在代码形式上非常像但语义差别很大。ContinueWith针对单个任务一个前序任务对应一个后续任务。它处理的是“一件事情做完之后做什么”。ContinueWhenAll针对多个任务所有前序任务都完成后才执行后续任务。它处理的是“多件事都做完之后做什么”。ContinueWhenAny同样针对多个任务但只要有任意一个前序任务完成后续任务就会启动。它处理的是“最先完成的那件事触发什么”。这三个方法共同构成了 TPL 延续机制的基础。ContinueWhenAll最适合“汇合”场景也就是并行分支最后要合并结果的场景。举个具体的例子你要从三个不同的服务拉取配置三个请求可以同时发出去但最终配置文件要合并后写入磁盘。这种情况ContinueWhenAll就是很自然的表达方式。1.3 有了 Task.WhenAll为什么还要了解 ContinueWhenAll说实话现在写 async/await 代码我几乎不会主动去用ContinueWhenAll。Task.WhenAll可谓它的“异步替代版”。当你写await Task.WhenAll(tasks)时编译器会帮你处理后续逻辑代码就像同步代码一样自然异常也能通过 try/catch 捕获可读性好太多了。但这不代表ContinueWhenAll没有价值。第一老代码里大量存在这种延续式写法尤其是一些历史遗留的批处理框架、工作流引擎、自定义任务调度组件。第二在某些高级场景里你可能需要把延续动作本身作为独立任务去调度、取消、监控。这时候ContinueWhenAll返回的Task对象就比 async/await 的隐式状态机更容易在运行时被“钩住”。第三ContinueWhenAll和ContinueWith系列在 TPL 内部的实现机制是相通的理解它有助于理解整个延续模型。所以我的态度很明确新代码优先 async/await但读懂ContinueWhenAll是基本功。接下来我把它的 API 细节和底层行为拆开看。2. 方法签名与调度细节先把 API 看透再写代码2.1 全部重载与参数含义Task.Factory.ContinueWhenAll不是只有一个简单版本它有一组重载。归纳起来主要分两类非泛型版本和泛型版本。非泛型版本的核心形式public Task ContinueWhenAll( Task[] tasks, ActionTask[] continuationAction )它的含义是等待tasks全部完成然后执行continuationAction。传给continuationAction的参数就是先前那一整个tasks数组。这看起来有点奇怪为什么回调参数不是单个Task而是整个数组因为延续的任务本身就是“汇总节点”你需要检查每个前序任务的结果状态自然要把数组传进来。泛型版本多针对TaskTResultpublic TaskTResult ContinueWhenAllTAntecedentResult, TResult( TaskTAntecedentResult[] tasks, FuncTaskTAntecedentResult[], TResult continuationFunction )这个版本返回一个泛型任务表示延续动作的最终结果。还有带CancellationToken、TaskContinuationOptions、TaskScheduler的完整重载public TaskTResult ContinueWhenAllTAntecedentResult, TResult( TaskTAntecedentResult[] tasks, FuncTaskTAntecedentResult[], TResult continuationFunction, CancellationToken cancellationToken, TaskContinuationOptions continuationOptions, TaskScheduler scheduler )参数多了之后行为就会变得复杂。cancellationToken不是用来取消前序任务的而是用来取消延续任务本身的。continuationOptions决定延续在什么条件下启动。scheduler决定延续任务被安排到哪个调度器上执行。这些参数组合起来是很多诡异问题的高发区。参数作用常见失误Task[] tasks待等待的前序任务集合传入空数组或混入未启动的任务continuationAction全部完成后的回调在回调里直接访问.Result触发异常cancellationToken取消延续任务本身误以为它能取消前序任务continuationOptions设置延续触发条件使用互斥选项导致运行时异常scheduler指定延续的调度器未指定时默认调度行为不直观2.2 TaskScheduler 和 TaskContinuationOptions 的行为细节先说TaskContinuationOptions。最常见的值是None也就是默认行为。这时候只要前序任务全部完成无论它们是成功、失败还是被取消延续任务都会执行。你可能会问既然如此那前序任务抛出的异常去哪了这里是新手最容易迷糊的地方TPL 延续默认不会自动传播前序任务的异常。如果你想在延续里拿到异常要么先通过Task.WaitAll(tasks)去触发AggregateException要么遍历tasks逐个检查Status和Exception。还要注意TaskContinuationOptions.OnlyOnRanToCompletion和NotOnRanToCompletion这类选项是基于“前序任务最终处于什么状态”来决定是否启动延续的。当你在代码里同时指定了NotOnCanceled和OnlyOnCanceled就会抛出ArgumentOutOfRangeException因为这两个条件互斥。这是很典型的运行时错误。再说TaskScheduler。TPL 的延续默认会使用当时的TaskScheduler.Current如果你在线程池线程上调用它一般就是线程池调度器。但如果你在某个自定义调度器上下文里调用延续就可能跳到那个上下文执行。这个行为对新手来说非常反直觉因为很多人以为延续一定会在线程池线程上运行。想要稳定预期最简单的办法就是显式传入TaskScheduler.Default。这样延续任务基本确定在线程池线程上执行避免被当前上下文“劫持”。如果你希望延续在特定同步上下文里跑可以传TaskScheduler.FromCurrentSynchronizationContext()不过对于典型的服务器应用我宁愿直接用显式的调度器或继续用 async/await。2.3 内部机制简析信号量式的完成计数理解ContinueWhenAll的底层行为我习惯把它想象成一个“倒计数信号量”。TPL 内部会为这组前序任务创建一个包装的延续代理每个前序任务完成时内部计数就减一。当计数归零也就是所有任务都完成时才触发真正的延续回调。为什么这样设计因为每个任务完成的时间不固定直接轮询所有任务状态太浪费资源倒计数方式只需每个任务完成时做一次判断。这个机制也解释了为什么ExecuteSynchronously选项存在。默认情况下最后一个前序任务完成时TPL 会把延续任务排队到目标调度器上执行这个排队过程有开销。加了ExecuteSynchronously后延续可能会直接在上一个完成任务的线程上同步执行。优点是省了调度开销缺点是执行线程不确定而且如果延续里写了耗时操作你会意想不到地拖慢最后一个前序任务所在线程。我的经验是只有延续动作极短时才用ExecuteSynchronously否则宁可多花一次调度开销也不要让延续阻塞未知线程。3. 实操三类最典型的 ContinueWhenAll 场景3.1 场景一聚合多个并行加载任务的结果假设你有三个数据源需要并行读取然后把读取结果合并成一个文件。这是典型的“并行读取、汇合写入”场景。先看一版基于ContinueWhenAll的写法var dataSources new[] { https://example.com/data/1.txt, https://example.com/data/2.txt, https://example.com/data/3.txt }; var loadTasks dataSources .Select(url Task.Run(() LoadData(url))) .ToArray(); var mergeTask Task.Factory.ContinueWhenAll( loadTasks, completedTasks { var buffer new StringBuilder(); foreach (var task in completedTasks) { // 不能直接 task.Result因为可能有失败任务 if (task.Status TaskStatus.RanToCompletion) { buffer.AppendLine(task.Result); } else { buffer.AppendLine($[{task.Status}]); } } File.WriteAllText(merged.txt, buffer.ToString()); }, CancellationToken.None, TaskContinuationOptions.OnlyOnRanToCompletion, TaskScheduler.Default ); // 等合并任务完成后检查结果 mergeTask.Wait();这段代码里我特意加上了OnlyOnRanToCompletion表示只有所有加载任务都成功完成时才执行合并。因为在“合并三个文件”的业务中任何一个数据源失败最终合并文件就是残缺的不如不合并。如果选择None那合并逻辑就必须自己判断每个任务的状态否则很容易在task.Result上崩掉。这里要注意一个细节Task.Run(() LoadData(url))返回的是Taskstring数组ContinueWhenAll会推断泛型参数所以你可以在回调里拿到带泛型结果的任务数组去访问.Result。当你只是把任务数组声明为Task[]结果就会被擦除成非泛型Task访问结果就比较麻烦了。3.2 场景二并行计算后合并输出再看一个 CPU 密集型场景把一个大数组拆成三段每段分别求和最后合并总和。用ContinueWhenAll可以写成这样var segments new[] { new[] { 1, 2, 3, 4, 5 }, new[] { 6, 7, 8, 9, 10 }, new[] { 11, 12, 13, 14, 15 } }; var sumTasks segments .Select(seg Task.Run(() SumArray(seg))) .ToArray(); var totalTask Task.Factory.ContinueWhenAll( sumTasks, completed completed.Sum(task task.Result), CancellationToken.None, TaskContinuationOptions.OnlyOnRanToCompletion, TaskScheduler.Default ); Console.WriteLine(总和: totalTask.Result);这里的ContinueWhenAll泛型版本返回Taskint你可以直接获取合并结果。将这个写法与Task.WhenAll对比var sums await Task.WhenAll(sumTasks); var total sums.Sum();后者明显简洁。前者更大的意义在于你可以把“计算总和”这个动作本身作为一个可等待的Task对象往下传甚至在它之后继续挂新的ContinueWith。如果你是在构建一个流水线式的任务链每个节点都需要统一处理那么ContinueWhenAll这种“动作即任务”的表达方式在动态构建链条时会比较顺手。3.3 场景三老代码 Task 链中的延续接力我有一次维护老项目里面有一段代码用ContinueWith实现了多阶段处理管道任务 A 完成后启动任务 B 和任务 C 并行等 B 和 C 都完成后再启动 D。代码结构大概是这样的var bAndC Task.Factory.ContinueWhenAll( new[] { taskB, taskC }, completed D(completed) ); var final bAndC.ContinueWith(bc Finalize(bc.Result));在这种老式任务链里ContinueWhenAll承担着“分支汇合点”的作用。你想用 async/await 重写当然可以但改动范围会很大。如果你能读懂这种链式结构维护起来心理压力会小很多。你只需要知道bAndC这个任务只有在taskB和taskC都完成时才会启动而final则是在bAndC完成后继续执行。它本质上就是一个基于回调的分步流水线。3.4 执行这段代码前必须知道的几个约定写ContinueWhenAll时有几个约定我想放在实操部分强调因为踩过坑印象太深了。第一前序任务数组必须已经处于被调度的状态。如果你把一堆new Task(...)尚未调度的任务传进去这些任务永远不启动延续也就永远不会触发。最常见的正确做法是先用Task.Run或Task.Factory.StartNew把任务启动再把它们作为数组传进去。第二延续回调里面不要直接做阻塞等待。延续回调本身就运行在一个后台线程上你在这个回调里再调用.Wait()如果线程池线程数不够很容易把线程池资源耗尽。正确做法是让回调本身“小且快”或者通过 async/await 配合WhenAll处理更复杂的等待逻辑。第三注意延续返回的任务需要被观察。ContinueWhenAll返回的也是一个Task如果它的回调里抛出异常而这个任务又没有被人等待或检查Exception异常可能被吞掉或触发TaskScheduler.UnobservedTaskException。我一直默认给延续任务挂一个观察逻辑至少要在合适的地方Wait()或者ContinueWith(t { var _ t.Exception; })一下。4. 踩坑实录ContinueWhenAll 的经典问题与排查方法4.1 问题一延续不执行这是群里被问得最多的问题。代码逻辑看着没问题但ContinueWhenAll里的回调就是没跑。排查顺序我一般是这样的先确认所有前序任务都已经启动。Task.Run返回的任务已经是热任务但如果你手动new Task后忘了Start()任务永远处于Created状态延续永远等不到完成信号。再确认前序任务没有因异常或取消而进入非预期状态。ContinueWhenAll默认是“全部完成就继续”这里的“完成”包含RanToCompletion、Faulted、Canceled。但如果你写了OnlyOnRanToCompletion只要有任何一个任务失败延续就不会执行。这未必是 bug可能是选项和业务意图不匹配。最后检查延续任务本身是否被取消。如果传入了已经取消的CancellationToken延续任务会直接进入取消状态回调不执行表面上看起来就像“延续没跑”。这种情况代码里通常不会有明显输出需要看Task.Status才能发现问题。我实际见过一个案例同事把CancellationTokenSource在别处提前Cancel()了然后一直怀疑是ContinueWhenAll的 bug。检查顺序对了马上就能定位。4.2 问题二在延续里访问 Result 抛出 AggregateException延续回调里经常会写task.Result因为大家觉得“既然所有任务都完成了拿结果应该没问题”。这种直觉在OnlyOnRanToCompletion场景下是对的但如果用了默认None就要小心了。当某个前序任务以Faulted状态结束时它内部记录的AggregateException不会自动抛出来可你一旦访问task.Result这个任务会重新抛出它的异常。结果就是你的延续回调会在一瞬间炸掉而且因为异常发生在延续任务内部你甚至不一定能在原地捕获到。排查的时候很容易绕弯路。我的建议是如果前序任务存在失败可能延续回调里先把task.Status检查一遍再决定是否访问Result。不要抱有侥幸心理。你可能会觉得OnlyOnRanToCompletion能规避这个问题但注意OnlyOnRanToCompletion是“所有前序任务都成功才执行”假如你只想让延续在某些子任务失败时也跑那还是得靠状态检查。4.3 问题三阻塞等待造成死锁或线程饥饿这是和 TPL 相关的经典问题。很多人写了延续之后会立刻在同一个线程里调用continuationTask.Wait()等结果。如果襟翼某些线程相关条件很容易死锁或线程饥饿。具体场景是这样线程池线程有限你把一批长任务都扔进线程池然后自己又Wait()在其中一个线程上。线程池发现线程不够会调度新线程但调度过程有延迟。如果任务里也有阻塞等待新线程又被占住整个系统就陷入饥饿状态。ContinueWhenAll根本得不到执行机会。我通常的处理原则是能用 async/await 就绝不用.Wait()因为await不会阻塞线程而是释放线程让其他任务执行。如果确实需要同步等待最好限定超时时间例如Wait(TimeSpan.FromSeconds(5))然后判断返回值不要无限期等下去。4.4 问题四空数组与取消行为的边界ContinueWhenAll传空数组的情况你可能以为会立即返回一个已完成的任务。我刚开始也是这么想的实际行为却出乎意料。TPL 的实现中空数组会被解释为“所有前序任务都已完成”的边界情况但延续任务真正执行与否还受continuationOptions和取消令牌影响。虽然从语义上可以把空数组视为“全部完成”但不同重载在组合OnlyOnRanToCompletion和取消令牌时的行为可能不一致最好不要依赖这种边界条件。我的建议非常简单在上层显式拦截空数组。如果前序任务列表为空等于这个合流点没有意义直接跳过或返回一个已完成任务而不是测试关于空的边界行为。4.5 问题排查速查表现象可能原因排查动作延续回调完全不执行前序任务未启动检查任务是否已Start()或来自Task.Run延续回调不执行但任务状态 Faulted启用了OnlyOnRanToCompletion查看前序任务异常调整选项回调执行后程序崩溃访问了Faulted任务的.Result改用状态检查或使用Exception属性等了很久没返回当前线程被阻塞延续任务线程饥饿改用await避免同步等待传入空数组结果不确定边界条件被依赖显式处理空列表取消令牌触发后回调仍执行取消令牌只作用于延续而非前序任务检查参数语义重新设计取消时机5. 一些个人建议这个 API 该怎么用5.1 现代写法的选择顺序如果让我从零开始写新功能我的选择顺序是这样的优先用 async/await 加Task.WhenAll这是可读性和异常处理最好的方案。只有当我要面对老代码的延续链、动态构建任务链节点、或者确实需要在一个延续任务上继续挂多个后续动作时我才会用ContinueWhenAll。为什么我不提倡新代码里大量使用它因为延续式代码的异常流实在不透明。前序任务失败、延续选项控制、回调里异常状态这些在 async/await 里都能用 try/catch 压掉但在延续模型里要操心的事情太多。可如果已经是延续式结构强行半路改成 async/await也会让代码风格很割裂。我的习惯是在一段代码里保持同一种并发风格不要混着来除非只是为了某个特定节点做隔离。5.2 沿用旧代码的适配策略老代码里的ContinueWhenAll往往不是孤立出现的它可能连着ContinueWith、ContinueWhenAny串成一长串。你要修改这部分逻辑时不要急着整体重写先画清楚任务依赖图每个任务的起点是谁合流点是谁分支点是谁。ContinueWhenAll就是合流点ContinueWhenAny是抢先点ContinueWith是接力点。画明白之后你才能确定每个任务的异常会流向哪里。我自己的实践是对遗漏异常的任务链加一层全局观察任务var wrapped continuationTask.ContinueWith( t { if (t.IsFaulted) { Log(t.Exception); } }, CancellationToken.None, TaskContinuationOptions.OnlyOnFaulted, TaskScheduler.Default );这样即使延续内部出了问题也至少有一个观察任务把异常记录下来不会让异常悄无声息地丢失。5.3 把延续任务当作观察点来用最后分享一个小技巧。ContinueWhenAll返回的延续任务本身可以当作一个“观察点”挂到其他逻辑上。比如你要监控一组后台任务的完成情况可以写var monitor Task.Factory.ContinueWhenAll( workerTasks, completed { var successCount completed.Count(t t.Status TaskStatus.RanToCompletion); var failureList completed .Where(t t.Status TaskStatus.Faulted) .Select(t t.Exception) .ToList(); Report(successCount, failureList); }, CancellationToken.None, TaskContinuationOptions.ExecuteSynchronously, TaskScheduler.Default );这里我用了ExecuteSynchronously因为回调逻辑很轻量数一下成功数量汇总一下异常列表。这么写的好处是主流程完全不用关心这组任务的内部细节只需要在一个全局监控任务里统一收集状态。等到你想重构的时候观察点又变成批量读取结果的地方。这种“把合流点做成监控钩子”的思路在处理后台批处理任务时非常实用。我在实际项目里见过很多ContinueWhenAll的误用也见过把它用得行云流水的代码。它就像一把老式扳手新项目里你可能更喜欢棘轮扳手的爽快感但当你在现场遇到那颗老螺母时老式扳手是唯一能卡住位置的工具。理解它、会用它是基本功至于什么时候掏出来用那就看你对整个任务依赖图的理解了。
返回列表