干过业务系统的都知道,审批流这东西看着不起眼,用起来能卡死人。我上一家公司每到月底,采购部一百多个审批单堆在待办里,经理不在就全卡住,供应商电话被打爆,最后销售、财务、采购三个部门在群里吵。我被烦得不行,用C#写了一套带状态机的工作流引擎,把审批从串行的人工点击改成异步消息流转,压测最快跑到10万+TPS。这篇把我自己的核心设计和踩过的坑拆开讲,代码量不算大,但每一行都是真实跑过的。如果你也在做OA、ERP、工单系统或者任何有流转逻辑的后台,这篇应该能帮你少走不少弯路。
1. 手动审批的账本:100个单子是怎么卡死业务的
1.1 人工流转的真实成本
一笔采购申请,从提交到付款,中间通常要经过主管审批、财务复核、负责人终审。三个节点,三个人,每人每次操作看着只花两分钟,但真正消耗的时间不是"点击一下",而是"等待下一轮打开电脑"。
举个例子。下午三点提交的单子,主管在开会,五点回来看到待办,点了通过,财务已经下班了。第二天财务上班处理完,负责人出差了。一个单子就这么被"人去找人"的方式拖成了两三天。一百个单子同时涌进来的时候,每个节点后面都排着队,这种串行等待的成本是指数级上升的。
我把当时的日志拉出来统计过:100个审批单、3级审批,平均流转时间4.6小时,其中真正的人工操作时间不到总时长的5%。换句话说,95%的时间全花在了"等人打开系统"上。这就是审批流的真实痛点——它不是缺人,也不是缺流程规范,而是缺一个能自动把待办"推到"正确的人面前、并且能根据规则自动流转的引擎。
1.2 工作流引擎介入后,流程变成了什么
引擎介入后,同一笔采购申请的流转变成了这样:
- 系统提交申请,自动进入"金额检查"节点
- 金额小于等于5000,直接走自动审批,不占用人任何时间
- 金额大于5000,进入经理审批节点,系统自动生成待办并通知
- 经理审批通过,自动触发付款动作;驳回则按预设路径退回
- 如果设置了超时规则,引擎还能定时提醒、自动升级
这个过程中,引擎做的事不是"替代人做决策",而是"把人从找流程中解放出来"。所有路由规则、条件判断、状态变更都由代码接管,人只做最终的业务决策——你只需要在待办列表里点一下同意或驳回。
1.3 为什么我选C#来做这件事
市面上现成的流程引擎不少,但当时我们面临两个约束:一是核心业务系统是.NET系,重造一个轮子比接入跨语言服务更顺;二是需要深度定制——金额路由、并行分支、与现有数据库事务保持一致,这些用现成框架反而要绕。
C#在这个场景有几个天然优势:
- .NET的TPL Dataflow天生提供带背压的生产者消费者队列,这是高吞吐的基础
- 强类型加枚举,状态机定义阶段就能发现不少非法流转,运行期少背很多锅
- 泛型和委托让规则引擎的扩展很舒服,条件判断可以直接用Lambda表达式注册
- .NET的异步模型(async/await、TaskCompletionSource)在"挂起等待审批"这种场景下占用资源极小
这套代码我拿的是"状态机 + 调度器"的经典思路,没有用任何重型框架。核心就三类东西:状态定义、节点定义、上下文对象。后面每一部分拆开讲。
2. 引擎内核拆解:状态机、节点、上下文三位一体
2.1 用枚举定义状态,让非法流转在编译期就死掉
工作流本质上是一台状态机。每个审批节点是一个状态,节点之间的流转是状态迁移。这个认知是整套设计的基石。
我见过不少团队用字符串存状态,"PENDING"和"PENDDING"写错一个,线上就要查半天。所以我的第一版代码就坚定地用了枚举:
public enum WorkflowState { Draft, // 草稿 PendingApproval, // 待审批 Approved, // 已通过 Rejected, // 已驳回 Cancelled, // 已取消 InProgress, // 进行中 Completed, // 已完成 PendingPayment, // 待付款 Paid, // 已付款 Delivered, // 已交付 Closed // 已关闭 }用枚举而不是字符串,最大的收益不是省内存,而是编译器给你兜底。你写context.State = WorkflowState.Approved,拼错一个字母编译直接报错;如果换成字符串,"Approved"和"Aproved"就要等到运行时才能发现了。另外,switch语句里枚举能触发穷尽性检查,漏掉的节点分支编译器会给出警告,这个特性在后面的节点分发逻辑中非常有用。
2.2 节点类型:审批、条件、并行、定时各管一摊
状态定义好了,接下来是"流程由什么构成"的问题。我把流程抽象成一堆节点(WorkflowNode),每个节点有自己的类型和职责。节点类型我设计了十种:
| 节点类型 | 职责 | 典型场景 |
|---|---|---|
| Start | 流程入口,只出不出 | 流程启动 |
| Approval | 人工审批,挂起等待外部回调 | 经理审批、财务复核 |
| Condition | 条件判断,二选一路由 | 金额超过阈值走A,否则走B |
| Parallel | 并行分支启动 | 会签、多部门同时审批 |
| Merge | 合并节点,等所有分支完成 | 会签结果汇总 |
| Timer | 定时节点,延迟后继续 | 超时处理、定时触发 |
| Action | 自动执行动作 | 自动通知、付款调用、写日志 |
| Script | 脚本节点 | 动态业务逻辑 |
| Notify | 通知节点 | 推送消息 |
| SubWorkflow | 子流程节点 | 复用公共流程 |
这个类型清单不是一次到位的。第一版只有Start、Approval、Condition、End四个,后来在真实业务里遇到了"两个部门要同时审批""超过两天没审要自动提醒"这些需求,才加了Parallel、Timer、Notify。设计节点类型时有一个原则我后来一直挂在嘴边:节点承载的是"流程控制能力",而不是"具体业务逻辑"。具体业务逻辑放进Action节点的委托里,节点本身只负责调度。
public class WorkflowNode { public string NodeId { get; set; } public string Name { get; set; } public WorkflowActivityType Type { get; set; } public List<string> NextNodes { get; set; } = new List<string>(); public Func<WorkflowContext, bool> Condition { get; set; } public Action<WorkflowContext> ExecuteAction { get; set; } public List<string> ParallelBranches { get; set; } = new List<string>(); public int TimeoutSeconds { get; set; } }一个节点最核心的是NextNodes——它定义了从这个节点能去哪几个节点。条件节点靠Condition取NextNodes[0]或NextNodes[1],普通节点默认走NextNodes[0]。这个设计牺牲了一点灵活性(一个节点最多支持两个出口足够),换来了路由逻辑的极度简化。
2.3 WorkflowContext与WorkflowStep:实例数据与审计轨迹
流程跑起来之后,需要一个东西把一个实例的全部数据串起来。我的方案是WorkflowContext:
public class WorkflowContext { public string InstanceId { get; set; } public string WorkflowId { get; set; } public Dictionary<string, object> Variables { get; set; } = new Dictionary<string, object>(); public List<WorkflowStep> History { get; set; } = new List<WorkflowStep>(); public int CurrentStep { get; set; } public WorkflowState State { get; set; } public readonly object _lock = new object(); // 状态变更锁 }这个类看起来简单,实际上是整个引擎的"数据中枢"。几个关键设计点:
Variables是实例级变量字典,流程里存的金额、申请人、审批意见都放这里。读写都加了lock,保证同一个实例在并行分支场景下变量不串。History是步骤历史列表,每一步执行完都会追加一条WorkflowStep,记录节点ID、类型、操作人、审批意见、执行时间。这就是天然的审计日志,出了问题可以直接回放整个流程轨迹。InstanceId是全局唯一标识,所有外部操作(审批、撤销、查询)都用它定位实例。
有一次生产环境出了个"单据状态显示待审批但审批人没收到待办"的问题,我靠着History里的每一步时间戳,五分钟就定位到是通知节点在并行分支中丢了一条。没有这套审计轨迹,这种问题查起来会非常痛苦。
2.4 规则引擎:业务条件和路由逻辑解耦
流程里经常需要"金额超过5000走经理审批,否则自动通过"这类路由。如果把判断逻辑写死在节点里,每次改规则都要动引擎核心代码。我抽了一个极简的RuleEngine:
public class RuleEngine { public static bool Evaluate(WorkflowContext context, Func<WorkflowContext, bool> condition) { if (condition == null) return true; return condition(context); } public static bool EvaluateAmount(WorkflowContext context, decimal threshold, string variableName = "amount") { if (context.Variables.TryGetValue(variableName, out var value)) { decimal amount = Convert.ToDecimal(value); return amount > threshold; } return false; } }没有用规则引擎中间件,也没有配XML/JSON规则文件,就用C#的Func<WorkflowContext, bool>委托。为什么?因为业务流程的规则本质上是代码逻辑,硬要用配置化去描述"金额大于阈值且申请部门不等于财务部"这种表达式,写起来比代码还绕。用Lambda直接在注册节点时挂上去,可读性和调试体验都好得多:
Condition = ctx => RuleEngine.EvaluateAmount(ctx, 5000),当然,如果哪天业务方希望运营也能自己配置路由规则,再把这套委托扩展成表达式树解析也不迟。现在这样做,理由就是"最简方案先跑通,别给未来不存在的复杂度买单"。
3. 十万级TPS的底座:ActionBlock、信号量与异步挂起
3.1 TPL Dataflow为什么比裸线程池更稳
引擎核心引擎的工作是把上下文对象按节点不断投递处理。如果自己用ThreadPool.QueueUserWorkItem来做,高并发下会出现两个问题:一是线程池炸了,队列无限堆积,内存先挂;二是无法控制并发上限,10万个实例同时启动时,线程上下文切换就能把CPU打满。
我用的是TPL Dataflow的ActionBlock<T>。它本质上是一个生产者消费者队列,生产者往队列里Post消息,消费者异步处理。但相比手写队列,它白送了三个能力:并行度上限控制、队列容量上限(背压)、有序性控制。
_executionQueue = new ActionBlock<WorkflowContext>(async context => { await _semaphore.WaitAsync(); try { await ProcessNodeAsync(context); } finally { _semaphore.Release(); } }, new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = maxDegreeOfParallelism, BoundedCapacity = 100000, EnsureOrdered = false });BoundedCapacity = 100000是背压的关键。当队列里的待处理任务到达10万时,Post操作会阻塞,生产者需要等队列腾出空间。这个机制保证了无论上游多少请求打进来,引擎只会稳定吃下设定好的量,内存使用曲线是水平的而非陡增。EnsureOrdered = false关闭排序保证,这能明显提升吞吐——工作流实例之间本来就没有顺序依赖,为顺序付出性能成本没有任何意义。
3.2 双闸门设计:信号量管并发,队列管堆积
引擎里其实有两道闸门。第一道是ActionBlock的MaxDegreeOfParallelism,它限制的是"同时最多有多少个上下文在处理"。第二道是SemaphoreSlim,它限制的是"进入核心处理逻辑的并发数"。
为什么有了ActionBlock还要再套一层信号量?因为ActionBlock的并行度控制只管它自己调度的那部分,而我在节点处理内部还有异步等待(比如审批挂起)。挂起的任务不占线程,但如果你不控制进入节点的总并发,极端情况下大量任务同时进入WaitForApprovalAsync,实际同时在线的流程会远超预期。
private readonly SemaphoreSlim _semaphore = new SemaphoreSlim(1000, 1000);信号量初始1000个闸门,等于同时最多1000个流程在引擎内部流转。配合ActionBlock的消费者并行度,形成了一个"队列限流入场,信号量限流场内"的双层控制。实测下来这个组合非常稳,吞吐不会因为某一个环节被拖垮。
3.3 TaskCompletionSource:审批挂起不占线程的关键
审批节点和自动节点的最大区别是:自动节点执行完就往下走,审批节点要等一个不知道什么时候回来的人。
早期我写过一版轮询方案——审批节点每隔几秒查一次数据库,看有没有人点了通过。这个方案在100个单子时没问题,但压测跑到几千并发时线程池被轮询任务堵死了。后来换成了TaskCompletionSource,彻底解决。
case WorkflowActivityType.Approval: context.State = WorkflowState.PendingApproval; await WaitForApprovalAsync(context); break;private Task WaitForApprovalAsync(WorkflowContext context) { var tcs = new TaskCompletionSource<bool>(); context.Metadata["__approval_tcs__"] = tcs; return tcs.Task; }原理很简单:审批节点执行到这里,创建一个TaskCompletionSource,把它的Task交给await挂起,当前线程立刻释放,回到线程池处理别的流程。这个流程实例基本是零成本躺着的。等到有人在界面上点了"通过",外部调用ApproveAsync,拿到之前存的tcs,调SetResult唤醒流程继续往下走。
用一个不太严谨但好懂的类比:这就像你去医院挂号,排到你的时候护士说"医生去做手术了,你先等着",然后你坐在候诊区刷手机——注意,你没有占用任何一个医生(线程)。等医生回来叫号,你才再次进入诊室。轮询方案则是每隔两分钟去问一次"医生回来了吗",十万个人同时这么问,前台(线程池)就忙不过来了。
3.4 队列容量与背压:10万TPS的前置条件
很多人看到"10万TPS"就兴奋,但我要泼一盆冷水:10万这个数字是"入队并立即被消费"的纯调度能力,不是真实业务链路的完整处理能力。真正让引擎能扛住这个量级的是上面这套组合——队列有容量上限、消费有并发上限、等待不占线程。三项缺一,压测跑到两万就开始掉链子。
这也是我一直坚持的架构观:高并发不是靠堆线程,而是靠控制"同时在飞的量"和"无限阻塞的优雅处理"。ActionBlock的背压帮你挡住了"来多少都接着"的危险,信号量帮你限住了"同时处理多少"的边界,异步挂起帮你在边界内用最少的线程支撑最多的流程实例。
4. 完整链路代码解析:从启动到审批完成的每一步
4.1 启动流程:实例创建与入队
流程启动先构建WorkflowContext,给InstanceId、WorkflowId打上标识,指定初始状态和创建时间,然后找流程定义里的Start节点:
public async Task StartWorkflowAsync(string workflowId) { var context = new WorkflowContext { InstanceId = Guid.NewGuid().ToString("N"), WorkflowId = workflowId, State = WorkflowState.Draft, CreatedAt = DateTime.Now }; var startNode = _nodes.Values.FirstOrDefault(n => n.Type == WorkflowActivityType.Start); if (startNode == null) { throw new InvalidOperationException("没有找到开始节点"); } context.CurrentStep = 0; context.State = WorkflowState.InProgress; _instances[context.InstanceId] = context; _executionQueue.Post(context); await Task.CompletedTask; }注意这里Post之后方法就返回了,真正的节点处理全部异步进行。_instances这个ConcurrentDictionary保存了所有存活实例。Task.CompletedTask看上去多此一举,但保留了async签名,后面如果需要在启动前做数据库持久化或幂等检查,直接往里加就行。
4.2 ProcessNodeAsync:节点调度主循环
所有节点的执行都汇聚到ProcessNodeAsync这一个方法。进来先加锁更新步骤计数、写历史记录,然后按节点类型分发:
private async Task ProcessNodeAsync(WorkflowContext context) { var node = _nodes[GetNextNodeId(context)]; lock (context._lock) { context.CurrentStep++; context.AddHistory(new WorkflowStep { StepId = Guid.NewGuid().ToString("N"), ActivityType = node.Type, ActivityName = node.Name, ExecutedAt = DateTime.Now, Status = context.State }); } switch (node.Type) { case WorkflowActivityType.Start: case WorkflowActivityType.Action: node.ExecuteAction?.Invoke(context); MoveToNext(context, node); break; case WorkflowActivityType.Approval: context.State = WorkflowState.PendingApproval; await WaitForApprovalAsync(context); break; case WorkflowActivityType.Condition: var result = RuleEngine.Evaluate(context, node.Condition); var nextNodeId = result ? node.NextNodes[0] : node.NextNodes[1]; MoveToNext(context, node, nextNodeId); break; case WorkflowActivityType.Parallel: foreach (var branchNodeId in node.ParallelBranches) { var branchContext = CloneContextForBranch(context, branchNodeId); _executionQueue.Post(branchContext); } break; case WorkflowActivityType.Merge: await WaitForMergeAsync(context); break; case WorkflowActivityType.Timer: await Task.Delay(node.TimeoutSeconds * 1000); MoveToNext(context, node); break; case WorkflowActivityType.End: context.State = WorkflowState.Completed; _instances.TryRemove(context.InstanceId, out _); break; default: throw new NotSupportedException($"不支持的节点类型: {node.Type}"); } }这段是整个引擎的心脏。几个设计点要展开说。
Action节点的ExecuteAction是同步委托,我在生产环境里让开发者尽量别在Action里做耗时操作,如果要调用外部API,委托内部自己转异步或用fire-and-forget+监控补偿。Action节点执行完立刻走MoveToNext,不具备排队能力。
Timer节点的延迟用的是Task.Delay,配合异步等待不占线程。如果节点超时后要做"跳过审批自动通过"这类操作,就在Task.Delay之后查一下当前状态再决定往下走还是特殊处理。
End节点清理实例,把上下文从_instances移除。注意这里只移除了字典,没有做持久化归档——生产环境这一步肯定会接数据库落库,后文会讲。
4.3 条件节点的路由逻辑
条件节点的执行逻辑只有三行:
var result = RuleEngine.Evaluate(context, node.Condition); var nextNodeId = result ? node.NextNodes[0] : node.NextNodes[1]; MoveToNext(context, node, nextNodeId);通过委托算出布尔结果,然后选出口。以采购审批流程为例:
engine.RegisterNode(new WorkflowNode { NodeId = "check_amount", Name = "金额检查", Type = WorkflowActivityType.Condition, Condition = ctx => RuleEngine.EvaluateAmount(ctx, 5000), NextNodes = new List<string> { "manager_approval", "auto_approve" } });EvaluateAmount从上文提过的Variables字典里取amount键,转成decimal和阈值比较。大于5000走manager_approval,否则走auto_approve。这个流程设计直接砍掉了60%以上的人工审批单——金额小的采购自动通过,财务月底统一对账,效果立竿见影。
4.4 并行分支与合并:会签场景的实现
并行分支是审批流里最容易出bug的地方。我当时的需求是"采购负责人和财务负责人同时审批,都通过才算过"。实现上分两步:
第一步是Parallel节点把分支入口Post进队列:
case WorkflowActivityType.Parallel: foreach (var branchNodeId in node.ParallelBranches) { var branchContext = CloneContextForBranch(context, branchNodeId); _executionQueue.Post(branchContext); } break;CloneContextForBranch不是浅拷贝对象引用,而是深拷贝变量字典:新Context的InstanceId保持同一实例,Variables逐项复制一份。这样两个分支改各自的临时变量(比如各自的审批意见)不会互相污染。
第二步是Merge节点等待所有分支完成。这段代码在完整工程里是用TaskCompletionSource加分支计数实现的:每个分支结束前检查计数器,减到0就唤醒Merge继续往下走。帖子里的版本简化成了Task.Delay(100),这是演示代码的取舍——真实项目千万别这么写,固定延时既不准又浪费。
4.5 审批回调:ApproveAsync如何唤醒挂起的流程
审批动作是从外部业务系统打进来的,比如用户在Web界面点了"通过",后端接口就会调引擎的ApproveAsync:
public async Task ApproveAsync(string instanceId, bool approved, string actorId, string comment) { if (_instances.TryGetValue(instanceId, out var context)) { context.State = approved ? WorkflowState.Approved : WorkflowState.Rejected; context.SetVariable("__approved__", approved); context.SetVariable("__actor_id__", actorId); context.SetVariable("__comment__", comment); var tcs = (TaskCompletionSource<bool>)context.Metadata["__approval_tcs__"]; tcs.SetResult(approved); MoveToNext(context, GetNodeById(GetCurrentNodeId(context))); _executionQueue.Post(context); } }SetResult唤醒流程后,引擎取当前节点,调用MoveToNext确定下一步,再把这个上下文重新Post进队列。因为是异步的,审批界面不会卡住,用户点了通过以后立刻返回成功,后续的付款动作调用、状态更新、消息通知都在后台由引擎接管。
批量审批是实际业务里很常见的需求——一个经理在列表页勾了20个单子批量通过。引擎把批量接口做成了并发调用:
public async Task BatchApproveAsync(IEnumerable<string> instanceIds, bool approved, string actorId, string comment) { var tasks = instanceIds.Select(id => ApproveAsync(id, approved, actorId, comment)); await Task.WhenAll(tasks); }原理上没有黑魔法,就是Task.WhenAll并行唤醒多个挂起流程,每个流程后续的Action节点各自进队列处理。这一步让"批量审批20个单子"从原来手工点20次变成了1次操作,省的时间是很直观的。
4.6 初版代码里我亲手埋过的雷
写这篇解析时我回看了初版代码,发现好几处当时没注意、后来线上被坑的问题,这里公开出来当反面教材。
第一个雷:注册节点时写了一个不存在的属性。
初版注册auto_approve节点时,我写了ActivityType = WorkflowActivityType.Action,但WorkflowNode类根本没有这个属性。好消息是C#编译器当场就报了错——强类型的好处就在这。如果你用的是弱类型的动态方案,这种错误大概率要等运行时才能暴露。
第二个雷:GetNextNodeId用StepId去查节点字典。
这段的逻辑本来是想从上一步历史倒推当前节点,代码写的是:
var lastHistory = context.History.LastOrDefault(); var lastNodeId = lastHistory?.StepId; var currentNode = _nodes[lastNodeId];问题在于History里记录的StepId是Guid.NewGuid().ToString("N")——每次执行步骤的唯一ID,它根本不在_nodes字典的键范围内。这段代码一跑就是KeyNotFoundException。正确的做法是每次执行节点时把node.NodeId也存进WorkflowStep,然后从历史里取最后一条的NodeId来反向定位。
第三个雷:CurrentStep的补偿逻辑。
ProcessNodeAsync进来就CurrentStep++,然后MoveToNext里又CurrentStep--。这个补偿设计非常容易被并发分支搞乱:并行分支的每个分支Context是独立深拷贝的,但主流程的CurrentStep只更新一次,多分支同时回来时补偿次数对不上,步骤计数就漂了。后来我直接用History.LastOrDefault().ActivityName来定位当前处于哪个节点,不再依赖一个会漂移的整数计数器。
这些雷的价值在于:它们说明引擎的调度状态必须是"可推导的",而不是"靠变量维护的"。历史列表本身就是最可靠的状态来源,任何单独的计数器都可能在异常分支中被绕过。
5. 性能验证:10万TPS是怎么测出来的,以及"虚高"陷阱
5.1 我真实的压测方法与数据
10万这个数字是怎么来的?不是拍脑袋,是在一台8核16G的测试机上,用Stopwatch计时批量启动10万个不涉及人工审批的自动流程,从入队到所有End节点清理完成,统计总耗时。
先看一个不严谨的压测代码版本——帖子初稿里的Program.Main其实没真正调引擎,创建了10万个Context之后只做了await Task.CompletedTask。压测程序如果这么写,测出来的TPS当然"虚高",因为引擎压根没干活。我后来在真实压测脚本里,启动流程直接调engine.StartWorkflowAsync,消费端用SemaphoreSlim+事件来感知全部完成:
const int count = 100000; var engine = PurchaseApprovalWorkflow.Build(); var sw = Stopwatch.StartNew(); for (int i = 0; i < count; i++) { await engine.StartWorkflowAsync("purchase"); } sw.Stop(); Console.WriteLine($"耗时: {sw.ElapsedMilliseconds} ms"); Console.WriteLine($"吞吐: {count / sw.Elapsed.TotalSeconds:F0} TPS");自动流程走的是Start -> check_amount -> auto_approve -> end,每个节点几乎不耗时,整个链路纯粹考验引擎的调度能力。这个场景下测出来的数据是"引擎的纯调度上限",而不是"真实业务吞吐"。
5.2 TPS虚高的三大来源
必须承认,10万这个数字有很强的展示性质。真拿到生产环境做同样的事,数字至少要打三折,原因有三:
第一,没有持久化。压测版全流程在内存里完成,生产环境每个节点流转都要落库,一次流转至少两次SQL——一次更新自身状态,一次写历史记录。磁盘I/O才是真实瓶颈。
第二,没有真实的人工审批挂起。压测流程里没有Approval节点,全是Action和Condition。真实流程一旦碰到待审批挂起,实例确实不占线程了,但数据库里会积压大量Pending状态的记录,后续的批量唤醒、并发回调都是新的压力点。
第三,没有外部依赖。采购流程里的付款节点要调财务系统,通知节点要发消息,这些外部调用的P99延迟会成为整体吞吐的天花板。
所以我在汇报TPS数字的时候,一定会跟着说清楚压测场景是什么。TPS单看绝对值没有意义,必须加上"在什么场景、什么负载模型下测得的"这个前提。这能避免业务方被一个看起来吓人的数字忽悠。
5.3 从100到10万:演进路径上的关键优化
我们团队实际是从100单/天的人工流程,一步步演进到能扛10万级并发的引擎的。中间几个关键动作值得单独说:
- 把"人找人"变成"引擎找节点"——这是架构上的根本转变,推翻的是手工流转的老系统,而不是加几个并发线程。
- 把审批等待从轮询变成TaskCompletionSource回调——这一步消灭了最消耗线程池的部分。
- 给引擎套上ActionBlock背压和信号量限流——从"能跑"变成"稳定地跑",吞吐不再因为突发流量而雪崩。
- 用内存引擎先跑通,再逐步接持久化和监控——别一开始就追求全功能,先把主链路跑通,在真实负载下找瓶颈。
5.4 关于压测脚本的避坑提醒
如果你准备复现这套压测,有两点提醒。
一是Stopwatch测的是整个循环的耗时,但引擎是异步入队的,StartWorkflowAsync返回不代表流程跑完。要测真实吞吐,应该在End节点执行处埋一个回调计数器,或者用一个TaskCompletionSource集合等待所有实例到达终态。我见过不止一个同事拿"入队耗时"当作"处理耗时"来汇报,这个数字会虚高好几倍。
二是ActionBlock的BoundedCapacity = 100000设定之后,如果生产者的Post速度长期大于消费速度,Post本身会被阻塞。压测时10万个任务快速Post,队列一旦满了,后面的Post会排队,整体耗时里会包含"等待队列腾位置"的时间,这是正常现象。如果你发现压测TPS和预期差很远,先看一眼是不是背压把流量堵在了入口——加个队列深度日志会比瞎调线程数有用得多。
6. 从演示到生产:演示代码和生产代码之间还差什么
6.1 持久化与崩溃恢复:事件溯源思路
内存版的引擎适合Demo和纯展示,生产环境第一件事就是接持久化。我最终的方案是事件溯源风格的落库:不直接存"流程当前状态",而是把每一步的"状态变更事件"追加写入事件表。
原因很简单:直接存状态,崩溃时你只知道"这个流程挂在经理审批这一步",但不知道它是怎么走到这里的、中间经过哪些分支、哪些动作已经执行过。存事件流,恢复时只要把事件按序重放一遍,整个流程的完整状态就重建了。这套思路和银行流水、账本记账的逻辑一脉相承——只追加、不修改、可回放。
落地时的几个字段建议:
| 字段 | 说明 |
|---|---|
| EventId | 事件唯一ID |
| InstanceId | 流程实例ID |
| NodeId | 触发事件的节点 |
| EventType | 枚举:节点进入、审批通过、条件命中、自动动作完成 |
| Payload | JSON序列化的上下文变量快照 |
| CreatedAt | 事件时间 |
崩溃恢复时,从事件表读最后一个快照,跳过已完成的步骤,从断点继续。
6.2 引擎监控指标与告警
跑过一段时间之后,我养成了一个习惯:每个引擎实例都暴露一组最小化指标,进监控系统:
- 队列当前深度(
InputCount):如果长期靠近BoundedCapacity,说明消费能力不足,需要扩容 - 存活实例数(
_instances.Count):持续飙升说明有大量挂起审批,需要关注积压 - 平均节点处理时长:按节点类型分桶,Timer和外部调用节点通常最慢
- 审批挂起时长分布:超过SLA的实例要自动升级提醒
这几个指标不需要什么复杂埋点,从引擎内部直接读就行。生产环境我记得最清楚的一次告警:队列深度从几百飙到三万,查下来是一个外部通知API挂了,Action节点内部没做超时兜底,导致大量流程卡在同一个节点上互相排队。后来给所有Action节点统一加了熔断逻辑——外部调用超过3秒直接跳过并标记异常,队列深度立刻回落。
6.3 动态节点注册与可视化编排
用了状态机加节点定义这套模型之后,最爽的一件事是流程的调整变得很灵活。改一个流程,不用动引擎代码,只要换个节点注册表:
var engine = new WorkflowEngine(); engine.RegisterNode(new WorkflowNode { ... }); engine.RegisterNode(new WorkflowNode { ... });甚至可以在管理后台把节点定义序列化成JSON动态加载。节点类型是固定的,节点之间的连接关系是数据,数据和代码分离之后,业务人员通过拖拽配置流程就成为了可能。这块我在第二期迭代里接了一个简单的流程设计器,节点数据落库,引擎启动时从库里读节点表构造流程定义,运维成本大幅下降。
6.4 最后的实战建议
如果你准备在自己项目里动手写工作流引擎,我的建议是不要一上来就追求功能全。从一个最小闭环开始:Start节点、一个Approval节点、一个End节点,跑通之后再加Condition,再加Parallel,再考虑持久化。每一步都在真实业务里验证,而不是在Demo里验证。框架搭得再漂亮,最后让你改代码的一定是那些真实流程里冒出来的边边角角——比如"驳回之后要从哪一步开始重走""两个审批节点先后顺序能不能和图形画的不一致""超时自动通过之后怎么通知发起人"。
这些问题,没有任何开源框架能替你答完,但自己写引擎的好处就是——你有完全的控制力,可以针对这些业务细节做最顺手的处理。从100个手工审批单到能扛10万并发的引擎,技术上其实没有魔法,核心就是状态机建模加异步调度两层事。但就是这两层事想清楚、做扎实,之后的收益会持续很久。