
好久没写工程向的分享了。前面一个项目里我负责做一个智能客服Agent一开始图省事所有流程控制全用if-else堆如果意图不明就追问如果知识库没命中就兜底如果LLM调用失败就重试……两周之后代码已经没法看了加一个节点要动三四处地方测试用例写了一堆还是漏。后来我花了一个周末把流程控制层整个重写成基于状态轮转的流程引擎——Java实现节点抽象状态显式定义流转用路由表驱动执行过程用SSE流式输出到前端。核心代码大概不到300行if-else基本被消灭了。这篇文章就围绕这个设计展开讲讲节点状态轮转和流式输出到底怎么落地顺便把踩过的坑一并整理出来适合正在做Agent开发、AI应用编排、或者被业务分支逻辑折磨的Java工程师参考。1. 为什么Agent工作流不能靠if-else硬写1.1 当Agent流程开始复杂if-else的账就算不清了Agent的工作流和普通后端接口最大的不同在于它的执行路径几乎不受你控制。一个典型Agent要经历意图识别、槽位补齐、知识库检索、工具调用、上下文记忆更新、答案生成这几个阶段而每个阶段都可能产生分支意图置信度低要追问信息不足要澄清检索结果为空要走兜底工具执行超时要重试。把这些分支全用if-else写本质上就是在用“顺序代码”去描述一张“有向图”。问题在前端接口还能忍一旦节点数超过五个代码就会变成这个样子外层if判断结果类型内层if判断状态码再内层if判断是否需要调用另一个方法。每新增一个节点你要找到所有上游节点的出口逻辑一处处补判断。更难受的是调试——线上用户触发了某条冷门分支你根本不知道当时走到了哪一步因为没有状态记录。我见过不少团队在Agent项目里维护一段超过500行的switch-case每加一个工具调用就往里塞一个case分支。这本质上和if-else没有区别只是换了个写法。流程控制一旦退化成这种形式它就不再是“可编排”的了而是“写死”的。1.2 状态机思想从路牌到流程图如果退一步看问题Agent的执行过程本质上就是一个状态机每个节点是一个状态节点的执行结果决定下一个状态节点在流转中会产生输出事件。用状态机的思路去建模比用代码堆逻辑要自然得多。我打个比方你开车去一个陌生的地方靠的不是在每个路口背诵“如果看到A就左转看到B就右转否则直行”这种if-else口诀而是一张路线图每个路口告诉你当前在哪个位置根据当前位置决定下一段怎么走。状态轮转就是这张路线图节点就是路口的指示牌执行结果就是你要做的选择。流程引擎的核心价值在这里体现出来了它把“业务逻辑”和“流程控制”彻底拆开。业务逻辑写在节点内部流程控制交给引擎的路由机制。以后哪怕要调整整个Agent的执行顺序你也不需要改业务代码在路由表里重新编排节点ID就可以。2. 流程引擎的整体设计先抽象再实现2.1 核心抽象Node节点和Workflow工作流设计一个流程引擎最重要的不是先写代码而是定义好抽象。我这边核心抽象只有两个Node节点和Workflow工作流。Node是流程中的最小执行单元。它负责做三件事声明自己的ID、执行具体的业务逻辑、明确执行完成之后下一步去哪个节点。业务逻辑五花八门没关系我把它统一收敛到一个抽象方法里返回执行状态。状态决定了路由方向。Workflow则是节点的集合和入口配置。它知道整个流程从哪个节点开始手里有一张注册表存着所有节点的ID到实例的映射。引擎执行的时候只需要拿到入口节点ID然后不停查询当前节点、执行当前节点、根据返回状态找下一个节点循环往复直到走到终止状态。这里有一个很重要的设计取舍节点之间互相不直接依赖。A节点执行完后它不需要知道B节点是什么只需要在路由表里声明“状态SUCCESS时去node_detail_query这个ID”。这种解耦带来的直接好处是节点可以独立测试可以随意替换也可以在不同工作流里复用。2.2 状态轮转如何替代if-else传统写法里控制权在“调用方”手里。一个方法调用另一个方法根据返回值决定下一步层层嵌套。状态轮转的设计里控制权被上交给“引擎”手里。每个节点执行完成后会返回一个NodeStatus引擎拿到这个状态去查当前节点的路由表拿到目标节点ID接着执行下一个节点。也就是说分支逻辑从“代码的调用关系”变成了“数据表里的映射关系”。if-else被彻底从代码里剥离你看到的是清晰的路由表SUCCESS → 下一步做什么FAILED → 是重试还是走兜底NEED_INPUT → 是否要跳转到追问节点有人会说这不就是把if-else搬了个位置吗其实区别很大。if-else是硬编码流程一旦写错或者要调整必须改代码重新发布路由表是结构化数据可以配置化、可视化、动态修改。而且对于流程引擎来说路由表天然适合做遍历检查——你可以启动时扫描所有节点的路由目标检测有没有指向不存在的节点ID有没有死循环。这些能力是普通if-else写法根本给不了的。2.3 为什么不用Activiti、Flowable这些现成引擎肯定有人要问Java生态里不是有Activiti、Flowable、Camunda这些成熟的工作流引擎吗为什么还要自己写我自己评估过这些引擎强在“人参与审批”的场景任务分配、会签、或签、超时提醒。但Agent场景完全是另一回事节点执行的是LLM推理调用、工具API调用、向量检索不是人工填表审批。引入Activiti这种重型BPM框架光是把流程定义文件和数据模型接进来就要花不少功夫再加上它们本身那套部署方式、历史表结构对一个小团队来说负担很重。更关键的一点是流式输出。Agent执行过程需要实时把中间状态推到前端——用户能看着“正在理解意图→正在检索知识库→正在生成回答”这是现代Agent应用的基本体验要求。Activiti这类引擎的监听机制是为系统内部事件设计的要接到SSE推送还得自己搭一层桥梁。与其绕一圈适配不如直接自己写一个轻量引擎保留最核心的节点路由能力把事件推送机制做成一等公民。实测下来核心引擎不到300行后面所有业务节点都在复用这一套骨架。3. 核心实现节点状态轮转引擎代码拆解3.1 状态枚举与节点抽象类先看最基本的定义状态枚举我一开始只设计了5个状态后来加了TIMEOUT和TERMINATED别嫌多实际写业务的时候你会发现每个状态都有对应场景。public enum NodeStatus { PENDING(待执行), RUNNING(执行中), SUCCESS(成功), FAILED(失败), SKIPPED(跳过), TIMEOUT(超时), TERMINATED(终止); private final String desc; NodeStatus(String desc) { this.desc desc; } public String getDesc() { return desc; } }然后是节点的抽象基类。核心就两个东西一个execute方法用于执行业务并返回状态一个路由表用于指明不同状态分别去哪个节点。public abstract class BaseNodeT { private final String id; private final String name; private final MapNodeStatus, String routeTable new EnumMap(NodeStatus.class); public BaseNode(String id, String name) { this.id id; this.name name; } public String getId() { return id; } public String getName() { return name; } public abstract NodeStatus execute(NodeContextT context); protected void on(NodeStatus status, String targetNodeId) { routeTable.put(status, targetNodeId); } public String route(NodeStatus status) { return routeTable.get(status); } }设计要点在route方法执行完一个节点引擎把返回状态交给路由表查询找到的就是下一站节点ID。如果某个状态没有配置路由说明这个状态就是终态引擎会终止循环。这里我建议一个实际优化在做验证时遍历注册表里的所有节点检查每个节点的路由目标是否存在于注册表中。这个检查我写在启动阶段一旦发现悬空引用直接报错比运行时走到一半才发现要好得多。3.2 执行上下文让所有节点共享数据和状态节点之间光靠参数传递是不够的Agent执行过程中会产生大量上下文数据比如用户的原始问题、意图识别的结果、知识库检索出来的片段、历史对话记录这些数据需要在不同节点之间流转。为此我设计了一个NodeContextpublic class NodeContextT { private final String traceId; private final T payload; private final MapString, Object variables new ConcurrentHashMap(); private String currentNodeId; private NodeStatus currentNodeStatus; public NodeContext(String traceId, T payload) { this.traceId traceId; this.payload payload; } public void setVariable(String key, Object value) { variables.put(key, value); } SuppressWarnings(unchecked) public V V getVariable(String key) { return (V) variables.get(key); } public String getTraceId() { return traceId; } public T getPayload() { return payload; } public String getCurrentNodeId() { return currentNodeId; } public void markNode(String nodeId, NodeStatus status) { this.currentNodeId nodeId; this.currentNodeStatus status; } }traceId是每次Agent请求的唯一标识贯穿整个执行链路既用于日志追踪也用于流式输出时关联事件。variables是节点间的数据交换层A节点把意图识别的结果放进去B节点直接取不需要方法签名级别的耦合。实际项目中这个上下文还可以扩展出上下文窗口管理、令牌消耗统计、重试次数记录等能力。它本质上是Agent的“工作记忆”记得在哪一步停过、拿过什么结果、下一步还缺什么信息。3.3 引擎主循环用路由表驱动节点跳转引擎是整个流程的中枢它负责调度节点、处理状态轮转、广播事件。我把它设计成可复用的通用组件public class FlowEngineT { private final MapString, BaseNodeT registry new HashMap(); private final ExecutorService executorService; private final ListWorkflowListener listeners new CopyOnWriteArrayList(); public FlowEngine(ExecutorService executorService) { this.executorService executorService; } public void registerNode(BaseNodeT node) { registry.put(node.getId(), node); } public void addListener(WorkflowListener listener) { listeners.add(listener); } public void start(String entryNodeId, T payload, String traceId) { executorService.submit(() - execute(entryNodeId, payload, traceId)); } private void execute(String entryNodeId, T payload, String traceId) { NodeContextT context new NodeContext(traceId, payload); String currentNodeId entryNodeId; try { while (currentNodeId ! null) { BaseNodeT node registry.get(currentNodeId); if (node null) { throw new IllegalStateException(节点不存在: currentNodeId); } context.markNode(currentNodeId, NodeStatus.RUNNING); emitEvent(NodeEvent.started(context, System.currentTimeMillis())); NodeStatus status node.execute(context); context.markNode(currentNodeId, status); emitEvent(NodeEvent.finished(context, System.currentTimeMillis())); currentNodeId node.route(status); } emitEvent(NodeEvent.completed(context, System.currentTimeMillis())); } catch (Exception e) { context.setVariable(lastError, e.getMessage()); emitEvent(NodeEvent.failed(context, e, System.currentTimeMillis())); } } private void emitEvent(NodeEvent event) { for (WorkflowListener listener : listeners) { listener.onEvent(event); } } }主循环逻辑很简洁根据当前节点ID拿到节点实例执行取状态查路由得到下一个节点ID。没有判断流程逻辑的if-else所有分支都收敛在节点的路由表里。这里我想强调一个容易忽略的点开始执行用的start方法是异步的把任务丢进线程池立即返回。为什么因为工作流要配合流式输出如果不异步化HTTP响应线程会被阻塞前端的SSE流根本没机会实时收到事件。异步化之后整个执行变成事件驱动的推送模型。4. 流式输出让Agent执行过程实时可见4.1 全量返回为什么会让用户焦虑Agent场景和普通接口最大的体验差异在于普通接口用户等一两秒拿到结果没问题但Agent执行一个复杂任务可能耗时几十秒甚至几分钟——比如多轮工具调用再加长文本生成。如果用户点完按钮之后页面一直转圈没有任何过程反馈体验是灾难性的。流式输出解决的就是这个问题。它的本质是把Agent执行过程中产生的每一个阶段性事件从服务端实时推送到客户端。注意是事件流不只是最终答案。用户能看到“正在判断意图”→“正在检索知识库”→“正在调用外部工具”→“正在生成回答”每一条状态变化都实时渲染到页面上。技术选型上我用了SSEServer-Sent Events而不是WebSocket。原因有两点Agent执行状态是服务端到客户端的单向数据流不需要客户端频繁上行消息SSE基于HTTP协议不需要额外的握手协议前端的EventSource对象可以直接消费实现成本低得多。4.2 基于SseEmitter的事件推送实现Spring Boot生态里SseEmitter是原生的SSE支持用起来很直接。我把它和流程引擎的事件监听器接在一起执行过程中每产生一个事件就通过emitter推送给前端RestController public class AgentController { private final FlowEngineAgentRequest flowEngine; private final ExecutorService executorService Executors.newCachedThreadPool(); public AgentController(FlowEngineAgentRequest flowEngine) { this.flowEngine flowEngine; } GetMapping(value /agent/run, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter run(RequestParam String question) { String traceId UUID.randomUUID().toString().substring(0, 8); SseEmitter emitter new SseEmitter(60_000L); WorkflowListener listener event - { try { MapString, Object data new HashMap(); data.put(traceId, event.getTraceId()); data.put(nodeId, event.getContext().getCurrentNodeId()); data.put(status, event.getContext().getCurrentNodeStatus().name()); data.put(type, event.getType().name()); if (event.getTargetNodeId() ! null) { data.put(next, event.getTargetNodeId()); } emitter.send(SseEmitter.event().name(agent-event).data(data)); } catch (IOException e) { emitter.completeWithError(e); } }; flowEngine.addListener(listener); emitter.onCompletion(() - flowEngine.removeListener(listener)); emitter.onTimeout(() - flowEngine.removeListener(listener)); flowEngine.start(ENTRY_NODE_ID, new AgentRequest(question), traceId); return emitter; } }这里有一个非常关键的设计SseEmitter推送给前端的不只是最终结果而是执行过程中的每一次状态轮转。节点开始执行推一个started事件节点结束推一个finished事件事件里带上当前节点ID和状态。前端拿到这些事件之后就可以在页面上逐个渲染状态变化。用户看到的是“流程在往前走”而不是一个死等loading。数据事件里还带了next字段如果有这个字段说明流程还要继续如果没有说明这个节点是终态。4.3 线程池、事件监听器与断连处理流式输出本身不复杂复杂在工程细节。第一个细节是线程隔离。Agent工作流的执行必须放在独立的线程池里不能占着Tomcat的请求线程。如果直接在请求线程里跑工作流SseEmitter就永远来不及发送数据因为响应还没返回。我用了独立的线程池请求线程只负责创建emitter和注册监听器真正的工作流执行在另一个线程里异步跑。第二个细节是监听器的生命周期管理。每个SSE连接对应一个监听器实例连接关闭或者超时之后监听器必须从引擎里移除。否则emitter已经关闭了工作流还继续往里面send数据会抛IOException。我在onCompletion和onTimeout回调里都调用了removeListener确保连接断开后不再有推送动作。第三个细节是超时设置。SseEmitter的构造参数是超时毫秒数我设置了60秒。如果Agent在这个时间内还没跑完连接会断掉用户需要重新发起请求。后续我改成了更合理的方案用定时器在超时前刷新一下连接或者把超时设成一个比较长的值配合心跳机制保持连接活性。对于生产环境建议单独评估每个Agent流程的最大耗时再决定超时时长。5. 实操案例从0搭建一个智能客服Agent5.1 节点设计与路由编排理论讲了一堆最终要落到一个能跑通的例子。我以智能客服Agent为例设计了5个业务节点intent_node意图识别节点。调用LLM判断用户意图返回四个状态BUY想买产品、AFTER_SALE售后咨询、OTHER闲聊、UNKNOWN无法识别。clarify_node追问节点。当意图识别返回UNKNOWN时追问用户具体需求然后重新路由回意图识别节点。knowledge_node知识库检索节点。根据意图和用户问题做向量检索从知识库中匹配相关内容。如果检索到的内容相似度低于阈值返回FAILED。fallback_node兜底节点。知识库检索不到内容时给出一个预设话术引导用户转接人工。answer_node答案生成节点。把检索结果和历史上下文拼进Prompt调用LLM生成最终回答。节点定义好了路由表就是整个Agent的“剧本”。我专门用一个配置类把路由关系集中管理Configuration public class AgentWorkflowConfig { public static final String ENTRY_NODE intent_node; Bean public FlowEngineAgentRequest agentFlowEngine() { FlowEngineAgentRequest engine new FlowEngine(Executors.newFixedThreadPool(8)); BaseNodeAgentRequest intentNode new BaseNodeAgentRequest(intent_node, 意图识别) { Override public NodeStatus execute(NodeContextAgentRequest context) { IntentResult intent llmService.recognizeIntent(context.getPayload().getQuestion()); context.setVariable(intent, intent); if (intent.confidence() 0.6) { return NodeStatus.NEED_INPUT; } return NodeStatus.SUCCESS; } }; intentNode.on(NodeStatus.SUCCESS, knowledge_node); intentNode.on(NodeStatus.NEED_INPUT, clarify_node); intentNode.on(NodeStatus.FAILED, fallback_node); BaseNodeAgentRequest clarifyNode new BaseNodeAgentRequest(clarify_node, 追问) { Override public NodeStatus execute(NodeContextAgentRequest context) { return confirmAnswer(context) ? NodeStatus.SUCCESS : NodeStatus.FAILED; } }; clarifyNode.on(NodeStatus.SUCCESS, intent_node); clarifyNode.on(NodeStatus.FAILED, fallback_node); // knowledge_node、fallback_node、answer_node 类似按流程编排 engine.registerNode(intentNode); engine.registerNode(clarifyNode); engine.registerNode(knowledgeNode); engine.registerNode(fallbackNode); engine.registerNode(answerNode); return engine; } }这个设计的好处你写几个节点就能体会出来。比如我想在答案生成之前加一个“敏感词检查”节点只需要新建一个节点类然后把answer_node的上游路由改一下其他节点的代码一个字都不用动。5.2 注册节点并启动工作流节点配置完成后启动一次Agent工作流只需要一行代码flowEngine.start(AgentWorkflowConfig.ENTRY_NODE, new AgentRequest(question), traceId);引擎会自动从intent_node开始一路轮转下去。每个节点的执行结果都会触发状态事件这些事件经过监听器、SseEmitter最终变成前端页面上的实时状态更新。我在实际部署中还加了一个内容审核节点放在最后专门检查LLM生成结果是否合规合规度低就重新生成一次最多重试两次。这些控制都通过路由表表达answer_node的下一步根据审核结果分别指向结束节点或者重试节点。流程的复杂度上来了但代码结构一直是同一套。5.3 前端收到的事件流长什么样配套的前端逻辑也很简单。用浏览器原生的EventSource监听事件每次收到消息就更新页面上的步骤状态const eventSource new EventSource(/agent/run?question encodeURIComponent(question)); eventSource.addEventListener(agent-event, (event) { const data JSON.parse(event.data); renderNodeStatus(data.nodeId, data.status); if (data.next) { appendNode(data.next, 等待执行); } });实际在浏览器里看到的执行过程是这样的init请求发出去后页面先显示“意图识别·执行中”1秒后变成“意图识别·成功”同时追加“知识库检索·执行中”检索完成后变成“知识库检索·成功”追加“答案生成·执行中”最后收到一个没有next字段的事件流程结束。整个过程是连续滚动的。这个体验比一个转圈loading强太多了。用户能清楚感知到Agent正在一步步处理问题而不是卡住了。6. 常见问题速查与踩坑记录6.1 节点异常与失败重试怎么处理第一版引擎设计里有个问题节点抛出异常就直接把整个工作流终止了。这在生产环境不行LLM调用经常因为网络抖动或者Token限制失败一次失败就终止整条流程用户体验很糟糕。我后来做了三层兜底。第一层节点内部自己捕获短期异常像LLM超时这种内部重试一次再决定返回什么状态。第二层路由表提供FAILED和TIMEOUT的流转配置失败时可以跳转到专门的重试节点而不是直接终结。第三层引擎捕获未处理异常后广播failed事件前端收到事件后可以提示用户重试同时记录traceId供排查。注意一个比较容易踩的坑节点执行失败后上下文里可能已经写入了一部分脏数据。重试之前必须考虑这些数据的覆盖问题。我一般建议变量写入采用“覆盖式”而不是“追加式”重试节点开始前把上一轮的中间变量清掉防止旧数据污染新结果。6.2 并发场景下的上下文隔离因为工作流是异步执行的同一个引擎实例会同时运行多个Agent请求。这时候最怕的就是节点里不小心把数据写进了共享内存。每个请求的NodeContext是按traceId隔离的节点操作variables里的变量时永远只操作自己上下文的数据这是最基本的原则。但还有一个隐性问题线程池里的线程被不同请求复用如果节点实现里用了ThreadLocal保存上下文一个请求结束后ThreadLocal没有清理下一个请求复用这个线程时就会读到上一个请求的残留数据。我踩过这个坑排查了很久最后强制规定节点内禁止使用ThreadLocal必须通过NodeContext传递数据。如果你的节点实在绕不开ThreadLocal请在finally块里彻底remove。同时建议给线程池设置一个明确的拒绝策略。Agent并发上来之后线程池排满任务新请求不能无限等待。我用的是AbortPolicy配合前端提示“当前服务繁忙”比把请求都堆在内存里等超时要好。6.3 我用这个方案后的一些体会最后说点个人的体会算不上结论更多是经验。状态轮转这套东西最大的收益不是代码变短了而是思考方式变了。以前设计Agent流程我脑子里是一段线性代码的执行顺序现在设计Agent流程我脑子里就是一张有向图——节点是图的顶点路由是图的边所有流程控制都体现在路由关系上。这个转变让流程变得可以观测。引擎里的每一个事件都记录了节点ID、状态、时间戳我可以轻松把这套事件流接到日志系统或者监控面板上。用户问“为什么我这个问题走了兜底”我查一下traceId对应的事件流立刻就能复现完整路径。这在if-else时代是做不到的。如果有朋友也想在项目里落地这套方案我的建议很简单不要一开始就追求大而全先把节点抽象、路由表、异步事件这三样做出来跑通一个最简单的三节点流程然后逐步往里面加节点、加状态、加监听器。等你真的把Agent流程跑起来再回头看那些纠缠在一起的if-else大概率会和我一样再也不想回去了。