多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

合并ExecuteFlowAsync:Parallel属性驱动的异步流程调度器重构

合并ExecuteFlowAsync:Parallel属性驱动的异步流程调度器重构 先说个背景。老项目里有两个版本的ExecuteFlowAsync一个在调度服务里乖乖串行跑一个在批处理模块里并行乱飞。两个方法同名不同命谁也不敢乱动。这次任务很明确把两个版本合并成一个最优解而且必须支持Parallel属性来决定模块是串行还是一个接一个地并行起飞。如果你也遇到过“同一个异步流程执行方法被拆成两版、最后要合并”的局面这篇文章应该能帮你少走不少弯路。我会从版本差异、方案选型、代码落地到现场踩坑完整捋一遍。1. 两个版本的设计冲突不只是方法重载那句“合并它们”合并代码之前先把两个版本为什么存在讲清楚。表面上是同名方法的重复实现本质上是两套不同的业务心智模型被硬塞进了同一个命名空间。1.1 串行版本流程的顺序性是它的生命线老调度服务里的串行版ExecuteFlowAsync处理的是那种“前一步产出没落定后一步就不能启动”的任务链。比如数据抽取后必须落盘、落盘成功才能触发血缘计算、血缘计算结束后才能回写状态。这种场景下模块的执行顺序是业务规则的一部分不是可选项。串行版本的核心保障有两个。一是严格的有序性for 循环按模块注册顺序执行前一模块的异常会短路后续所有模块保证流程不会在一个坏结果之上继续叠加执行。二是上下文的可变性安全因为同一时刻只有一个模块在写上下文模块之间可以直接共享状态对象不需要加锁也不需要并发容器。但这个版本的难受之处在于吞吐如果流程里有 5 个模块它们之间没有真正的数据依赖只是历史习惯把它们排成了一条线那么串行执行的时间就是所有模块耗时之和。接口检测、批量数据上报、推送通知这类 IO 密集型模块叠在一起一个 3 秒、一个 2 秒、一个 4 秒整条链就要 9 秒但其实并行起来可能只要 4 秒。1.2 并行版本为独立任务而生的“放养模式”批处理模块里的并行版ExecuteFlowAsync则是另一个极端。当时的业务场景是一批互相独立的运维任务比如同时去拉取多个业务系统的健康状态、同时刷新多张缓存、同时触发多份报表重算。这些模块唯一的共性就是“都归我这个服务管”彼此之间没有任何数据往来于是设计者很自然地用Task.WhenAll把所有模块一把梭跑起来。并行版本的优势很直白能并行的一定并行吃满 IO 并发的红利。但它几乎把所有维护问题都留给了后来人模块共享的上下文对象线程不安全、某个模块抛异常会导致其他模块被连带取消或者被整体吞掉、模块启动顺序和完成顺序都不确定排查问题时没法靠“看日志顺序”快速还原现场。更麻烦的是业务后来演进出了“部分依赖”流程 A 依赖流程 B 的结果但流程 C、D、E 又完全独立。老的并行版本没法表达这种中间状态只能要么全部并行跑错要么把 A、B 打包成一个模块勉强能跑但代码丑陋。1.3 合并执行器面临的三个核心矛盾把两个版本摆到一起合并的核心矛盾其实就三条第一个矛盾是执行语义的统一。两个版本都叫ExecuteFlowAsync但一个承诺“按顺序执行、异常短路”一个承诺“一起触发、结果聚合”。调用方没法只凭方法名知道行为必须去看内部实现。合并的目标就是让行为由显式配置决定而不是由“你在哪个项目里”决定。第二个矛盾是上下文共享模式的冲突。串行场景下模块之间可以直接共享可变对象并行场景下共享同一个引用必然导致脏读脏写。合并版本必须既保留串行时那种“随便读写上下文”的爽快感又保证并行时数据不乱。第三个矛盾是重试与排错的模型。串行版本只需要关心“哪个模块挂了它前面那些模块是否要回滚”并行版本要关心“5 个模块挂了俩另外 3 个要不要继续跑已经跑完的结果是否有效”。这两个响应策略很难用一套代码无脑兼容必须把决策权交还给调用方。合并不是把两个方法的代码并到一个文件里而是重新设计一个执行器让“串行”和“并行”成为同一个执行器在不同配置下的两种运行姿态。2. 方案选型用 Parallel 属性驱动一个执行引擎既然目标是提供最优解我第一个排除的就是“保留两个方法内部各写各的”。原因很直接你一旦保留两个入口就永远有调用方问“我该用哪个”也永远有人会在串行需要的场景里误用并行方法。正确的做法是统一收敛到单一的ExecuteFlowAsync入口让Parallel属性作为执行策略的开关。这个开关决定了引擎内部走串行调度器还是并行调度器但两个调度器共用同一套模块解析、上下文传递、异常处理框架。2.1 为什么坚持“一个入口 属性开关”从调用方的视角看ExecuteFlowAsync变成了一句话就能描述的任务执行器“给我一个流程定义我按你的要求跑完它。”至于流程内部是串行还是并行是调用方在定义流程时就已经确定的事情而不是执行器自作主张。这里还有一层额外的好处调用方可以用同一个方法名去承载不同配置的流程甚至未来可以在不改方法签名的情况下扩展第三种执行模式比如按依赖关系分阶段执行。接口不膨胀API 面稳定对下游调用者来说这就是“最优解”的一个重要维度。从排查问题的视角看单入口单日志也远优于双入口双日志。线上只需要搜索ExecuteFlowAsync一个调用点就能列全所有流程执行记录不用在串行日志和并行日志之间人工对齐时间线。2.2 两个调度器之间哪些逻辑必须统一真正容易踩坑的地方在这里设计师很容易把“串行实现”和“并行实现”想象成两个毫无交集的独立方法结果就是同样的错误处理、同样的超时控制、同样的上下文初始化逻辑被复制了两份改一处漏一处。我在合并时强制抽出了三层公共逻辑第一层是流程定义解析。无论Parallel是 true 还是 false都要先做同样的模块实例化、依赖注入、参数绑定和配置校验。这些动作不应该因为并行与否而有所不同。第二层是模块执行外观。每个模块都实现同一个IExecutableModule接口对外统一表现为ExecuteAsync(context, cancellationToken)。调度器只关心这个外观不关心模块内部是调接口、读数据库还是算算法。第三层是异常与结果收集。串行模式需要记录“第几个模块失败、失败原因、前面哪些模块已成功”并行模式需要记录“哪些模块成功、哪些模块失败、失败的完整异常堆栈”。两套数据格式必须一致这样上层报警和重试逻辑可以复用同一套代码。2.3 串行与并行分支处真正决策的是“调度策略”Parallel属性表面上只会引导执行走if / else两个分支但枝分背后的调度策略值得按业务场景细化。我按以下三种形态来做选择第一种是纯顺序短路调度Parallelfalse模块按索引逐个执行某个模块抛异常或主动设置StopRequested后立即中断适合事务性场景。第二种是全量并行调度Paralleltrue且流程定义中没有依赖关系所有模块直接批量触发等待全部完成适合互不干涉的独立任务。第三种是依赖感知的分阶段并行调度Paralleltrue且流程定义中存在DependsOn关系时执行器把模块拆成多个批次同一批次内并行执行批次之间按依赖关系串行推进。这一种最接近我们常说的“并行执行 linux 命令”的思维模型先把能同时跑的命令一次性发到后台等它们陆续完成后再执行依赖它们结果的后续命令。xargs -P就是经典的并发控制模型C# 里用批处理循环实现同样的效果。3. 拆解核心细节上下文隔离、依赖关系与并发上限执行器合并不是把foreach改成Task.WhenAll就行真正的工程量在于把并行模式下“共用同一套上下文”的隐患逐个排掉。3.1 FlowContext 线程安全改造从普通字典到 ConcurrentDictionary串行版本的上下文最初就是一个普通Dictionarystring, object模块 A 写进去的东西模块 B 立刻就能读没有任何并发问题。但并行模式下5 个模块同时往同一个字典里写context.Data[cost] xxx轻则互相覆盖重则直接引发异常。合并时我把Data换成了ConcurrentDictionarystring, object?并约定并行模式下不允许模块之间通过共享引用来传递可变对象。如果模块 A 产出一个结果列表需要通过context.Data[biz_list]传给模块 C 使用模块 A 必须写入一个不可变或者深拷贝后的对象否则 C 在读取时可能在遍历列表的同时被 A 修改。这里有一个实际的取舍ConcurrentDictionary保证的是单个键值对的原子性不保证读写对象的线程安全性。所以我在流程定义里增加了一个开关SharedStateIsolation默认开启开启时执行器在串行模式下直接用普通共享对象并行模式下要求模块传入的共享对象实现ICloneable或者使用不可变数据结构。如果你在并行场景里发现偶发的“上下文数据错乱”八成就是这个位置出了问题。3.2 并行中的依赖表达DependsOn 与分阶段执行并行和依赖看似矛盾但在实际业务里必须共存。比如“并行采集 5 个数据源全部采完后统一做清洗入库”采集阶段可以并行清洗阶段必须等采集全部结束。如果执行器只有“全并行”和“全串行”两种模式这种流程就只能写成“先并行执行采集模块组再串行执行清洗模块组”硬拆两个阶段。我的做法是给流程定义增加一个DependsOn字典值是被依赖模块的名称列表。执行器在并行模式启动时先做一次依赖关系解析构建一个按批次划分的执行计划第一批没有被任何模块依赖或者无依赖前置的模块全部并行执行。后续批次只有当前批次全部成功后才启动依赖已满足的下一批模块。如果解析后发现依赖循环直接抛出InvalidOperationException避免运行时死锁。批次之间用Task.WhenAll聚合批内并发数量还要受到并发上限约束。这个设计最初只是顺手加的后来成了整个方案里最被业务方称道的一点因为很多本来需要“拆流程、分两次跑”的需求直接被一个流程定义覆盖了。3.3 并发上限控制并行不等于无限并发无脑并行是上线前最喜欢炸的一颗雷。老的并行版本曾经同时拉起 100 个模块线程直接把数据库连接池打爆整条链路的成功率从 99% 掉到 70%。合并时我给执行器加了一个重要的概念最大并行度MaxParallelism。推进时时按照一定规则运行。实现上我用了SemaphoreSlim控制活动模块数。默认值设置为Environment.ProcessorCount * 2这个值不是拍脑袋定的而是结合“模块多是 IO 密集CPU 密集只占少数”这个经验法则。如果业务方确认所有模块都是 IO 密集可以把值放宽到 4 倍核数如果有 CPU 密集模块混入则保守设为核数 1。线程安全上还有个小细节值得注意SemaphoreSlim的WaitAsync必须通过块的形式传入模块重试逻辑否则模块内部抛错时信号量可能没有被释放导致并发上限慢慢耗尽、后续所有模块全部死等。这个坑我在第 5 节还会详细展开。4. 实操记录合并重构 ExecuteFlowAsync 的完整演进前面讲的是理念从这节开始是能直接抄作业的落地过程。我按重构时的实际操作顺序来写先理清调用方再实现核心代码之后做异常聚合最后验证性能。4.1 重构前先做调用方盘点而不是先改执行器很多人一拿到合并任务就急着写代码结果改完执行器发现调用方传参方式五花八门不得不回头改接口。我强烈建议第一步先用代码搜索把两个版本的所有调用点拉出来按业务场景分一下类纯顺序型调用模块之间明显有先后的这类调用方基本不需要改配置。纯并行型调用模块之间无依赖原来已经用Task.WhenAll的改成Paralleltrue即可。混合型调用模块之间部分有依赖、部分独立需要重新设计依赖表达是改动最大的类型。盘调用方的主要收益是能提前确认新配置项的默认值应该怎么设如果 80% 的旧调用都是串行场景新执行器的Parallel默认值就应该是 false避免老调用方在不改代码的情况下被“悄悄并行化”。4.2 核心代码实现模块接口、流程定义与执行引擎先把最核心的模块接口定义出来。这里不搞花活儿一个接口只负责一件事调度器靠这个接口统一驱动所有模块public interface IExecutableModule { string Name { get; } Task ExecuteAsync(FlowContext context, CancellationToken cancellationToken default); }接下来是流程定义。Parallel属性是本节的核心DependsOn用于支持依赖感知的并行MaxParallelism可以在流程级别覆盖引擎默认并发上限public class FlowDefinition { public string FlowId { get; set; } string.Empty; public bool Parallel { get; set; } public ListIExecutableModule Modules { get; set; } new(); public Dictionarystring, Liststring DependsOn { get; set; } new(); public int? MaxParallelism { get; set; } }执行上下文用ConcurrentDictionary作为共享数据存储同时维护一个取消信号便于并行场景下某个模块失败时通知其他模块停止无意义的工作public class FlowContext { public string FlowId { get; init; } string.Empty; public ConcurrentDictionarystring, object? Data { get; } new(); public bool StopRequested { get; set; } public CancellationTokenSource LinkedCts { get; set; } new(); }最后是合并后的执行引擎。核心思路已经在前面说过了一个入口按Parallel属性分流到两个调度器但两个调度器都通过RunSingleModuleAsync这个统一外观来执行模块。这样无论串行还是并行模块的超时控制、日志埋点、计数器自增都只有一份实现public class FlowEngine { private readonly ILoggerFlowEngine _logger; private readonly SemaphoreSlim _engineLimit; public FlowEngine(ILoggerFlowEngine logger) { _logger logger; _engineLimit new SemaphoreSlim(Environment.ProcessorCount * 2); } public async Task ExecuteFlowAsync( FlowDefinition flow, FlowContext context, CancellationToken cancellationToken default) { ArgumentNullException.ThrowIfNull(flow); ArgumentNullException.ThrowIfNull(context); if (flow.Modules.Count 0) { _logger.LogWarning(流程 {FlowId} 未配置任何模块, flow.FlowId); return; } ValidateDependencies(flow); if (!flow.Parallel) { await ExecuteSequentialAsync(flow, context, cancellationToken).ConfigureAwait(false); } else { await ExecuteParallelAsync(flow, context, cancellationToken).ConfigureAwait(false); } } private async Task ExecuteSequentialAsync( FlowDefinition flow, FlowContext context, CancellationToken cancellationToken) { for (int i 0; i flow.Modules.Count; i) { cancellationToken.ThrowIfCancellationRequested(); var module flow.Modules[i]; await RunSingleModuleAsync(module, context, cancellationToken).ConfigureAwait(false); if (context.StopRequested) { _logger.LogInformation(流程 {FlowId} 收到停止信号在第 {Index} 个模块后中断, flow.FlowId, i); break; } } } private async Task ExecuteParallelAsync( FlowDefinition flow, FlowContext context, CancellationToken cancellationToken) { var tasks flow.Modules.Select(async module { await _engineLimit.WaitAsync(cancellationToken).ConfigureAwait(false); try { await RunSingleModuleAsync(module, context, cancellationToken).ConfigureAwait(false); } finally { _engineLimit.Release(); } }); await Task.WhenAll(tasks).ConfigureAwait(false); } private async Task RunSingleModuleAsync( IExecutableModule module, FlowContext context, CancellationToken cancellationToken) { var sw ValueStopwatch.StartNew(); try { _logger.LogInformation(开始执行模块 {ModuleName}, module.Name); await module.ExecuteAsync(context, cancellationToken).ConfigureAwait(false); _logger.LogInformation(模块 {ModuleName} 执行成功耗时 {Elapsed}, module.Name, sw.Elapsed); } catch (Exception ex) { _logger.LogError(ex, 模块 {ModuleName} 执行失败耗时 {Elapsed}, module.Name, sw.Elapsed); throw; } } }这段代码看起来不算长但它在合并过程中的价值很高。串行版本和并行版本都收敛到了同一个RunSingleModuleAsync后续加模块级重试、加熔断、加自定义超时都只需要动这一个点。4.3 依赖感知的并行按批次执行而不是无脑 WhenAll全量Task.WhenAll在“所有模块都独立”的场景下够用但碰上部分依赖就必须升级。我给出的实现方案是按依赖关系把模块分批每批内并行执行批次之间串行等待。这个模式很像在 shell 里用xargs -P控制并发一样语义简单、行为可预测。private async Task ExecuteDependencyAwareParallelAsync( FlowDefinition flow, FlowContext context, CancellationToken cancellationToken) { var pending new HashSetIExecutableModule(flow.Modules); var completed new HashSetobject(); while (pending.Count 0) { cancellationToken.ThrowIfCancellationRequested(); var ready pending .Where(m !flow.DependsOn.TryGetValue(m.Name, out var deps) || deps.All(d completed.Contains(d))) .ToList(); if (ready.Count 0) { throw new InvalidOperationException($流程 {flow.FlowId} 存在循环依赖无法继续执行); } var batchTasks ready.Select(async module { await _engineLimit.WaitAsync(cancellationToken).ConfigureAwait(false); try { await RunSingleModuleAsync(module, context, cancellationToken).ConfigureAwait(false); } finally { _engineLimit.Release(); } }); await Task.WhenAll(batchTasks).ConfigureAwait(false); foreach (var module in ready) { pending.Remove(module); completed.Add(module.Name); } } }这里有一个容易被忽略的业务设计依赖关系应该是模块粒度而不是阶段粒度。如果你只定义“阶段1 并行、阶段2 并行”那么每次加模块都要重新划分阶段非常容易错。改为DependsOn之后新增一个模块只需要声明“我依赖谁”执行器自动把它调度到正确的批次里维护成本大大降低。4.4 异常聚合与失败语义负责任地“炸”串行版本的异常处理很简单某个模块抛了直接向外抛调用方捕获后知道是第几个模块干的。并行版本的问题在于Task.WhenAll的默认行为是抛出第一个异常剩下 4 个模块的异常全被丢弃了。这在生产环境是灾难性的因为你看到的第一条异常不一定是最根因而其他模块的失败可能才是压垮服务的真正原因。我的处理方案是双重保险RunSingleModuleAsync内部已经把异常和模块名写进日志调度器则捕获Task.WhenAll整体异常并把每个模块的异常包装成ModuleExecutionException异常消息里包含模块名。这样外层感知到的是“流程执行失败”但通过日志和异常对象能精确还原每个失败模块的明细。并行模式下还有一个“失败后其他模块是否继续”的决策要交给调用方。默认行为是继续跑因为模块已经启动了中途取消反而可能造成数据半提交执行器只负责把已发生的失败如实汇报。如果调用方希望一损俱损可以在流程定义里打开FailFast开关该模式下会启用LinkedCts某个模块一失败就触发整体取消令牌其他模块在下一个取消检查点停止。5. 常见问题与排查实录并行模式坑人实录合并重构上线一周我差不多把下面这些坑挨个踩了一遍。有些是设计层面的有些是运行时才暴露的写出来给大家当体检清单。5.1 并行导致共享状态变量互相覆盖现象并行模式下有个模块执行结果时灵时不灵串行模式下完全正常。排查后发现两个模块执行完都会往context.Data[final_data]写值原本设计是模块 A 写创建时间模块 B 写更新时间但并行时 B 可能先写完A 再覆盖掉它。处理办法就是前面提到的并行模式下一条数据只允许一个模块写其他模块要么放在不同 key 下要么等待具有写权的模块完成后通过依赖关系再执行。不要靠“执行顺序”来保护共享变量的写入用依赖关系来保护顺序是最脆弱的约定。5.2 Task.WhenAll 异常被吞日志里只有一条异常现象5 个模块并行执行3 个失败但日志里只看到最早抛出的那一条异常。通过Task.WhenAll的Result属性可以发现它其实包含了InnerExceptions但太多人只catch到最外层那个就结束了。排查建议在并行模式的catch (Exception ex)里额外记录ex is AggregateException agg ? string.Join(;, agg.InnerExceptions.Select(e e.Message)) : ex.Message把每次异常聚合的全貌打印出来。我已经把这一步放进了前面的RunSingleModuleAsync日志里但如果你用的是自己的执行器实现请务必把这一步加进去。5.3 并发上限把数据库连接池打满现象并行开关一开数据库连接池报timeout expired但单模块单独跑毫无压力。原因很典型8 个模块同时执行每个模块又各自开了 10 个数据库连接瞬间 80 个连接向连接池申请全部排队等资源。处理办法是在引擎层限制模块并发数的同时给每个模块内部的连接池也设置上限或者在流程定义里根据模块类型调整MaxParallelism。更稳妥的做法是给数据库访问模块加独立的局部SemaphoreSlim避免执行器层的上限控制不了模块内部的资源消耗。5.4 快速排查速查表现象可能原因优先排查方向并行执行结果随机不一致共享可变对象被多模块并发写检查context.Data写入方是否唯一异常只看到一条其余丢失Task.WhenAll只抛第一个异常改用异常聚合并记录全部InnerExceptions连接池超时模块并发数超过数据库连接池上限降低引擎并发度或模块内部加局部信号量偶发死锁模块内部同步等待异步任务全链路去掉.Result/.Wait()改用 async 协程等待没设 stop 信号时也中断了默认取消令牌与流程级取消混淆区分引擎默认 Ct 和流程私有 Ct 的作用域6. 关于“最优解”的一点个人体会合并完两个版本后我回头看这个任务感触最深的一点是最优解不是“代码看起来多优雅”而是执行器的行为能被调用方确定性预测。串行还是并行由一个Parallel属性显式控制比让调用方猜“这个方法内部是串行还是并行”要确定得多。与我预想的相反这个执行器在真实业务里跑得最漂亮的地方反而是混合模式——模块个数少、依赖复杂、但又想尽量并行。依赖感知的分阶段并行刚好踩中了那个需求区间没有让流程调度变成一锅粥。如果让我给后来者一条建议合并执行器之前别急着抠代码细节先把不同调度策略对上下文、异常、资源上限的影响捋清楚。你可以在上线前靠代码评审解决 80% 的问题剩下 20% 的坑会让你的监控和日志设计自动补上。最后别把Parallel属性当成银弹。它只是个开关真正决定执行质量的是开关背后那些关于依赖、隔离和上限的小决策。希望这篇记录能帮你少走几个弯路。
返回列表