上个月接了一个简历筛选的 Agent 需求。一开始我图省事,直接在业务类里用 if-else 把流程串了起来:解析简历,过滤关键词,命中条件就调 LLM 生成摘要,然后保存报告、通知 HR。前两三个节点的时候还挺清爽,等流程加到十个节点、十几个分支判断之后,代码彻底没法看了——每个方法里都是"下一步该调谁"的判断,加一个新人来改需求,得在线索链里翻半天才知道某个分支在哪儿接出去的。
后来我用了三天时间,用 Java 从 0 到 1 写了一个轻量级的 Agent 工作流引擎,核心就两件事:节点状态轮转和流式输出。状态机负责让节点按定义好的流转规则自动推进,业务代码里不再出现"如果 a 成功就调 b,否则调 c"这类硬编码路由;流式输出则让前端能实时看到流程跑到了哪个节点、Agent 当前正在做什么。这篇文章就把整个设计思路、关键代码和踩坑记录完整地写出来,适合已经会用 Spring Boot、但对 Agent 工作流底层实现还好奇的 Java 工程师。
1. Agent工作流为什么需要一个真正的引擎:if-else的失控与Flowable的过度设计
1.1 靠if-else串流程,问题出在状态和路由缝在一起
先还原一下大多数人的第一版写法。假设流程是:接收简历 -> 解析简历 -> 关键词过滤 -> LLM 摘要 -> 保存报告 -> 通知 HR,用 if-else 大概长这样:
void execute(String cvId) { Resume resume = parseResume(cvId); if (resume == null) { retryOrFail(cvId); return; } boolean matched = keywordMatch(resume); if (matched) { String summary = llmSummarize(resume); saveReport(cvId, summary); notifyHr(cvId, "matched"); } else { saveReport(cvId, "no match"); notifyHr(cvId, "not matched"); } }这段代码在需求不变的时候还能忍受。问题是 Agent 工作流几乎每周都在变。比如加一个"解析失败先重试两次再走人工兜底"的规则,加一个"LLM 超时走降级模板"的逻辑,加一个"关键词匹配命中的同时并行生成多份不同视角的摘要"。每加一个规则,你就得在代码里再嵌套一层 if-else,然后方法越来越长,状态越来越多,最后每个节点到底在什么情况下走到哪个分支,只能靠人肉回忆。
我用一段时间后发现,if-else 串流程的核心弊病是把两件事焊死在一起了:路由规则和业务逻辑。节点本身只应该做一件事,比如"解析简历""过滤关键词",但 if-else 写法里,本节点必须在代码里知道"下一个节点是谁、满足什么条件才去"。这套东西一旦数量多了,本质就是在地狱级地维护一张肉眼不可见的状态机。
1.2 Flowable/Camunda为什么不合适
既然 if-else 不行,那直接上 Flowable、Activiti 这类成熟的 BPM 引擎行不行?我的答案是:行,但不是为 Agent 工作流准备的。
Flowable 的优点非常硬核:完整的 BPMN 2.0 规范支持、数据库持久化、历史流程审计、人工审批节点、会签或签等企业级能力。如果你的流程很大概率会涉及人工审批、需要严格合规审计,那直接选 Flowable 没毛病。但 Agent 工作流的真实画风不是这样的——它更像一条动态编排的图:节点是"调用大模型""调工具""做条件判断",边是"成功才能往下走""某条件满足时走这条分支"。你并不需要一个人工审批列表,也不需要完整的历史流程实例库,你真正想要的只是让十几个 Agent 节点按规则有条不紊地跑完,并且把过程实时暴露给前端。
用 Flowable 跑 Agent 工作流,意味着你得把每个 Agent 节点包成 JavaDelegate,在 BPMN XML 里定义各种 serviceTask、sequenceFlow,再配一个流程引擎的 History 库。为了串起十个节点,引入一套完整的事务和任务机制,调度复杂度和运维成本都上去了。而且 Flowable 的默认输出模型是"流程结束时拿结果",和 Agent 场景里想要的"一边跑一边把 token 和节点状态推给前端"天然有张力。
1.3 自研轻量级引擎的适用边界
再往上一层的对比对象,其实是很多人用过的 Coze、Dify 这类平台。它们帮你处理了流程编排、运行界面、多租户,体验确实好。但这里的痛点是:它是一个托管黑盒,流程定义存在平台侧,只能通过 API 调用,流程中间你要是想改上下文管理策略、想把执行事件推到自己的业务系统里观察、想接自己公司的内部 RPC,都得看平台能力脸色。
所以我建议的判断标准是这样的:如果你的流程规模可控(几十个节点以内)、需要深度嵌进现有 Spring Boot 业务系统、希望自定义流式输出事件、有精力维护一点基础设施代码,那自研轻量级引擎非常值得;反之,如果你需要真正的 BPMN 标准、需要人工审批流转,或者你的诉求是"今天就要可视化编排上线",直接上 Flowable 或者用平台,别自己折腾。我的选择是前者,下面开始讲怎么设计。
2. 节点状态机的核心设计:让状态流转与业务逻辑彻底解耦
2.1 单节点状态枚举与状态转移表
工作流引擎里的第一个核心概念是节点状态机。每个节点不会只有"开始/结束"两个状态,它至少应该有这样一套状态:
| 状态 | 含义 | 触发条件 |
|---|---|---|
| PENDING | 等待执行,已入队但还没轮到 | 流程启动时初始状态 |
| RUNNING | 执行中,可能正在调 LLM 或工具 | 被执行器取走并开始运行 |
| SUCCESS | 执行成功,已产出结果 | 执行器返回成功结果 |
| FAILED | 执行失败 | 执行器返回失败且不再重试 |
| SKIPPED | 被跳过 | 上游条件不满足,或降级策略生效 |
| TERMINATED | 被强制终止 | 全局取消或流程终止 |
配套的还有流程级状态:READY、RUNNING、COMPLETED、FAILED、TERMINATED。节点状态和流程状态要分清楚,前者是引擎里单个节点的生命周期,后者是整个流程实例的生命周期。
这套状态转移关系用一个表格给出来,就当是状态机定义:
| 当前状态 | 事件 | 下一状态 | 说明 |
|---|---|---|---|
| PENDING | START | RUNNING | 执行器开始干活 |
| RUNNING | COMPLETED | SUCCESS | 节点正常结束 |
| RUNNING | RETRY | RUNNING | 失败但允许重试,重新执行 |
| RUNNING | FAILED | FAILED | 失败且不再重试 |
| PENDING 或 RUNNING | SKIP | SKIPPED | 路由判定不满足,跳过 |
| 任意状态 | TERMINATE | TERMINATED | 流程被取消时全局终止 |
实际代码里,我建议把这套状态转移逻辑收敛到一个小的状态机工具类里,里面维护一张"当前状态 -> 事件 -> 目标状态"的映射表,而不是散落在各节点里自己 setState。
2.2 用状态机代替散落各处的if-else判断
这一节是整个引擎的灵魂。以前我们用 if-else 管流程,是让每个节点"自己决定下一步去哪"。改成状态机之后,节点执行器彻底变成"哑节点"——它只负责处理自己的业务逻辑,返回一个 ExecutionResult,里面包含 status 和数据,绝不包含"下一个节点是谁"这个信息:
public interface NodeExecutor { String type(); ExecutionResult execute(NodeContext ctx); }ExecutionResult 大概长这样:
public class ExecutionResult { private NodeStatus status; // SUCCESS / FAILED / SKIPPED private Map<String, Object> data; // 节点输出,写入共享上下文 private int retryCount; // 配合 maxRetries 用 }那么"下一个节点是谁"这个路由决策交给谁?交给引擎,由引擎根据流程定义里的 Transition 数组去判定。流程定义大概是:
{ "id": "keyword-filter", "type": "keywordFilter", "transitions": [ { "condition": "#ctx.get('matched') == true", "target": "llm-summarize" }, { "condition": "#ctx.get('matched') == false", "target": "save-report" } ] }节点跑完,引擎把 ExecutionResult 里的数据合并进共享上下文,然后逐个评估 transitions 里的 condition,命中哪条就走哪个 target。业务代码里那个"如果匹配就走摘要,否则保存报告"的判断,彻底从 Java 方法里消失,变成了 JSON 里的一行路由规则。这就把"状态轮转"从"业务代码"里解耦出来了。以后要加一个新分支,改 JSON 就好,不用动 Java 代码。
2.3 节点执行器的注册表设计
状态机定了,引擎还需要知道"type 为 keywordFilter 时,究竟调用哪个 Java 类"。这里本来可以写一个 switch-case,但那就又把类型和实现耦合死了。我用的是 Spring 环境下最常见的注册表思路:定义注解,启动时扫描注册。
@Target(ElementType.TYPE) @Retention(RetentionPolicy.RUNTIME) public @interface NodeExecutorType { String value(); }每个执行器加注解,例如:
@NodeExecutorType("keywordFilter") @Component public class KeywordFilterExecutor implements NodeExecutor { @Override public ExecutionResult execute(NodeContext ctx) { // 业务逻辑... } }引擎启动时通过 ApplicationContext 拿到所有带注解的 Bean,构建成 Map<String, NodeExecutor>:
@Component public class NodeExecutorRegistry { private final Map<String, NodeExecutor> registry = new HashMap<>(); public NodeExecutorRegistry(ApplicationContext applicationContext) { Map<String, Object> beans = applicationContext.getBeansWithAnnotation(NodeExecutorType.class); for (Object bean : beans.values()) { NodeExecutorType annotation = bean.getClass().getAnnotation(NodeExecutorType.class); registry.put(annotation.value(), (NodeExecutor) bean); } } public NodeExecutor get(String type) { NodeExecutor executor = registry.get(type); if (executor == null) { throw new IllegalStateException("unknown node executor type: " + type); } return executor; } }这样新增一种节点类型,只需要新写一个执行器类并打上注解,引擎的循环逻辑完全不用动。这个"注册表"模式,本质上就是干掉了把所有执行器塞进一个巨型 if-else 或 switch-case 的做法。
3. 流程定义层:用JSON编排一张Agent执行图
3.1 流程定义JSON的最小模型
流程引擎能不能用得舒服,流程定义的表达力占一大半。我设计的 JSON 模型保持最小化,但足够支撑 Agent 工作流的常见需求:
{ "flowId": "resume-screening-flow", "startNode": "receive-resume", "nodes": [ { "id": "receive-resume", "type": "inputReceiver", "name": "接收简历", "params": { "requiredFields": ["cvId", "source"] }, "next": "parse-resume" }, { "id": "keyword-filter", "type": "keywordFilter", "params": { "keywords": ["Java", "Agent", "工作流"] }, "transitions": [ { "condition": "#ctx.get('matched') == true", "target": "llm-summarize" }, { "condition": "#ctx.get('matched') == false", "target": "save-report" } ] } ] }这个模型里我故意把next和transitions都保留:无条件的直连边用next,表达起来更简洁;有条件的边用transitions,里面每一项都是"条件 + 目标节点"。一个节点上多个 transitions 会被按顺序评估,命中第一个就停止,这样天然支持"优先级路由"。params是节点执行器自己消费的参数,比如关键词列表、超时时间、模型 ID 等,从 JSON 里解析出来装入 NodeContext。
这套模型的核心思想是:流程定义是"图"数据,不是"代码"数据。每个节点只知道自己的出边,不需要知道上游是谁;整张图可以被解析、校验、可视化,甚至前端可以直接拿这棵"节点关系树"渲染一个简版流程图。
3.2 用SpEL作为条件路由的表达式引擎
条件表达式我直接用 Spring 的 SpEL,这是 Java 世界里现成的、表达能力足够、又不需要引入额外重依赖的表达式方案。引擎侧写一个很小的 ConditionEvaluator:
@Component public class ConditionEvaluator { private final ExpressionParser parser = new SpelExpressionParser(); public boolean evaluate(String expression, NodeContext ctx) { if (expression == null || expression.isBlank()) { return true; } StandardEvaluationContext context = new StandardEvaluationContext(); context.setVariable("ctx", ctx); try { return Boolean.TRUE.equals(parser.parseExpression(expression).getValue(context, Boolean.class)); } catch (Exception e) { // 表达式异常时,按不满足条件处理,走兜底 return false; } } }然后在 JSON 里写"condition": "#ctx.get('matched') == true",引擎解析时把整个 NodeContext 作为变量ctx暴露给表达式。这里有几个够用就行的小技巧:
- 表达式里统一走
#ctx.get('key'),不要在表达式里直接拿 Java 对象的复杂字段,否则 JSON 定义很容易和苏式化。 - 不满足条件就"静默走 else",不会因为表达式解析失败让整个流程崩掉。
- 如果想调试,可以在表达式前后打印一下变量快照,后面我会讲一个调试套路。
用 SpEL 的意义在于,流程边界条件的调整不再需要发版,普通开发也能读懂 JSON 里的条件路由,反正就是"匹配了吗"这类简单判断。
3.3 上下文对象:Agent工作流的黑板
Agent 工作流和传统表单审批流最大的不同,是它有一个非常"胖"的共享上下文——各节点要传递的不只是一个个审批结果,还有文档内容、解析结果、LLM 输出、工具调用痕迹。这里我采用底层叫法叫Blackboard(黑板)模式:所有节点共享同一个黑板,每个节点往黑板上写入自己的输出,后续节点按 key 读取需要的内容。
public class NodeContext { private final String flowId; private final String nodeId; private volatile Map<String, Object> variables = new ConcurrentHashMap<>(); // getVariable / setVariable / removeVariable / snapshot }比较关键的一点是主动控制上下文体积。Dify 这类产品里经常遇到"上下文超长"的问题,那是因为平台把所有中间结果都堆在一起传给大模型。自研最大的好处是可以在节点执行前做"上下文摘要":比如 LLM 节点只需要最近两条关键数据,就用一个小的上下文组装器,从黑板上只挑需要的 key 拼成 prompt,而不是把整张黑板倒给模型。这块让我尝到甜头的地方,后面在踩坑章节里会再展开。
4. 引擎执行核心:图遍历、分支汇聚与并发控制
4.1 队列驱动的图遍历:为什么不用递归
拿到流程 JSON,也拿到了注册表,接下来是引擎主循环。不少新手会自然想到用递归去做深度优先遍历:从 startNode 开始递归执行,执行完 a 节点就递归调 b。但我不建议这么做,原因是 Agent 工作流里一个节点可能要调 LLM 等几十秒,递归版本的调用栈会挂起一长串,中断、恢复、超时控制都很别扭,而且节点多了还容易栈溢出。
我采用的是队列驱动的遍历:先用一个待执行队列把 startNode 放进去,循环里取一个节点、执行、根据路由把后继节点放回队列,直到队列为空或流程被终止。核心逻辑大概是:
public void start(WorkflowDefinition definition, Map<String, Object> initialData) { Deque<String> readyQueue = new ArrayDeque<>(); readyQueue.push(definition.getStartNode()); while (!readyQueue.isEmpty()) { String nodeId = readyQueue.poll(); WorkflowNode node = definition.getNode(nodeId); if (!stateMachine.casState(nodeId, PENDING, RUNNING)) { continue; // 状态推进失败,说明被其他线程抢跑了,跳过 } ExecutionResult result = executeNode(node); mergeResultToContext(nodeId, result); List<String> nextNodes = route(node, result); nextNodes.forEach(readyQueue::push); } }这个 while+队列的版本好处非常实际:执行顺序可预测,容易做超时和终止控制,线程模型也更贴近"任务调度器"而不是"递归嵌套"。而且流程做到一半如果想取消,直接把流程状态置为 TERMINATED,循环下一次取节点时发现终止标记,就直接收尾,不需要处理一串递归栈。
4.2 ALL/ANY汇聚机制与并行分支
Agent 工作流经常会遇到"一个节点成功后,两个分支并行执行"的场景。最典型的就是:解析完简历,一边调 LLM 生成摘要,一边生成简历结构化数据文件,两个互不依赖。这时候单线程按顺序跑没问题,但没有发挥效率,所以我给流程节点加了两个汇聚策略:
ALL:所有上游完成后,本节点才允许执行ANY:任意一个上游完成后,本节点即可执行
实现方式我用一个简单的"输入计数"来推进。每个节点在运行时会有 inputCount,每当一条上游边完成,就对该节点的计数做一次原子减一;减到 0(ALL 场景)或减到 1(ANY 场景)时,就把该节点投放入待执行队列。核心数据结构是:
ConcurrentHashMap<String, AtomicInteger> dependencyCounter = new ConcurrentHashMap<>();并行部分的线程池,我建议如果 JDK 版本允许直接用虚拟线程池,老版本就用固定大小线程池。另外补充一个很多人会问的"Agent 怎么扛并发"的问题:引擎本身的并发能力只是一部分,更关键的是要在调用 LLM 的节点上做好信号量限流、给事件队列设置上限,避免一个流程跑十条分支时把 LLM 供应商的配额打爆、把内存打爆。
4.3 循环、重试与防死循环
有向图只要不是纯 DAG,就可能出现环。在 Agent 工作流里,环通常是两种意图:一是"重试直到成功",二是"循环直到用户输入确认"。我在节点上配置maxRetries和maxIterations两个参数来分别控制。
重试逻辑放在状态机里:RUNNING 状态收到 FAILED 事件时,如果重试次数没到阈值,状态回到 RUNNING,节点被再次投放入队列;如果到了,直接 FAILED,然后走降级路由。循环逻辑则用一个"剩余迭代次数"放进上下文,每经过一次循环节点就减一,减到 0 之后,路由不再回跳到循环头,而是跳到循环出口或降级节点。
这个防呆设计非常重要。我见过同事直接用 while 循环写递归式工作流,最后因为一次 LLM 输出的条件表达式没命中,流程在"确认 -> 重试 -> 确认"这个环里卡了一整夜。有界循环 + 条件路由,才能保证流程不会变成生产事故。
5. 流式输出:让Agent的每一步都可以实时推送到前端
5.1 事件的类别设计
流程引擎跑起来只是第一步,Agent 工作流动辄几十秒甚至几分钟,用户看不到中间过程就等于"死页面"。所以流式输出是这次重构里我最看重的部分。我设计了一套事件模型,引擎在关键节点都会发事件出来:
public abstract class FlowEvent { private final String flowId; private final Instant timestamp; // 状态流转事件:NODE_STARTED / NODE_COMPLETED / NODE_FAILED / FLOW_FINISHED // 内容增量事件:TOKEN_PRODUCED }具体事件大概有这几种:
NodeStartedEvent:某个节点开始执行,前端可以渲染"正在解析简历"NodeCompletedEvent:某个节点执行完成,前端可以渲染"简历解析完成,耗时 1.2 秒"TokenProducedEvent:LLM 节点产生的增量 token,前端可以做打字机效果FlowTerminatedEvent:整个流程结束或失败
这里要提一下,如果你在 Agent 节点里接了 MCP 工具,工具运行的时候往往会产生持续的增量输出——比如一边调用工具一边把内容写到文件、或者一边流式生成一段报告。这些增量输出我也会统一包装成TokenProducedEvent发布出去,这样前端的流式渲染逻辑不用分心去处理不同来源。
5.2 事件总线与SSE对接
事件生产出来了,怎么推给前端?我用的方案是:引擎内部维护一个事件总线,外部注册监听者。在 Spring Boot 里可以直接用一个简单的实现:
@Component public class FlowEventBus { private final ConcurrentHashMap<String, CopyOnWriteArrayList<FlowEventListener>> listeners = new ConcurrentHashMap<>(); public void subscribe(String flowId, FlowEventListener listener) { listeners.computeIfAbsent(flowId, k -> new CopyOnWriteArrayList<>()).add(listener); } public void publish(FlowEvent event) { String flowId = event.getFlowId(); CopyOnWriteArrayList<FlowEventListener> list = listeners.get(flowId); if (list != null) { for (FlowEventListener listener : list) { listener.onEvent(event); } } } }流式输出到前端的传输协议,我直接选 SSE,理由后面讲。Controller 里接 SseEmitter:
@GetMapping(value = "/flows/{flowId}/events", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter streamEvents(@PathVariable String flowId) { SseEmitter emitter = new SseEmitter(0L); FlowEventBus.subscribe(flowId, event -> { try { emitter.send(SseEmitter.event().name(event.getType()).data(event)); } catch (IOException e) { emitter.completeWithError(e); } }); return emitter; }前端一个 EventSource 就接住了。事件驱动和 SSE 结合之后,整个引擎的运行过程对用户是完全可见的:流程走到哪一步、哪个节点在跑、大模型生成了什么内容、最后结果是什么,全都能在页面上实时滚出来。
5.3 为什么选SSE而不是WebSocket
在 Agent 工作流这个场景,几乎是服务器单向推送为主。用户在页面发起一个请求,然后整个流程的进展都是服务器主动往客户端推,客户端基本不需要往服务器发什么业务消息。这种"单向推送"模型下,SSE 是最合适的选择:基于 HTTP,实现简单,有自动重连,前端使用成本极低,一个大大的EventSource对象就完事。
WebSocket 更适合双方高频双向交互的场景,比如聊天、白板协同。如果为了流式输出硬上 WebSocket,你得处理连接管理、心跳、双通道消息协议,复杂度高出不少,收益在这个场景里却几乎为零。真要说需要客户端做点什么,也就是"取消流程",那把取消单独做成一个 POST 接口就够了,完全没有必要为这一个交互升级成 WebSocket。
5.4 心跳、背压与事件降频的几个小细节
这几个细节看着小,实际影响体验很大。
第一是心跳。SSE 连接长时间没数据,部分网关和浏览器会自动断开。我在 SseEmitter 里加了个定时任务,每 15 秒发一条注释型心跳,保证连接不假死。第二是有界事件队列。如果某个流程瞬时生成大量 token 事件,而客户端消费不及,无界队列就会内存暴涨。我统一用有界队列(ArrayBlockingQueue),满了之后对TokenProducedEvent做合并降频,把多个小片段合成一个大片段;但NodeStartedEvent、NodeCompletedEvent这类核心事件只允许阻塞等待,不能丢失。这样体验不至于变成"卡帧",底层也稳得住。
6. 全链路演示:简历筛选Agent,从JSON到SSE跑通
6.1 业务场景与节点设计
用一个完整的简历筛选 Agent 来演示这套引擎。场景是候选人提交简历后,系统自动完成解析、过滤、摘要、报告、通知的全流程,前端的 HR 页面能看到流程实时进度。
节点清单如下:
| 节点 ID | 执行器 type | 职责 | 路由 |
|---|---|---|---|
| receive-resume | inputReceiver | 接收上传的简历文件 | next -> parse-resume |
| parse-resume | resumeParser | 解析简历文本与结构化数据 | next -> keyword-filter |
| keyword-filter | keywordFilter | 检查是否包含目标关键词 | 条件路由到 llm-summarize 或 save-report |
| llm-summarize | llmSummarizer | 调用 LLM 生成简历摘要 | next -> save-report |
| save-report | reportSaver | 保存最终报告 | next -> notify-hr |
| notify-hr | hrNotifier | 通知 HR 处理 | 结束 |
6.2 流程定义JSON完整示例
完整 JSON 长这样。这里我为了演示方便,把 LLM 节点也放进路由里,方便你理解条件路由和并行分支的真实形态:
{ "flowId": "resume-screening-flow", "startNode": "receive-resume", "nodes": [ { "id": "receive-resume", "type": "inputReceiver", "name": "接收简历", "params": { "requiredFields": ["cvId", "source"] }, "next": "parse-resume" }, { "id": "parse-resume", "type": "resumeParser", "params": { "key": "resumeRaw" }, "next": "keyword-filter" }, { "id": "keyword-filter", "type": "keywordFilter", "params": { "keywords": ["Java", "Agent", "工作流"] }, "transitions": [ { "condition": "#ctx.get('matched') == true", "target": "llm-summarize" }, { "condition": "#ctx.get('matched') == false", "target": "save-report" } ] }, { "id": "llm-summarize", "type": "llmSummarizer", "params": { "promptTemplate": "请总结这份简历…", "maxTokens": 1024 }, "next": "save-report" }, { "id": "save-report", "type": "reportSaver", "params": { "outputKey": "reportPath" }, "next": "notify-hr" }, { "id": "notify-hr", "type": "hrNotifier", "params": { "channel": "email" } } ] }我开发时喜欢用"JSON 里永远不写死下一个节点 ID 之外的业务判断"这个原则,所以上面的 keyword-filter 里你看到的是 transitions 条件路由,而不是在 Java 里写 if。这样调整过滤条件或增加一个"AI 初筛不通过进人工池"的节点,只需要改 JSON 的 transitions 和新增一个节点,不动任何别人写好的执行器。
6.3 引擎启动与SSE接入代码
启动流程的入口,对业务方来说应该简单到几个方法:
String flowId = workflowEngine.start("resume-screening-flow", initialData); // 返回 flowId,前端拿这个 flowId 去订阅 SSE引擎内部 start 方法的逻辑就是把流程定义做一次快照,创建 NodeContext,初始化状态机,把 startNode 放入待执行队列,然后异步开始跑。SSE 订阅代码我在 5.2 已经给出,两个部分拼起来就是一个完整闭环:POST /flows/{flowId}/start 开启流程,GET /flows/{flowId}/events 接流,等待 FlowTerminatedEvent 出现后前端关闭 EventSource,任务完成。
为了演示不依赖真实 LLM 也能跑通,我在 llmSummarizer 执行器里做了一个 fake 版本,生成一个固定模板摘要并按每 150 毫秒一个 chunk 抛 token 事件。这样本地启动项目之后,打开页面就能看到"正在解析简历 -> 关键词命中 -> LLM 逐字输出摘要 -> 保存报告 -> 通知 HR"一条龙推进。
6.4 本地调试的实操技巧
最后一个环节讲讲调试。流程引擎类项目最怕的就是"节点多了之后不知道卡在哪儿"。我自己常用的几个手段:
- 在 NodeContext 里加一个
snapshot()方法,每个节点执行前后打印关键变量,观察数据在每个环节怎么演化。 - 设置"慢放模式":给每个节点执行前插一个
Thread.sleep(200),配合 SSE 页面看事件顺序是否合理。这比断点调试更能看出路由问题。 - 给 LLM 节点做一个 Mock 开关,切到 mock 模式后不真实调用外部 API,返回固定数据。调试和 CI 都稳。
- 所有事件日志统一格式,按 flowId 分组,排查问题一条命令 grep 出来。
这些手段不需要任何额外的监控系统,纯粹靠引擎事件体系带来的透明性就足够跑通大部分场景。
7. 实战踩坑与边界:比if-else版本低了多少维护成本
7.1 并发下重复执行:状态CAS推进
第一个坑是并行分支跑起来之后发现的:两个线程同时从队列里取到了同一个 PENDING 节点,在没有保护的情况下,这个节点会被同时执行两次。如果节点是"发通知",那用户会收到两条一模一样的邮件;如果是"扣费"节点,那就是生产事故。
修复思路是在状态机上做 CAS 推进。节点从 PENDING 推进到 RUNNING 是一个"只允许一次"的操作,实现上是这样:
AtomicReference<NodeStatus> statusRef = new AtomicReference<>(PENDING); boolean ok = stateMachine.casState(nodeId, PENDING, RUNNING); if (!ok) { // 说明别的线程已经抢跑了,当前线程直接忽略 }所以看到我前面主循环里的continue逻辑了吗?那个不是随便写的,就是为了配合 CAS 状态推进,拦截掉重复执行的节点。状态推进的原子性是并发流程引擎的底线。
7.2 流式输出的背压与内存控制
第二个坑是事件队列。刚开始我把所有事件丢进一个无界队列,结果是:某个流程里 LLM 节点疯狂输出 token,而前端消费速度跟不上,JVM 堆一路涨上去,最后 OOM。排查的时候看堆 dump,里面全是待发送的 SseEmitter 消息对象。
优化方案前面提过:事件通道换 ArrayBlockingQueue,并且区分重要事件和低频事件。NodeCompletedEvent这类必须保证送达,采用阻塞投递;TokenProducedEvent这类高频低价值事件,采用合并批量投递。我还在中间做了一个小压缩策略:同一节点产生的 token,50 毫秒内的合并成一段再推送,前端用打字机效果完全感知不到延迟,但内存占用一下子降了一个数量级。
7.3 我踩过的三个真实坑
再列三个小但真实的坑,你们遇到了能少走弯路:
第一个是 SpEL 表达式写错导致的静默路由错误。比如把#ctx.get('matched') == true写成了#ctx.get("matched") == true,JSON 里外层用了双引号,SpEL 解析直接报错,我的 ConditionEvaluator 又做了"解析异常返回 false"处理,结果所有匹配到的简历全走了不匹配分支。排查了很久才反应过来。后面我在解析异常时增加了 error 日志,表达式错误必须抛出来,绝不静默吞掉。
第二个是 LLM 节点超时导致整个流程挂起。Agent 工作流里最慢的就是模型调用,如果不给节点设置超时,一次上游模型接口的偶发卡顿会让整个流程几分钟没人管。我现在每个慢节点都配上超时时间,超时后走 FAILED 路由,宁可降级也不能挂死。
第三个是日志里的敏感内容。简历里有手机号、邮箱、薪资期望,流程日志和事件推送如果不做脱敏,面试者隐私和公司数据合规都会有隐患。我在 NodeContext 的 snapshot 里增加了脱敏处理,涉及个人信息的字段一律打码,事件推送也只传摘要不进原文。这个细节你们用真实业务数据时一定要留意。
7.4 千万别自研引擎的情况
最后说点泼冷水的话。有些场景你要冷静判断,别因为看了这篇文章就什么都自研:
- 你的流程需要严格遵循 BPMN 标准,需要走人工审批、会签、或签,直接上 Flowable,别自己写。
- 你需要一个可视化编排界面,让非技术人员自己拖拽改流程,而你的团队没有前端人力去搭这个配置台,用 Coze/Dify 这类平台更合适。
- 你的流程规模大(上百个节点)、需要完整的历史实例查询和审计,自研引擎的成本会指数级上升。
- 团队本身没有 Java 工程师长期维护基础设施代码,只是为了"不用 if-else"就把引擎写出来,那反而引入一个更复杂的长期维护负担。
我自己写这套轻量级引擎,核心动机是做一个可以被业务代码深嵌、可以自定义事件输出、运行模型完全可控的运行时。这套引擎我已经用在内部工具链里大半年了,从维护成本来看,比之前 if-else 版本低了不止一个量级——改流程变成改 JSON,加节点变成新增一个执行器类,排查问题变成看事件流,整个工作流就是"状态轮转 + 流式输出"两个词。如果你也在做 Java 侧的 Agent 编排、还没想好怎么组织流程代码,我建议你先别急着上重框架,画清楚节点图、定好状态机,写一个几百行的小引擎,你对工作流的理解会完全不一样。最后一个小建议:先设计好 NodeContext 的变量结构,再谈引擎选型,这是我吃亏之后最想提醒你的。