
做后端开发的同学大概都有过这种经历一个业务审批流、一个数据清洗管道、一个多步骤的AI Agent任务最开始用if-else硬编码三层嵌套还能忍等到步骤加到七八个、还要支持中途暂停、失败重试、状态查询的时候代码就彻底变成了一团乱麻。我接手过一个类似的项目前任开发者用switch-case写了一个六阶段的文档处理流程光是状态判断就写了四百多行加一个新节点要改五个地方测试覆盖率死活上不去。后来我花了两天时间把它重构成一套轻量级的流程引擎核心代码不到三百行新增节点只需要实现一个接口、注册一下就行。这篇文章就把这套从0到1的实现思路完整拆开讲一遍包括节点状态怎么轮转、流式输出怎么接、为什么这样设计而不是那样设计。适合有Java基础、正在做Agent开发或者工作流相关需求的同学参考也适合想理解流程引擎底层原理的朋友。1. 为什么if-else写工作流迟早要崩1.1 硬编码流程的四个典型症状先说说我遇到的那个真实场景。需求是一个文档处理Agent上传文档、解析内容、提取关键信息、调用大模型生成摘要、格式化输出、存档。听起来六个步骤挺清晰的对吧前任的写法是在一个service方法里顺序调用六个private方法每个方法返回一个boolean表示成功失败主方法里用if-else串联。第一个问题是状态不可见。用户想知道现在处理到哪一步了代码里根本没有状态记录只能靠日志grep。第二个问题是无法中断恢复。处理到第四步服务重启了前面三步的成果全丢只能从头再来。第三个问题是分支逻辑爆炸。后来产品要求如果文档是PDF走A路径如果是Word走B路径如果解析失败要重试两次再走降级路径if-else的嵌套层级直接到了五层。第四个问题是无法复用。另一个业务线要做类似的流程只能复制粘贴再改改完两边就不同步了。这四个症状的本质是同一个问题流程的控制逻辑和业务逻辑耦合在了一起。if-else既是流程调度器又是业务执行器职责不清自然越写越乱。1.2 流程引擎要解决的核心命题一个合格的流程引擎本质上要回答四个问题。第一当前在哪个节点这需要状态管理。第二下一个节点是谁这需要路由决策。第三节点执行的结果怎么传递这需要上下文Context机制。第四执行过程怎么让外部感知这需要事件通知或流式输出。把这四个问题抽象出来流程引擎的骨架就清晰了一个节点接口定义做什么一个上下文对象承载数据一个执行器负责调度一个监听器负责通知。剩下的都是在这四个骨架上做扩展。我选择不引入Activiti、Flowable这类重型工作流引擎原因很直接它们是为BPMN规范设计的面向的是人工审批、表单流转这类场景配置复杂、学习成本高而且和AI Agent这种代码驱动、动态决策的场景并不契合。轻量级自研反而更灵活几百行代码就能覆盖绝大多数Agent场景。1.3 节点状态轮转整个引擎的心跳状态轮转是流程引擎最核心的机制。每个节点在执行过程中会经历一系列状态变迁理解这个变迁过程就理解了引擎的运转逻辑。我定义的状态枚举包含六个值PENDING待执行、RUNNING执行中、SUCCESS成功、FAILED失败、SKIPPED跳过、RETRYING重试中。一个节点从PENDING开始被调度器选中后变为RUNNING执行成功转SUCCESS抛异常转FAILED如果配置了重试策略则先转RETRYING再回到RUNNING。这里有个容易忽略的细节状态变迁必须是原子的。如果节点执行成功但状态更新失败或者状态更新了但结果没保存就会出现数据不一致。我的做法是把状态更新和结果保存放在同一个同步块里虽然牺牲了一点并发性能但保证了正确性。对于Agent场景来说流程执行本身就不是高并发操作这个取舍是值得的。状态轮转的另一个价值是可观测性。每个节点的状态变化都通过监听器广播出去前端可以实时展示当前处理到第3步正在调用大模型用户体验直接上了一个台阶。2. 节点抽象与上下文设计的关键取舍2.1 Node接口应该定义几个方法节点接口的设计直接决定了整个引擎的易用性。我见过有的实现把接口设计得特别复杂五六个方法要全部实现写一个简单节点要几十行样板代码。我的原则是接口方法数量控制在两个以内其余用默认方法提供。最终的核心接口是这样的public interface Node { String getName(); NodeResult execute(FlowContext context) throws Exception; default ListString nextNodes(FlowContext context) { return Collections.emptyList(); } default RetryPolicy retryPolicy() { return RetryPolicy.noRetry(); } }getName()用于标识节点execute()是核心执行逻辑nextNodes()决定后续路由默认返回空表示流程结束retryPolicy()定义重试策略。后两个都是默认方法简单节点只需要实现前两个。为什么把路由决策放在节点里而不是引擎里因为路由逻辑往往依赖节点的执行结果。比如解析文档节点执行后要根据文档类型决定走哪个分支这个判断逻辑放在节点内部最自然引擎不需要知道业务细节。这就是所谓的节点自治原则。2.2 FlowContext数据传递的载体上下文对象是节点之间传递数据的唯一通道。我见过两种极端设计一种是把上下文做成一个Map什么都能塞灵活但类型不安全另一种是每个流程定义一个强类型的Context类安全但复用性差。我的方案是折中FlowContext内部持有一个MapString, Object作为通用存储同时提供泛型getter方法public class FlowContext { private final MapString, Object data new ConcurrentHashMap(); private final String flowId; private volatile String currentNode; public T T get(String key, ClassT type) { Object value data.get(key); return value null ? null : type.cast(value); } public void put(String key, Object value) { data.put(key, value); } public T T getOrDefault(String key, ClassT type, T defaultValue) { T value get(key, type); return value ! null ? value : defaultValue; } }用ConcurrentHashMap是因为流式输出场景下可能有多个线程往上下文里写数据比如异步的LLM调用回调。currentNode用volatile修饰保证状态变更对其他线程立即可见。注意上下文里不要存大对象。我踩过一个坑把整个文档的字节数组塞进上下文结果流程跑一百个并发就OOM了。正确做法是存引用或路径需要时再加载。2.3 节点结果的标准化封装节点执行结果需要标准化否则引擎没法统一处理。我定义了NodeResult类包含三个字段success是否成功、message描述信息、output输出数据可选。public class NodeResult { private final boolean success; private final String message; private final MapString, Object output; public static NodeResult success(String message) { return new NodeResult(true, message, Collections.emptyMap()); } public static NodeResult success(String message, MapString, Object output) { return new NodeResult(true, message, output); } public static NodeResult failure(String message) { return new NodeResult(false, message, Collections.emptyMap()); } }output字段的设计有个讲究它会被自动合并到FlowContext里。这样节点之间传递数据就不需要手动put了节点A返回的output节点B直接从context里get就行。这个小小的自动化省掉了大量样板代码。3. 执行器与状态轮转的完整实现3.1 主循环从入口节点到终止执行器是引擎的心脏负责驱动整个流程。核心逻辑是一个while循环从当前节点开始执行、判断结果、决定下一个节点、更新状态直到没有下一个节点或遇到终止条件。public class FlowExecutor { private final MapString, Node nodeRegistry new ConcurrentHashMap(); private final ListFlowListener listeners new CopyOnWriteArrayList(); private final int maxSteps; public FlowExecutor(int maxSteps) { this.maxSteps maxSteps; } public void register(Node node) { nodeRegistry.put(node.getName(), node); } public FlowContext execute(String startNode, FlowContext context) { String current startNode; int stepCount 0; while (current ! null stepCount maxSteps) { stepCount; Node node nodeRegistry.get(current); if (node null) { throw new IllegalStateException(节点未注册: current); } context.setCurrentNode(current); notifyListeners(FlowEvent.nodeStarted(current)); NodeResult result executeWithRetry(node, context); if (!result.isSuccess()) { notifyListeners(FlowEvent.nodeFailed(current, result.getMessage())); context.setFailed(true); break; } context.mergeOutput(result.getOutput()); notifyListeners(FlowEvent.nodeCompleted(current, result)); ListString nextNodes node.nextNodes(context); current nextNodes.isEmpty() ? null : nextNodes.get(0); } if (stepCount maxSteps) { throw new IllegalStateException(流程执行步数超过上限可能存在死循环); } return context; } }maxSteps这个参数非常重要它是死循环的保险丝。Agent场景下节点路由可能依赖LLM的输出如果LLM返回了意料之外的节点名或者路由逻辑有bug流程可能无限循环。设置一个上限比如100步超过就抛异常能避免服务被拖死。3.2 重试机制指数退避与最大次数重试是Agent场景的刚需因为LLM调用、网络请求都可能偶发失败。我的重试策略支持三个参数最大重试次数、初始退避时间、退避倍数。public class RetryPolicy { private final int maxRetries; private final long initialBackoffMs; private final double backoffMultiplier; public static RetryPolicy noRetry() { return new RetryPolicy(0, 0, 1.0); } public static RetryPolicy exponential(int maxRetries, long initialMs) { return new RetryPolicy(maxRetries, initialMs, 2.0); } }执行时的重试逻辑private NodeResult executeWithRetry(Node node, FlowContext context) { RetryPolicy policy node.retryPolicy(); int attempt 0; long backoff policy.getInitialBackoffMs(); while (true) { try { NodeResult result node.execute(context); if (result.isSuccess() || attempt policy.getMaxRetries()) { return result; } } catch (Exception e) { if (attempt policy.getMaxRetries()) { return NodeResult.failure(执行异常: e.getMessage()); } } attempt; context.setNodeState(node.getName(), NodeState.RETRYING); notifyListeners(FlowEvent.nodeRetrying(node.getName(), attempt)); try { Thread.sleep(backoff); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); return NodeResult.failure(重试被中断); } backoff (long) (backoff * policy.getBackoffMultiplier()); } }指数退避的意义在于如果是服务端限流导致的失败立即重试只会加重限流等待时间翻倍增长能给服务端喘息的机会。实测下来初始500ms、倍数2.0、最多3次重试能覆盖90%以上的偶发失败场景。3.3 状态持久化让流程可以断点续跑状态轮转如果只存在内存里服务一重启就全丢了。要做到断点续跑需要把状态持久化。我的方案是定义一个StateStore接口提供save和load方法默认实现是内存版生产环境可以换成Redis或数据库版。public interface StateStore { void save(FlowContext context); FlowContext load(String flowId); }持久化的时机很关键。太频繁每个节点状态变化都存会影响性能太少又可能丢状态。我的做法是在节点执行完成后持久化一次因为节点内部的状态变化对恢复来说意义不大恢复时从节点边界重新执行即可。这就要求节点执行是幂等的或者至少是可重入的。提示如果你的节点有副作用比如发邮件、扣款一定要做幂等设计。我一般会在上下文里记录一个executedNodes集合节点执行前先检查是否已执行过。4. 流式输出让Agent的思考过程可见4.1 为什么流式输出对Agent至关重要传统的工作流是黑盒提交任务等结果。但Agent场景不一样用户需要看到Agent的思考过程。调用大模型生成摘要如果等全部生成完再返回用户要盯着loading转十几秒如果流式输出用户能看到文字一个个蹦出来体验完全不同。流式输出的技术本质是生产者-消费者模型。节点是生产者不断产生数据块输出通道是消费者把数据块推送给前端。中间需要一个缓冲区来解耦两者速度差异。4.2 基于回调的流式接口设计我设计了一个StreamCallback接口节点在执行过程中通过它推送数据public interface StreamCallback { void onChunk(String chunk); void onComplete(); void onError(Throwable error); }节点执行时如果需要流式输出就从上下文里取出callbackpublic class LlmSummaryNode implements Node { Override public NodeResult execute(FlowContext context) { StreamCallback callback context.get(streamCallback, StreamCallback.class); String content context.get(documentContent, String.class); StringBuilder fullResponse new StringBuilder(); llmClient.streamGenerate(buildPrompt(content), chunk - { fullResponse.append(chunk); if (callback ! null) { callback.onChunk(chunk); } }); context.put(summary, fullResponse.toString()); return NodeResult.success(摘要生成完成); } }这个设计的好处是节点不需要关心输出到哪里。callback可能是推送到SSE连接、可能是写入文件、可能是发到消息队列节点只管调用onChunk具体实现由外部注入。这就是依赖倒置原则的实际应用。4.3 SSE推送与背压处理实际项目中流式输出最常见的落地方式是SSEServer-Sent Events。Spring Boot里用SseEmitter就能实现。但这里有个坑如果生产速度大于消费速度数据会堆积在内存里。我遇到过一次LLM生成速度很快但前端网络慢SseEmitter的缓冲区越积越大最后OOM。解决方案是加背压控制给callback加一个带容量限制的阻塞队列队列满了就阻塞生产者。public class BackpressureStreamCallback implements StreamCallback { private final BlockingQueueString queue; private final SseEmitter emitter; public BackpressureStreamCallback(SseEmitter emitter, int bufferSize) { this.emitter emitter; this.queue new ArrayBlockingQueue(bufferSize); startConsumer(); } Override public void onChunk(String chunk) { try { queue.put(chunk); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } private void startConsumer() { new Thread(() - { try { String chunk; while ((chunk queue.take()) ! null) { if (__END__.equals(chunk)) break; emitter.send(chunk); } emitter.complete(); } catch (Exception e) { emitter.completeWithError(e); } }).start(); } }队列容量我一般设100太小容易阻塞影响生成速度太大起不到背压作用。这个值需要根据实际的前端消费速度和LLM生成速度来调。4.4 流式输出与状态轮转的协同流式输出和状态轮转需要协同工作。一个节点在RUNNING状态下开始流式输出输出完成后转SUCCESS。如果输出中途出错转FAILED并通知callback的onError。这里有个细节流式输出的内容要不要存进上下文我的做法是存但只存最终完整结果不存每个chunk。因为chunk数量可能成千上万全存下来内存扛不住。节点内部用StringBuilder累积最后一次性put进上下文。5. 从if-else迁移到流程引擎的实操路径5.1 识别可抽取的节点边界迁移不是推倒重来而是渐进式重构。第一步是识别现有代码里的节点边界。判断标准很简单一个方法如果满足输入明确、输出明确、副作用可控就可以抽成一个节点。拿我那个文档处理流程举例原来的六个private方法天然就是六个节点。但有些方法内部还有分支比如解析文档方法里根据文件类型走不同逻辑这时候有两种选择抽成一个节点内部做分支或者抽成多个节点由路由决定。我的经验是如果分支后的逻辑差异很大超过20行就拆成多个节点如果只是参数不同就留在一个节点里。5.2 上下文数据的迁移策略原来的if-else写法数据通常通过方法参数和返回值传递。迁移到流程引擎后所有数据都要走FlowContext。这里有个迁移技巧先做一层适配把原来的参数打包进context方法签名改成接收context。比如原来是parseDocument(byte[] bytes, String type)迁移后变成execute(FlowContext context)内部从context取bytes和type。这样改动最小每个方法独立迁移迁完一个测一个风险可控。5.3 灰度切换与回滚方案生产环境迁移一定要有灰度方案。我的做法是加一个开关通过配置决定走老逻辑还是新引擎。新引擎先跑影子模式同样的输入两边都执行对比结果但不影响主流程。跑一周没问题后再切流量。回滚方案也要准备好。因为流程引擎的状态是持久化的回滚时要注意新引擎产生的状态数据老逻辑不认识。所以回滚前要确保没有正在执行中的流程或者做一个状态转换层。5.4 迁移后的收益量化迁移完成后我统计了几个指标。代码行数从原来的1200行降到400行引擎300行节点100行。新增一个节点的改动量从改5个地方降到实现一个类注册一行。测试覆盖率从45%提升到82%因为每个节点可以独立单测。最直观的是排查问题的时间以前用户报卡住了要翻日志找现在直接查状态表一眼看到卡在哪个节点、什么状态。6. 几个容易踩的坑与排查思路6.1 节点注册顺序导致的空指针第一个坑是节点注册顺序。如果节点A的nextNodes返回了节点B但B还没注册执行时就会抛节点未注册。这个问题在启动阶段就能发现我的做法是在引擎初始化完成后加一个校验遍历所有节点的nextNodes检查返回的节点名是否都已注册。public void validate() { for (Node node : nodeRegistry.values()) { for (String next : node.nextNodes(new FlowContext(validate))) { if (!nodeRegistry.containsKey(next)) { throw new IllegalStateException( 节点 node.getName() 指向未注册的节点: next); } } } }注意这里传了一个空的context给nextNodes所以nextNodes的实现不能依赖context里的数据否则校验时会NPE。如果确实需要依赖就把校验改成运行时校验。6.2 上下文并发修改的隐蔽bug第二个坑是并发修改。流式输出场景下LLM的回调线程可能和主执行线程同时操作context。我遇到过一次ConcurrentModificationException原因是遍历context的keySet时另一个线程put了新key。解决方案是遍历时用快照。new HashMap(context.getData())复制一份再遍历。或者用ConcurrentHashMap的弱一致性迭代器它不会抛异常但可能看不到最新的修改。具体用哪种取决于业务能否接受看不到最新修改。6.3 重试导致的重复副作用第三个坑最隐蔽重试导致的重复执行。节点执行到一半失败了重试时从头执行如果前半段有副作用比如已经发了一条消息就会重复。我的解决方案是两阶段提交节点内部把操作分成准备和提交两步准备阶段无副作用提交阶段才产生副作用。重试时只重跑准备阶段提交阶段用幂等键保证只执行一次。这个模式稍微复杂但对于有副作用的节点是必须的。6.4 流式输出中断的处理第四个坑是流式输出中断。用户关闭了浏览器SSE连接断了但后端的LLM还在生成。如果不处理就是白白浪费token。我的做法是给callback加一个isCancelled()检查节点在每次onChunk前检查一下如果已取消就抛异常终止。同时注册SSE的onCompletion和onTimeout回调触发取消信号。emitter.onCompletion(() - callback.cancel()); emitter.onTimeout(() - callback.cancel());这个细节看起来小但能省下不少API调用费用尤其是用按token计费的大模型时。7. 引擎的扩展方向与个人实践体会7.1 并行节点与条件分支的扩展当前实现是单线执行的但实际场景经常需要并行。比如同时调用三个模型生成摘要取最好的那个。扩展思路是把nextNodes的返回值从List改成支持并行组的结构执行器用CompletableFuture并发执行全部完成后合并结果再继续。条件分支的扩展更简单在nextNodes里根据context的数据返回不同的节点名即可。比如if (context.get(docType).equals(PDF)) return List.of(pdfNode); else return List.of(wordNode);。这比if-else清晰得多因为路由逻辑被隔离在节点内部不会污染主流程。7.2 与Spring生态的集成方式生产项目里节点通常需要依赖Spring容器里的Bean比如Service、Mapper。集成方式有两种一种是把节点也注册成Spring Bean引擎从容器里取另一种是节点内部通过ApplicationContext手动获取。我推荐第一种因为可以利用Spring的依赖注入和AOP。具体做法是给Node接口加Component注解引擎注入MapString, NodeSpring会自动把所有Node实现类注入进来key是bean名称。这样连手动注册都省了。Component public class FlowExecutor { Autowired private MapString, Node nodeRegistry; }7.3 监控指标的埋点建议流程引擎的监控非常重要。我一般会埋这几类指标每个节点的执行次数、成功率、平均耗时、P99耗时整个流程的执行次数、成功率、平均步数重试次数、失败原因分布。这些指标用Micrometer暴露成Prometheus格式配合Grafana看板能快速定位问题。比如某个节点成功率突然下降或者P99耗时飙升都是明显的异常信号。7.4 我个人的几点实践体会最后分享几点踩坑踩出来的体会。第一引擎要简单。我见过有人把流程引擎做成了支持BPMN、支持可视化编排、支持热更新的庞然大物结果维护成本比业务代码还高。轻量级引擎的价值就在于简单可控几百行代码谁都能看懂、能改。第二节点要无状态。节点的所有状态都应该在context里节点实例本身不持有可变状态。这样节点可以被多线程安全地复用也方便做单元测试。第三日志要打全。每个节点的开始、结束、异常都要打日志带上flowId和nodeName。排查问题时一个flowId就能串起整个流程的所有日志效率提升非常明显。第四别过度设计。一开始不要想着支持所有场景先把最核心的顺序执行状态轮转流式输出做扎实等真正遇到并行、分支、子流程的需求时再扩展。我第一版引擎只有200行后来根据实际需求逐步加到了300多行每一步扩展都有明确的业务驱动没有一行是凭空想象的。这套引擎后来在我们团队内部推广开了三个业务线都在用累计跑了上百万次流程。最让我欣慰的是新来的同事看半天代码就能上手写节点再也不用去啃那四百行的switch-case了。如果你也在被if-else写的工作流折磨不妨试试这个思路从抽出一个最简单的节点开始慢慢把流程引擎搭起来。