
1. 为什么我决定把 Agent 工作流从 if-else 里彻底抽出来先说结论如果你现在写 Agent 还在用一长串if (step 1) {...} else if (step 2) {...}来驱动那这套代码撑不过三个业务迭代。我接手过一个内部工具项目最初就是一个智能问答 工具调用的小 Agent逻辑简单到用 switch 就能写完。结果两个月后需求变成了多轮对话、工具调用失败重试、人工审核节点、流式输出、上下文压缩、节点超时降级。原来的 switch 膨胀到 800 多行改一个分支要通读全文测试根本覆盖不过来。这不是代码风格问题是架构问题。if-else 的本质是把流程控制权交给代码顺序而 Agent 的本质是把流程控制权交给状态。这两者天生冲突。Agent 的每一步执行结果都是不确定的——工具可能超时、模型可能返回空、用户可能中途打断——用顺序代码去描述一个状态机等于用直线去画圆。所以这篇文章要讲的事情很明确用 Java 从 0 到 1 实现一个 Agent 工作流的流程引擎核心是三件事——节点状态轮转、流式输出、以及把流程定义从代码里剥离出来。关键词里的 Java、Agent、工作流、流程引擎、流式输出我会一个一个拆开讲透。适合谁看写过 Java 后端、对 Agent 有基本概念、但被流程控制折磨过的同学。不需要你懂什么高深框架Spring Boot 基础够用就行。我先把整体思路摆出来避免你读到一半不知道在干嘛。整个引擎分四层流程定义层节点、边、条件、状态机层节点状态轮转、执行层节点执行器 上下文、输出层流式事件推送。这四层里状态机是心脏流式输出是血管流程定义是骨架。下面逐层拆。提示本文所有代码都是可运行的最小实现思路不是伪代码。你可以直接照着搭一个 demo跑通之后再往生产环境加东西。2. 流程引擎的骨架节点、边与上下文到底怎么设计2.1 为什么不用现成的工作流引擎你可能会问Java 生态里不是有 Activiti、Flowable 这些成熟的工作流引擎吗我试过。结论是它们是为确定性业务流程设计的不是为 Agent 设计的。Activiti 的节点执行是同步阻塞的一个 UserTask 要等人审批而 Agent 的节点执行是异步的、可能流式返回的、可能被模型动态决定的。你硬套 Activiti会发现它的 BPMN 模型根本表达不了根据模型输出动态选择下一个节点这种语义。轻量级工作流才是正解。所谓轻量级就是不引入 BPMN 那套 XML 规范用纯 Java 对象描述流程。一个流程就是一组节点加一组转移规则节点执行完返回一个结果引擎根据结果决定下一步去哪。这个模型简单到你可以用 200 行代码实现核心但它能覆盖 90% 的 Agent 场景。2.2 三个核心抽象Node、Edge、Context先定义数据结构。我用的是最朴素的 POJO没有注解没有继承树因为 Agent 流程需要的是灵活不是规范。public interface Node { String getId(); NodeResult execute(WorkflowContext ctx) throws NodeException; } public class NodeResult { private final String nextNodeId; // 下一个节点null 表示结束 private final boolean stream; // 是否流式输出 private final Object payload; // 节点产出 // 构造、getter 省略 } public class WorkflowContext { private final MapString, Object variables new ConcurrentHashMap(); private final ListStreamEvent events new CopyOnWriteArrayList(); // put/get/appendEvent 省略 }Node接口只有一个execute方法返回NodeResult。这个设计的关键在于节点不关心自己下一步去哪它只负责执行并给出建议。真正的路由决策交给引擎。这样节点之间彻底解耦你可以随意替换、复用、单测。WorkflowContext是贯穿整个流程的上下文所有节点共享。这里有个坑我踩过不要用 ThreadLocal 存上下文。Agent 执行经常跨线程异步工具调用、流式回调ThreadLocal 会丢数据。用显式的 Context 对象传递虽然啰嗦但可靠。2.3 边的表达条件路由而不是硬编码节点之间的连接我用一个Edge来描述public class Edge { private final String from; private final String to; private final PredicateWorkflowContext condition; private final int priority; }condition是个谓词引擎在节点执行完后遍历所有从当前节点出发的边按 priority 排序选第一个条件为真的边跳过去。如果所有条件都不满足就走默认边priority 最低、condition 恒真。这套设计的好处是新增一条分支只需要加一个 Edge不用改任何节点代码。比如工具调用失败要重试这个需求就是加一条fromtoolNode, totoolNode, conditionctx - ctx.get(lastError) ! null retryCount 3的边。原来的 if-else 里你得在 toolNode 内部写重试逻辑节点职责就不纯了。2.4 流程定义的 DSL 化光有对象还不够流程定义最好能配置化。我做了个简单的 BuilderWorkflow workflow Workflow.builder() .node(new LlmNode(llm)) .node(new ToolNode(tool)) .node(new AnswerNode(answer)) .edge(llm, tool, ctx - ctx.get(needTool) Boolean.TRUE) .edge(llm, answer, ctx - ctx.get(needTool) ! Boolean.TRUE) .edge(tool, llm) // 工具执行完回到 llm .start(llm) .build();这段代码读起来就是流程本身比 800 行 switch 直观太多。而且这个 Builder 可以进一步从 JSON 或数据库加载做到流程热更新——改流程不用重新发版。这是 if-else 永远做不到的。注意Builder 里的边顺序会影响匹配优先级我建议显式指定 priority别依赖插入顺序否则重构时容易出诡异 bug。3. 节点状态轮转把下一步去哪交给状态机3.1 状态轮转的本质是一个循环流程引擎的主循环其实非常短短到你可能不信public WorkflowResult run(Workflow workflow, WorkflowContext ctx) { String current workflow.getStartNodeId(); int step 0; while (current ! null) { if (step MAX_STEPS) { throw new WorkflowException(超过最大步数疑似死循环); } Node node workflow.getNode(current); ctx.setCurrentNode(current); NodeResult result node.execute(ctx); ctx.appendEvent(StreamEvent.nodeDone(current, result)); current route(workflow, current, ctx, result); } return WorkflowResult.from(ctx); }核心就三行取当前节点、执行、路由到下一个。所有复杂度都被封装在node.execute和route里了。这就是状态机的威力——主循环稳定不变变化的是节点和路由规则。MAX_STEPS这个保护必须有。Agent 流程里节点互相跳转很容易写出 A→B→A 的死循环。我设的默认值是 50超过就抛异常并打印完整路径方便定位。3.2 节点状态不是运行中/完成这么简单很多人以为节点状态就是待执行、执行中、已完成。实际做下来Agent 节点至少需要这些状态状态含义触发条件PENDING等待执行节点被路由选中但还没开始RUNNING执行中execute 方法已调用STREAMING流式输出中节点正在推送 tokenWAITING等待外部输入需要人工审核或用户补充SUCCEEDED成功完成正常返回FAILED执行失败抛异常或返回错误SKIPPED被跳过条件不满足为什么要区分 STREAMING 和 RUNNING因为流式输出时节点还没结束但已经产出了部分内容。前端需要知道这个节点正在吐字而不是这个节点卡住了。这个区分对用户体验影响很大。3.3 状态轮转的驱动方式事件驱动而非轮询状态怎么从 RUNNING 变成 SUCCEEDED有两种做法轮询和事件驱动。轮询就是主线程不断查节点状态简单但浪费 CPU而且流式场景下延迟高。我用的是事件驱动节点执行过程中通过 Context 推送事件引擎监听事件并更新状态。public interface StreamEvent { String getNodeId(); EventType getType(); // NODE_START, TOKEN, NODE_END, ERROR Object getData(); long getTimestamp(); }节点在流式输出时每产出一个 token 就ctx.appendEvent(StreamEvent.token(nodeId, token))。引擎侧有个EventDispatcher负责把这些事件推给订阅者比如 SSE 连接。这样状态更新和输出推送是同一套机制不用维护两套逻辑。3.4 状态持久化让流程可以暂停和恢复Agent 流程经常需要暂停——比如等用户确认、等外部系统回调。这时候状态必须能持久化。我的做法是每次节点执行完把 Context 序列化存库。恢复时反序列化从currentNodeId继续跑。public class WorkflowSnapshot { private String workflowId; private String currentNodeId; private MapString, Object variables; private ListStreamEvent events; private long version; // 乐观锁 }这里有个经验Context 里的变量必须可序列化。我一开始往 Context 里塞了个HttpServletRequest结果持久化直接报错。后来定了个规矩Context 只存基本类型、String、List、Map 和自定义的可序列化 POJO。工具句柄、连接池这类东西用的时候现取别往 Context 里放。提示快照的 version 字段是给乐观锁用的。多个请求同时恢复同一个流程时version 不匹配就拒绝避免状态覆盖。4. 流式输出从节点到前端的完整链路4.1 为什么流式输出是 Agent 的刚需先说个反直觉的事实流式输出不只是为了看起来快它是 Agent 可用性的前提。一个 LLM 节点生成 500 字回答非流式要等 8 秒才返回用户早就以为卡死了。流式输出让首字延迟降到 500ms 以内用户能立刻看到反馈。但流式输出的难点不在推而在和流程引擎的配合。节点在流式输出时流程还没结束下一个节点还没确定。这时候如果用户刷新页面状态怎么恢复如果流式中途出错已经推出去的内容怎么办这些问题不解决流式输出就是个玩具。4.2 用 SSE 做传输层但别绑死传输层我选的是 SSEServer-Sent Events因为它是 HTTP 原生的、浏览器支持好、实现简单。但引擎内部不依赖 SSE而是依赖一个抽象的EventSinkpublic interface EventSink { void emit(StreamEvent event); void complete(); void error(Throwable t); }SSE 只是EventSink的一个实现。这样你换成 WebSocket、gRPC stream、甚至写文件引擎代码一行不用改。关键词里提到的流式输出内容到文件就是实现一个FileEventSink而已。public class SseEventSink implements EventSink { private final SseEmitter emitter; public void emit(StreamEvent event) { try { emitter.send(SseEmitter.event() .name(event.getType().name()) .data(event, MediaType.APPLICATION_JSON)); } catch (IOException e) { // 客户端断开标记取消 throw new ClientDisconnectedException(); } } }4.3 背压处理别让快节点撑爆慢客户端流式输出有个隐蔽的坑生产速度远大于消费速度。LLM 节点每秒能吐几十个 token如果客户端网络慢事件会在内存里堆积最后 OOM。这就是背压问题。我的处理方式是有界队列 丢弃策略。EventSink内部维护一个容量 1000 的队列满了之后对于 TOKEN 类型事件直接丢弃反正后面还有对于 NODE_START、NODE_END 这种关键事件阻塞等待。这样既不会 OOM也不会丢关键状态。private final BlockingQueueStreamEvent queue new ArrayBlockingQueue(1000); public void emit(StreamEvent event) { if (event.getType() EventType.TOKEN) { queue.offer(event); // 满了就丢不阻塞 } else { queue.put(event); // 关键事件必须送达 } }实测下来这个策略在弱网环境下很稳。用户看到的是文字偶尔跳一下而不是页面卡死然后报错。4.4 流式输出的顺序保证还有个容易忽略的点事件顺序。多线程环境下节点 A 的 token 和节点 B 的 token 可能交错。前端拿到乱序的 token拼出来的句子就是乱的。解决办法是给每个事件打上nodeId sequence前端按 sequence 排序后再渲染。引擎侧保证同一个节点的事件是顺序产生的节点内部单线程跨节点的事件用nodeStartTime排序。这样即使传输乱序前端也能还原正确顺序。注意别在节点内部开多线程并发 emit token那样顺序无法保证。如果确实需要并发比如同时调多个工具让每个工具的输出走独立的事件流用 nodeId 区分。5. 把 if-else 彻底赶出去节点执行器的设计模式5.1 每个节点只做一件事if-else 最大的问题是一个代码块干了太多事。比如一个处理用户输入的分支里面可能同时做了意图识别、参数提取、工具选择、结果格式化。这四件事耦合在一起改一个影响全部。流程引擎的思路是拆成四个节点IntentNode、ExtractNode、SelectToolNode、FormatNode。每个节点只做一件事通过 Context 传递中间结果。这样每个节点都能单独测试、单独替换。public class IntentNode implements Node { private final LlmClient llm; public NodeResult execute(WorkflowContext ctx) { String input ctx.get(userInput, String.class); String intent llm.classify(input); ctx.put(intent, intent); return NodeResult.next(extract); // 固定路由 } }注意NodeResult.next(extract)是固定路由因为意图识别完必然要提取参数。而 SelectToolNode 就是动态路由了它根据 intent 决定去哪个工具节点。5.2 用策略模式替代条件分支节点内部如果还有 if-else用策略模式消掉。比如 FormatNode 要根据不同工具的输出格式做不同处理public interface Formatter { boolean supports(String toolType); String format(Object raw); } public class FormatNode implements Node { private final ListFormatter formatters; // Spring 注入 public NodeResult execute(WorkflowContext ctx) { String toolType ctx.get(toolType, String.class); Object raw ctx.get(toolResult); String formatted formatters.stream() .filter(f - f.supports(toolType)) .findFirst() .orElseThrow(() - new NodeException(无匹配格式化器)) .format(raw); ctx.put(formatted, formatted); return NodeResult.next(answer); } }新增一种工具类型只需要加一个 Formatter 实现FormatNode 一行不改。这就是开闭原则在 Agent 场景的落地。5.3 节点注册与依赖注入节点多了之后管理是个问题。我用 Spring 的ApplicationContext做节点注册每个节点是一个Component启动时扫描所有Node实现按getId()注册到NodeRegistry。Component public class NodeRegistry { private final MapString, Node nodes new HashMap(); Autowired public NodeRegistry(ListNode nodeList) { nodeList.forEach(n - nodes.put(n.getId(), n)); } public Node get(String id) { Node n nodes.get(id); if (n null) throw new WorkflowException(节点未注册: id); return n; } }这样节点之间的依赖比如 LlmClient、ToolClient都由 Spring 管理你只管写业务逻辑。实测下来一个中等复杂度的 Agent 流程大概 15-20 个节点用这套注册机制管理起来很清爽。5.4 节点执行器的横切关注点节点执行前后有些通用逻辑日志、耗时统计、异常包装、重试。这些不该写在每个节点里。我用装饰器模式包一层public class RetryNodeDecorator implements Node { private final Node delegate; private final int maxRetry; public NodeResult execute(WorkflowContext ctx) { int attempt 0; while (true) { try { long start System.currentTimeMillis(); NodeResult r delegate.execute(ctx); log.info(节点 {} 耗时 {}ms, delegate.getId(), System.currentTimeMillis() - start); return r; } catch (Exception e) { if (attempt maxRetry) throw e; log.warn(节点 {} 第 {} 次重试, delegate.getId(), attempt); } } } }注册的时候决定哪些节点需要重试装饰。这样重试逻辑和业务逻辑彻底分离节点代码保持干净。提示不是所有节点都适合重试。LLM 节点重试可能产生重复计费工具节点重试可能产生副作用。我一般只对幂等的查询类节点加重试装饰。6. 实测中的坑状态丢失、流式乱序与死循环6.1 状态丢失异步回调里的 Context 陷阱第一个大坑是异步工具调用导致 Context 丢失。我有个工具节点调外部 API用了CompletableFuture异步执行回调里想往 Context 写结果结果发现 Context 是空的。原因是回调线程拿到的 Context 引用和主线程不是同一个我早期用了 ThreadLocal。修复方案前面提过Context 显式传递不用 ThreadLocal。异步回调里通过闭包捕获 Context 引用CompletableFuture.supplyAsync(() - tool.call(params)) .thenAccept(result - ctx.put(toolResult, result));但这里还有个隐患ctx.put用的是ConcurrentHashMap多线程写是安全的但写顺序不确定。如果两个异步任务都写同一个 key后写的覆盖先写的。我的做法是给每个异步任务分配独立的 key 前缀最后在汇聚节点统一合并。6.2 流式乱序一次真实的排查过程有次线上反馈回答的句子是乱的。我按这个链路排查先看前端渲染逻辑——按接收顺序拼接没问题。再看 SSE 传输——抓包发现事件顺序确实乱了。定位到引擎侧——EventSink用了线程池推送多线程导致乱序。根因节点内部为了提高吞吐开了线程池并发 emit token。修复很简单token 的 emit 必须单线程。我把节点内的并发 emit 改成串行吞吐确实降了一点但顺序对了。后来优化方案是并发只在计算阶段emit 阶段串行化。// 错误做法并发 emit tokens.parallelStream().forEach(t - sink.emit(token(t))); // 正确做法并发计算串行 emit ListString results tokens.parallelStream().map(this::compute).toList(); results.forEach(t - sink.emit(token(t)));6.3 死循环A→B→A 的隐蔽陷阱前面提过MAX_STEPS保护但那只防无限循环防不了有限但错误的循环。我遇到过一个 caseLLM 节点判断信息不足路由到追问节点追问节点拿到用户补充后又回到 LLM 节点LLM 还是判断信息不足……循环 5 次后触发 MAX_STEPS。排查这种问题光看代码没用得打印完整路径。我在引擎里加了个path记录ListString path new ArrayList(); while (current ! null) { path.add(current); // ... } // 异常时打印 path有了路径一眼就能看出是哪个节点在反复跳。后来发现根因是追问节点没把用户补充写进 ContextLLM 每次看到的都是旧信息。6.4 排查清单遇到问题先看这几项现象优先排查常见根因状态丢失Context 传递方式用了 ThreadLocal流式乱序emit 是否单线程节点内并发 emit死循环打印完整 path节点未更新 Context内存暴涨事件队列容量无背压事件堆积恢复失败快照可序列化Context 含不可序列化对象节点不执行路由条件条件谓词恒假这张表是我踩坑踩出来的建议你直接抄进项目 wiki。7. 从 Demo 到生产还需要补哪些能力7.1 可观测性没有 trace 的流程引擎是黑盒Demo 跑通不代表能用。生产环境第一件事是加 trace。我给每个流程实例分配一个traceId每个节点执行打一条结构化日志包含 traceId、nodeId、耗时、状态。这样出问题能快速定位到具体节点。MDC.put(traceId, ctx.getTraceId()); MDC.put(nodeId, node.getId()); log.info(node execute start); // ... MDC.clear();配合日志采集系统可以做出流程执行瀑布图一眼看出哪个节点慢。7.2 超时控制每个节点都要有超时Agent 节点调外部服务必须设超时。我的做法是在RetryNodeDecorator里加超时控制FutureNodeResult future executor.submit(() - delegate.execute(ctx)); try { return future.get(node.getTimeoutMs(), TimeUnit.MILLISECONDS); } catch (TimeoutException e) { future.cancel(true); throw new NodeTimeoutException(node.getId()); }超时时间按节点类型配置LLM 节点 30 秒工具节点 10 秒格式化节点 1 秒。别用统一超时否则要么 LLM 被误杀要么工具节点拖死整个流程。7.3 并发控制同一流程实例不能并发执行一个流程实例比如一个会话同时只能有一个执行线程。否则两个请求同时推进同一个流程状态会乱。我用分布式锁Redis 实现保证String lockKey workflow:lock: workflowId; if (!lock.tryLock(lockKey, 5, TimeUnit.SECONDS)) { throw new WorkflowException(流程正在执行中); } try { // 执行流程 } finally { lock.unlock(lockKey); }锁的粒度是 workflowId不是全局锁这样不同会话可以并行。7.4 版本管理流程定义变更怎么办流程定义会变但正在执行的流程不能受影响。我的做法是流程定义带版本号流程实例启动时锁定版本执行过程中用锁定版本的定义。新请求用新版本。这样变更平滑不会出现执行到一半流程定义变了的问题。public class WorkflowDefinition { private String id; private int version; private MapString, Node nodes; private ListEdge edges; }数据库里存多版本WorkflowRegistry按 id version 取。这个设计在流程频繁迭代的场景下特别重要。7.5 一个最小可用的生产配置最后给个我实际用的配置参考workflow: max-steps: 50 default-timeout-ms: 30000 event-queue-capacity: 1000 snapshot: enabled: true store: redis ttl-hours: 24 lock: enabled: true timeout-seconds: 5这套配置支撑过日均十万级的流程执行稳定性没问题。当然具体数值要根据你的业务调比如 max-steps 如果流程确实很长可以放宽到 100但要配合路径监控。8. 我在这套引擎上的一些个人体会写到这里核心的东西都讲完了。最后分享几个我自己的体会不算总结就是踩坑之后的真实感受。第一流程引擎的价值不在引擎本身而在约束。它强制你把流程拆成节点强制你显式定义路由强制你处理状态。这些约束一开始觉得麻烦但正是它们让代码在需求膨胀时还能保持可控。if-else 的自由是假自由它把复杂度藏起来了等到爆发时你连从哪下手都不知道。第二流式输出要早做。我第一个版本没做流式后来加的时候发现要改节点接口、改 Context、改传输层牵一发动全身。如果一开始就把EventSink抽象出来后面加流式就是实现一个接口的事。第三别过度设计。我见过有人一上来就搞分布式流程引擎、搞 BPMN、搞可视化编排结果核心业务逻辑还没跑通。先用 200 行把单机版跑通验证流程模型对不对再考虑扩展。这套引擎我从 200 行起步迭代了半年才到现在这个规模每一步都是被真实需求推着走的。如果你现在手上正好有个 if-else 堆成山的 Agent 项目我的建议是别急着全量重构先挑一个最复杂的流程用这套思路重写一遍跑通之后你自然知道该怎么推广。代码这东西看一百篇文章不如自己写一遍。