ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Java Agent流程引擎实战:状态轮转替代if-else与流式输出

Java Agent流程引擎实战:状态轮转替代if-else与流式输出 1. 从if-else泥潭到Agent工作流引擎我为什么重写流程调度最近在做一个Java Agent项目接了意图识别、工具调用、RAG检索、人工确认好几套逻辑一开始图省事所有流程调度全用if-else硬写。等需求改到第三版的时候我终于意识到自己给自己埋了一个巨大的雷方法里叠了六层判断每一层还要处理成功、失败、超时、需要人工介入这些分支加一个新环节就得在关键路径上找半天该插到哪个else里改一个顺序牵一发动全身。这不是代码质量问题是流程建模的方式从一开始就选错了。Agent的核心本质是一个能自主决策、分步执行的程序它的每一次执行都像一条流水线先判断用户意图再调用工具然后生成回复可能中间还要停下来问用户一句确认执行吗。这种场景天然适合用流程引擎来描述而不是用if-else来描述。于是我把原来的调度逻辑全部推翻用一套非常轻量的流程引擎重做了。这套引擎的核心就两件事节点状态轮转和流式输出。节点状态轮转让流程调度变成一张可配置的图节点自己决定我完成后下一步该交给谁引擎只负责按状态推进流式输出则让Agent执行的每一个中间结果、每一个工具调用日志、每一条token增量都能实时推给前端。做完之后最直观的感受是新增一个Agent节点只需要写一个节点类然后注册进去调度逻辑完全不用动代码里再也没出现过一堆if-else嵌套。这篇文章就是把我从0到1实现这套流程引擎的全过程写出来包括核心设计、状态机模型、流式输出方案、完整代码实操以及我踩过的坑。如果你也在用Java搞Agent开发或者想把自己项目里的硬编码流程改成可编排的引擎这篇可以直接当参考。2. 流程引擎整体设计与状态轮转机制2.1 节点模型与状态机的设计想把流程交给引擎调度第一步是定义清楚节点长什么样。在我的设计里一个流程是由多个节点组成的有向图每个节点做一件独立的事节点执行完会返回一个结果结果里包含下一步要去哪个节点以及当前节点的状态。节点本身是一个抽象概念实际业务里它可以是调用大模型生成意图、执行某个工具方法、发送一条MQ消息、等待用户确认等等。我最开始犯过一个错误就是让节点直接持有下一个节点的引用结果整个流程还是硬编码的。后来改成节点通过返回值声明下一步引擎从注册表里按名字找下一个节点这才把节点之间的耦合彻底解开。节点的生命周期状态我用一个枚举来定义public enum NodeStatus { PENDING, // 待执行 RUNNING, // 执行中 SUCCESS, // 执行成功 FAILED, // 执行失败可重试或终止 SKIPPED, // 跳过条件不成立 WAITING // 等待外部输入如人工确认 }这个状态机就是整个流程引擎的底层骨架。节点最初都是PENDING引擎从Start节点开始跑拿到节点后先把状态置为RUNNING执行一段业务代码根据返回结果把状态改成SUCCESS或FAILED再依据结果里的targetNodeId找到下一个节点继续跑。这就是所谓的状态轮转——流程的推进不是靠代码里的条件分支而是靠节点当前状态 返回的下一跳自然流转。2.2 状态轮转如何替代if-else很多人一开始不理解你用状态轮转和用if-else最终不都是要判断下一步走哪个分支吗区别到底在哪我用一个例子说明。假设一个Agent任务要经历意图识别 - 调用工具 - 生成回复 - 人工确认。用传统if-else写代码长这样if (intent 查询天气) { String weather callWeatherTool(city); if (weather ! null) { String answer generateAnswer(weather); if (needManualReview(answer)) { waitingForManual(); } else { send(answer); } } else { handleError(); } } else if (intent 订机票) { // 又是一大坨 }看着好像还行但实际上每个if里面还要加各种异常分支、重试分支、日志埋点一旦扩展到十几个节点这个方法的行数会迅速膨胀到几百上千行。最要命的是流程的形状被写死在代码里你想在工具调用和生成回复之间插一个工具结果格式化节点必须改这段嵌套逻辑。改用流程引擎后同样的流程只需要定义一张节点关系映射MapString, ListTransition graph new HashMap(); graph.put(intent, List.of(new Transition(TOOL, intenttool))); graph.put(TOOL, List.of(new Transition(GENERATE, success))); graph.put(GENERATE, List.of(new Transition(REVIEW, need_review)));引擎拿到当前节点执行完返回的状态和下一跳目标去这张表里查消息走到对应节点。节点之间的关系是数据不是代码。新增一个节点就是往图里加一个条目任何时候想改变流程走向只需要改配置或改节点返回的targetNodeId。if-else判断的是当前这一刻怎么走状态轮转描述的是整张图怎么走后者才能驾驭复杂的Agent编排。2.3 引擎核心调度器实现调度器是整个流程引擎的心脏我把它实现成一个非常精简的循环。核心逻辑就三件事拿到当前节点、执行、根据结果决定下一个节点。看代码public class WorkflowEngine { private final MapString, Node nodeRegistry new ConcurrentHashMap(); private final WorkflowContext context new WorkflowContext(); public void registerNode(Node node) { nodeRegistry.put(node.getName(), node); } public void run(String startNodeName) { String currentNodeName startNodeName; while (currentNodeName ! null !Thread.currentThread().isInterrupted()) { Node currentNode nodeRegistry.get(currentNodeName); if (currentNode null) { throw new IllegalStateException(Node not found: currentNodeName); } currentNode.updateStatus(NodeStatus.RUNNING); NodeResult result currentNode.execute(context); currentNode.updateStatus(result.getStatus()); // 根据状态决定是否重试、跳过或结束 if (result.getStatus() NodeStatus.FAILED result.isAllowRetry()) { currentNodeName currentNodeName; // 原地重试 continue; } if (result.getStatus() NodeStatus.SKIPPED) { currentNodeName result.getTargetNodeId(); continue; } if (result.getStatus() NodeStatus.WAITING) { // 挂起等待外部信号这里简化处理为结束 break; } currentNodeName result.getTargetNodeId(); } } }这个循环写起来简单难点在于考虑完整。比如人工确认节点它执行完不会立即给出结果而是要挂起流程等用户在外部点击确认后再触发引擎继续跑。所以我给NodeStatus加了WAITING遇到WAITING就先把当前节点挂起引擎主循环退出后续通过一个resumeWorkflow方法从挂起节点继续推进。3. 流式输出让Agent的每个动作都实时可见3.1 为什么Agent场景必须用流式输出Agent和普通后端接口有个巨大的差异执行时间长且过程信息有展示价值。一次完整的Agent调用可能涉及模型推理、多个工具调用、多轮内部思考如果像普通接口一样等到全部执行完再一次性返回结果用户在前端就是白屏等待十几秒甚至几分钟体验极其糟糕。我之前试过轮询方案Agent每执行一步就往数据库里写一条状态前端每隔两秒查一次。能用但太蠢第一轮询有延迟第二数据库被高频写入第三执行过程的一些半结构化日志比如大模型流式吐出来的token很难通过轮询还原成自然的打字机效果。后来我把方案换成了服务器推送事件SSE也就是基于HTTP长连接的服务端流式推送。SSE实现简单、兼容性好、Java生态支持完善特别适合Agent这种服务和前端是单向数据流的场景。流式输出解决了两个问题一是让用户看到Agent正在做什么而不是干等二是让中间产物工具调用参数、检索到的资料、推理片段能被前端实时渲染比如在界面上展示当前正在调用天气工具参数北京这样的过程卡片。3.2 基于Spring WebFlux的流式推送实现我项目里用的Spring Boot 3直接依赖了WebFlux的响应式栈。这里我不会把整个响应式编程讲一遍只说怎么在一个普通接口里实现SSE流式输出。核心思路是用一个广播源作为事件总线。Agent引擎里的任何节点都能往这个广播源里发消息HTTP响应通过Flux订阅这个广播源把消息实时转发给前端。import org.springframework.http.codec.ServerSentEvent; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; import reactor.core.publisher.Flux; import reactor.core.publisher.Sinks; import java.util.concurrent.ConcurrentHashMap; RestController public class AgentFlowController { // 广播源支持多订阅者 private final Sinks.ManyServerSentEventString sink Sinks.many().multicast().onBackpressureBuffer(); private final WorkflowEngine engine; public AgentFlowController(WorkflowEngine engine) { this.engine engine; } GetMapping(/execute) public FluxServerSentEventString executeAgent() { // 这里把当前会话的sink实例传给引擎让引擎在节点执行时向这个sink发送事件 engine.setEventSink(sink); // 异步启动引擎 engine.runAsync(START); return sink.asFlux(); } }在Node执行过程中只需要调用sink.tryEmitNext()就能推送一条事件。比如在执行工具调用节点时就会发一条类似tool_start的事件sink.tryEmitNext(ServerSentEvent.builder(开始调用天气工具城市北京) .event(tool_start) .build());这样前端就能实时收到文本片段配合CSS就能实现打字机效果。注意WebFlux的Flux是无限流所以需要确保当Agent流程结束后要主动向sink发一个flow_end事件然后调用sink.tryEmitComplete()关闭流否则前端会一直等待。3.3 节点状态与流式事件如何联动流程引擎有了流式输出能力之后我做了个很顺手的扩展把节点状态变化也作为事件推出去。也就是说前端不仅能看到Agent输出了什么还能看到一个实时的流程状态面板类似意图识别 - 识别完成 - 工具调用中 - 工具调用成功 - 生成回复中的可视化进度条。我定义了一个FlowEvent类型包含事件类型、节点名、节点状态、数据载荷public class FlowEvent { private String type; // NODE_STATUS, TOKEN, TOOL_CALL, ERROR, FLOW_END private String nodeName; private NodeStatus status; private Object data; // 构造器、getter省略 }节点执行前、执行中、执行后分别emit不同的事件。在节点基类里统一做这件事public abstract class Node { protected Sinks.ManyServerSentEventString sink; public void setEventSink(Sinks.ManyServerSentEventString sink) { this.sink sink; } public final NodeResult execute(WorkflowContext context) { emitStatus(NodeStatus.RUNNING); try { NodeResult result doExecute(context); emitStatus(result.getStatus()); return result; } catch (Exception e) { emitError(e.getMessage()); return NodeResult.failure(ERROR, e.getMessage()); } } private void emitStatus(NodeStatus status) { if (sink ! null) { sink.tryEmitNext(ServerSentEvent.builder() .event(node_status) .comment(nodeName : status) .build()); } } protected abstract NodeResult doExecute(WorkflowContext context); }这里有个细节ServerSentEvent的comment字段不会显示给前端但可以用来传递元数据前端通过解析event和data来更新状态面板。节点状态和数据流通过两个不同的事件类型推送前端各做各的渲染互不干扰。4. 从0到1的实操过程一个可用的Agent流程引擎4.1 环境准备与项目结构实际操作环境建议组件版本说明JDK17用到了record、sealed等新特性Spring Boot3.2WebFlux支持内置SSEMaven3.8依赖管理Reactor3.6WebFlux自带项目结构很清爽src/main/java/com/example/agentflow ├── engine │ ├── WorkflowEngine.java │ ├── WorkflowContext.java │ ├── NodeResult.java │ ├── NodeStatus.java │ └── Node.java ├── node │ ├── StartNode.java │ ├── IntentNode.java │ ├── ToolCallNode.java │ ├── GenerateNode.java │ └── EndNode.java ├── event │ ├── FlowEvent.java │ └── AgentFlowController.java └── Application.java这是一个极其简化的结构实际项目里我还会把NodeResult和WorkflowContext做更丰富的扩展比如WorkflowContext里放会话ID、用户ID、traceId方便追踪整条链路。4.2 节点与流程定义我用一个查询天气并生成回复的Agent流程来演示。流程是START - INTENT - TOOL_CALL - GENERATE - END。如果意图不是查询天气TOOL_CALL节点会被跳过直接走向GENERATE。先定义各个节点的核心逻辑。StartNode流程入口只有一个作用就是决定第一个实际节点public class StartNode extends Node { Override protected NodeResult doExecute(WorkflowContext context) { return NodeResult.success(INTENT); } }IntentNode调用大模型(这里简化成硬编码判断)识别意图public class IntentNode extends Node { Override protected NodeResult doExecute(WorkflowContext context) { String userInput context.get(userInput, String.class); if (userInput.contains(天气)) { context.put(intent, weather); return NodeResult.success(TOOL_CALL); } context.put(intent, chat); return NodeResult.success(GENERATE); } }ToolCallNode调用天气工具public class ToolCallNode extends Node { Override protected NodeResult doExecute(WorkflowContext context) { String city context.get(city, String.class); // 这里模拟工具调用实际项目中会换成HTTP调用或函数方法 String weather 晴25度; context.put(weather, weather); // 推送工具调用事件 if (sink ! null) { sink.tryEmitNext(ServerSentEvent.builder(调用天气工具完成 weather) .event(tool_result) .build()); } return NodeResult.success(GENERATE); } }这里注意所有节点都必须通过NodeResult指定下一步节点名这是状态轮转的关键约束。如果一个节点没有下一步了就返回null或者用一个END节点兜底。GenerateNode生成回复public class GenerateNode extends Node { Override protected NodeResult doExecute(WorkflowContext context) { String intent context.get(intent, String.class); if (weather.equals(intent)) { String weather context.get(weather, String.class); String answer 当前城市天气 weather; context.put(answer, answer); // 模拟流式输出逐字推送 for (char c : answer.toCharArray()) { if (sink ! null) { sink.tryEmitNext(ServerSentEvent.builder(String.valueOf(c)) .event(token) .build()); } try { Thread.sleep(50); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } } else { context.put(answer, 我现在只能查询天气哦); } return NodeResult.success(END); } }EndNode什么都不做返回null让引擎结束public class EndNode extends Node { Override protected NodeResult doExecute(WorkflowContext context) { return NodeResult.end(); } }4.3 核心代码实现WorkflowContext就是一个简单的Map封装但为了防止在并发节点执行时数据错乱我用的是ConcurrentHashMappublic class WorkflowContext { private final MapString, Object data new ConcurrentHashMap(); public void put(String key, Object value) { data.put(key, value); } public T T get(String key, ClassT type) { return type.cast(data.get(key)); } }NodeResult负责承载节点执行结果。我给它加了几个静态工厂方法让节点代码更简洁public class NodeResult { private final String targetNodeId; private final NodeStatus status; private final boolean allowRetry; private NodeResult(String targetNodeId, NodeStatus status, boolean allowRetry) { this.targetNodeId targetNodeId; this.status status; this.allowRetry allowRetry; } public static NodeResult success(String targetNodeId) { return new NodeResult(targetNodeId, NodeStatus.SUCCESS, false); } public static NodeResult failure(String targetNodeId, boolean allowRetry) { return new NodeResult(targetNodeId, NodeStatus.FAILED, allowRetry); } public static NodeResult end() { return new NodeResult(null, NodeStatus.SUCCESS, false); } public static NodeResult skip(String targetNodeId) { return new NodeResult(targetNodeId, NodeStatus.SKIPPED, false); } public static NodeResult waitForExternal(String targetNodeId) { return new NodeResult(targetNodeId, NodeStatus.WAITING, false); } // getter省略 }接着把上面的引擎主循环补全加入事件sink和异步执行能力。异步执行我用一个简单的线程池封装public class WorkflowEngine { private final MapString, Node nodeRegistry new ConcurrentHashMap(); private final ExecutorService executor Executors.newCachedThreadPool(); private Sinks.ManyServerSentEventString sink; public void setEventSink(Sinks.ManyServerSentEventString sink) { this.sink sink; nodeRegistry.values().forEach(node - node.setEventSink(sink)); } public void registerNode(Node node) { node.setEventSink(sink); nodeRegistry.put(node.getName(), node); } public void runAsync(String startNodeName) { executor.submit(() - run(startNodeName)); } public void run(String startNodeName) { String currentNodeName startNodeName; int maxStep 100; int step 0; while (currentNodeName ! null) { if (step maxStep) { throw new IllegalStateException(流程疑似死循环超过最大步数); } Node currentNode nodeRegistry.get(currentNodeName); if (currentNode null) { throw new IllegalStateException(未找到节点: currentNodeName); } NodeResult result currentNode.execute(createContext()); NodeStatus status result.getStatus(); if (status NodeStatus.FAILED result.isAllowRetry()) { continue; } if (status NodeStatus.WAITING) { // 流程挂起等待外部resume break; } currentNodeName result.getTargetNodeId(); } if (sink ! null) { sink.tryEmitNext(ServerSentEvent.builder(流程结束).event(flow_end).build()); sink.tryEmitComplete(); } } }4.4 运行效果演示启动应用后用curl模拟前端发起SSE请求curl -N -H Accept: text/event-stream http://localhost:8080/execute?userInput北京天气怎么样响应流会依次出现event: node_status data: START : RUNNING event: node_status data: INTENT : RUNNING event: token data: 当 event: token data: 前 event: token data: 城 event: token data: 市 event: node_status data: TOOL_CALL : SUCCESS event: tool_result data: 调用天气工具完成晴25度 event: flow_end data: 流程结束前端可以根据事件类型做不同的UI渲染node_status更新状态面板token追加到正文tool_result展示工具卡片。这套联动让Agent执行过程可视化程度非常高实际体验很接近ChatGPT那种思考中...然后开始输出答案的交互。5. 常见问题与排查技巧实录5.1 节点重复执行问题我在实现初期犯过一个非常经典的错在run循环里当一个节点执行完并返回SUCCESS后如果目标节点ID正好还是它自己比如配错了图就会导致同一个节点无限循环执行。后来我在引擎里加了最大步数兜底但这只是防止死循环并没有真正解决重复执行的隐患。真正需要重视的是节点幂等性。Agent流程里经常有发送通知写入日志这类节点如果在流程恢复或重试时被执行两次就会产生重复副作用。我的解决方案是在节点基类提供isIdempotent()方法默认返回false对于非幂等节点引擎收到外部resume请求时不会从当前节点重新执行而是从它的下一个节点开始。另外节点执行前在WorkflowContext里记录一个executionId同一executionId的节点即使被再次触发也会被跳过用执行标记保证单次流程内每个节点最多执行一次。5.2 SSE流式输出中断和乱序SSE本身是基于HTTP长连接如果中间有Nginx代理默认可能有缓冲和超时问题。我本地调试好好的一上测试环境就发现前端收不到消息最后查出来是Nginx没关缓冲proxy_buffering off; proxy_cache off; proxy_read_timeout 3600s;这一点是SSE部署的经典坑遇到本地能用走代理不行的问题优先检查这行配置。乱序问题是我在并发节点时候踩的坑。虽然我现在演示的流程是串行的但Agent实际运行中有些节点可以并行比如同时调用两个工具。如果两个并行节点同时往sink里emit事件前端拿到的顺序可能跟期望不一致。我的做法是给FlowEvent加一个递增的sequence字段前端收到事件后按sequence重排。简单有效不用引入复杂的分布式消息。5.3 引擎如何扩展新的节点类型很多朋友会问新需求来了我要加一个发送邮件节点怎么办非常简单写一个类继承Node实现doExecute方法在doExecute里指定它的下一步节点通常是某个后续节点名在启动配置里注册这个节点。整个过程不需要改动WorkflowEngine和任何既有节点。这就是状态轮转带来的核心收益——新增节点是加法操作不是修改操作。如果某个新的Agent流程需要不同的编排我甚至可以直接用配置文件定义节点顺序引擎启动时读取配置构建注册表。这个我已经在做了流程定义完全外置业务上写流程编排就像搭乐高。5.4 并发场景下的引擎安全Agent服务不可能只有一个用户所以引擎必须支持多流程并发。我刚开始把WorkflowContext定义成单例多个流程同时跑的时候数据互相污染。后来改成每次流程启动时创建独立的Context并把Context作为NodeResult执行时的参数传递进去。具体做法是在注册节点时只存节点实例但节点内部不能持有任何用户态的成员变量。所有数据都放进Context。如果一个节点是有状态的比如需要缓存某次工具调用的结果就把缓存数据放到Context里而不是节点的字段里。这个经验对我来说是血泪教训最初那个节点里塞了个List存放临时结果结果两个用户并发请求时数据串得一塌糊涂排查了一下午。6. 写在最后一点实战体会整套流程引擎做下来最大的体会是不用if-else不是目的流程可编排才是目的。if-else是图灵完备的什么流程都能写但问题在于它不可观测、不可复用、不可动态调整。状态轮转把流程的形状从代码中抽离出来让每个节点只关心自己的事这是Agent这种复杂场景下更优雅的建模方式。最后分享一个小技巧我引擎里加了一个调试模式开关开启后每个节点执行前后的状态变化、目标节点、耗时都会通过流式事件推给前端调试面板。这个功能在联调时帮了大忙——团队成员肉眼就能看到流程走到哪一步、为什么走到这一步排查问题的速度翻倍。如果你也在设计Agent工作流建议先从最精简的引擎开始只做到节点注册和状态轮转再逐步加上流式输出、并行节点、人工确认、重试策略。这套架构的扩展方向非常多但核心模型一旦稳固后面加功能都是水到渠成的事。
返回列表