轻量级Java Agent工作流引擎:告别if-else,实现动态任务编排

📅 发布时间:2026/10/7 5:52:46
轻量级Java Agent工作流引擎:告别if-else,实现动态任务编排
1. 这不是又一个“流程引擎Demo”而是一套能直接塞进生产环境的Agent工作流内核你有没有写过这样的代码一个审批逻辑if-else嵌套五层每个分支里还要调用不同服务、发不同通知、更新不同状态加个新节点得翻三四个类改七八处判断测试不敢动上线心发慌。更别提现在AI Agent场景下用户一句话进来要并行调用知识库、调用LLM、调用工具、做结果校验、再流式返回——传统硬编码根本扛不住。标题里说“再也不用if-else写法了”真不是标题党。我去年在给一家智能客服中台重构对话路由模块时就彻底扔掉了所有状态机switch-case用一套纯Java手写的轻量级流程引擎替换了原来2000行胶水代码。它不依赖Spring Flow、不引入Activiti这种重型框架核心就三个类Node节点、Workflow流程定义、Engine执行器但支持节点状态轮转、异步编排、失败重试、上下文透传、流式输出回调——最关键的是所有逻辑都靠配置驱动新增一个“语音转文字”节点只需写一个实现类注册到配置里前端拖拽连线后端零代码改动。这背后不是炫技而是对Agent工作流本质的理解它不是业务流程自动化BPM而是意图驱动的动态任务编排系统。节点不是“审批”“归档”这种静态动作而是“调用RAG检索”“执行Python沙箱”“生成SVG图表”这种原子能力封装状态轮转不是“待审→已审→归档”而是“pending→executing→streaming→completed”或“pending→executing→failed→retrying”。所以本文不讲理论模型只拆解我压在生产线上跑了一年多的这套引擎怎么从0写起、怎么让每个节点真正“活”起来、怎么把LLM的token流稳稳接住再推给前端——所有代码都在JDK8上跑没用任何AI框架连Jackson都只用基础序列化就是纯Java的肌肉感。2. 核心设计哲学为什么放弃Activiti/Spring Flow选择手写2.1 传统工作流引擎的“水土不服”Activiti、Flowable这类BPM引擎设计初衷是管理企业级审批流固定节点、预设路径、人工干预多、事务强一致性。但Agent工作流完全相反节点高度动态今天用Qwen3做摘要明天换成DeepSeek-V3节点实现类随时替换但流程拓扑不变状态粒度极细LLM调用不能只分“成功/失败”得区分“请求已发”“首token到达”“流式输出中”“超时中断”上下文非结构化不是简单的MapString, Object而是可能包含大段文本、嵌套JSON、二进制文件句柄、甚至WebSocket会话引用执行模式异构有的节点CPU密集本地模型推理有的IO密集API调用有的需长连接SSE流混在一个线程池里调度必崩。我试过把Coze工作流导出的JSON强行喂给Activiti结果发现Activiti的Execution对象根本存不下10MB的PDF解析结果它的TaskService设计用来查“张三待审批”不是查“用户ID123456的第3次重试记录”最致命的是它没有onTokenReceived()这种回调钩子——而这是流式输出的生命线。Spring Flow更糟它把流程当成“状态机DSL”但Agent需要的是“能力编排事件驱动”两者范式冲突。2.2 我们要的不是“引擎”而是“执行总线”所以最终定下三个设计铁律节点即服务Node as Service每个节点必须实现NodeExecutor接口只暴露execute(Context ctx)方法内部自己决定同步/异步、重试策略、超时控制。引擎不干涉具体执行只管调度和状态流转。状态即事实State as Fact节点状态不是枚举值而是带时间戳、错误堆栈、输出快照的NodeStatus对象。比如streaming状态必须包含lastTokenTime和tokenCount方便前端做心跳保活。流式即契约Streaming as Contract所有支持流式输出的节点必须实现StreamableNode接口提供subscribe(ConsumerString tokenHandler)方法。引擎不处理token拼接只负责把handler透传给节点——这样节点内部可以用OkHttp的SSE解析器也可以用WebClient的Flux完全解耦。这个设计让引擎体积压到不到200行核心代码。Engine.run(Workflow workflow, Context initContext)方法里真正的调度逻辑只有解析流程图DAG拓扑排序节点按顺序触发node.execute(ctx)捕获异常更新节点状态为failed记录errorStackTrace若节点实现StreamableNode则在其execute内主动调用tokenHandler.accept(token)。没有状态机、没有持久化、没有事务管理——因为Agent工作流的“事务”由LLM自身保证如函数调用失败自动重试引擎只做可靠传递。2.3 为什么选Java不是因为“企业级”而是因为“可控性”看到热搜词里一堆“java面试题”“java基础”可能有人觉得Java太重。但恰恰相反在Agent场景下Java的确定性是救命稻草。内存可控LLM输出动辄几万token用Python的GC你永远不知道什么时候OOMJava可以精确控制堆外内存如用ByteBuffer.allocateDirect存原始token流GC日志一目了然线程模型清晰ForkJoinPool.commonPool()处理CPU密集型节点如本地embeddingExecutors.newCachedThreadPool()处理IO密集型如HTTP调用ScheduledThreadPoolExecutor管重试——每种负载有专属线程池不会互相拖垮调试友好当某个节点流式输出卡住jstack一把抓出线程栈立刻定位是OkHttp的EventSource没关闭还是Netty的ChannelHandlerContext写半包。Python的asyncio调试试试看。我们团队用这套引擎支撑日均50万次Agent调用平均延迟1.2秒99.9%成功率。上线后第一件事不是压测而是把所有节点的execute方法加上Timed注解Micrometer用Prometheus看每个节点P95耗时——这才是Java工程师该干的事不是在YAML里调参数。3. 核心细节解析节点状态轮转与流式输出的落地陷阱3.1 状态轮转不是状态机而是“状态快照链”传统状态机用state State.EXECUTING这种变量但Agent需要知道“为什么是EXECUTING”。所以我们定义NodeStatus为不可变对象public final class NodeStatus { public final String nodeId; public final State state; // PENDING, EXECUTING, STREAMING, COMPLETED, FAILED, CANCELLED public final Instant startTime; public final Instant endTime; public final String output; // 首次输出的前100字符用于快速预览 public final Throwable error; // 仅FAILED状态有 public final MapString, Object metrics; // 自定义指标如tokenCount, apiCost }关键点在于每次状态变更都生成新对象旧状态存入Context的statusHistory列表。这样做的好处是调试时可回溯“节点A为什么卡在STREAMING”——查历史状态发现上一次STREAMING的endTime为空但startTime是3分钟前说明流式连接断了没触发COMPLETED监控时可聚合统计所有FAILED状态的error.getClass().getName()发现70%是TimeoutException立刻优化超时配置前端可渲染把statusHistory转成时间轴用户看到“10:02:15 开始调用知识库 → 10:02:18 收到首token → 10:02:22 流式输出中...”。提示不要用enum State直接赋值我踩过的坑某次升级OkHttp到4.xEventSource的onComplete()回调有时不触发导致状态卡在STREAMING。后来改成状态变更必须显式调用context.updateNodeStatus(nodeId, newStatus)并在updateNodeStatus里加日志埋点问题立刻暴露。3.2 流式输出的“三明治”结构前端、网关、引擎协同流式输出不是引擎单方面的事而是三层协作层级职责关键实现前端接收token、拼接、渲染、防抖EventSource监听/api/workflow/{id}/stream用textContent event.data每50ms触发一次DOM更新避免频繁重排网关协议转换、连接保活、错误透传Spring Boot WebMvc中GetMapping(value /stream, produces MediaType.TEXT_EVENT_STREAM_VALUE)用SseEmitter包装引擎的tokenHandler设置setTimeout(30000)防止Nginx断连引擎生成token、触发回调、状态同步节点执行时tokenHandler.accept(token)被调用引擎立即更新NodeStatus的stateSTREAMING并记录lastTokenTime最易错的是网关层。很多人直接emitter.send(SseEmitter.event().data(token))但忘了SseEmitter默认缓冲区1024字节大token如base64图片会阻塞必须emitter.setBufferSize(8192)Nginx默认proxy_buffering on会攒满才推给前端必须在location里加proxy_buffering off; proxy_cache off;如果节点执行超时SseEmitter会抛IllegalStateException必须捕获并调用emitter.complete()否则连接悬空。我们实测下来用SseEmitter比WebSocket更稳——因为SSE天然支持自动重连前端new EventSource(url)自带重试而WebSocket要自己实现心跳和断线重连。3.3 Context上下文不是Map而是“能力注入容器”Context是引擎的血液但它绝不是HashMapString, Object。我们定义public class Context { private final MapString, Object data new ConcurrentHashMap(); private final MapString, NodeStatus statusHistory new ConcurrentHashMap(); private final ExecutorService ioExecutor; // IO专用线程池 private final ExecutorService cpuExecutor; // CPU专用线程池 private final ScheduledExecutorService retryScheduler; // 重试调度器 // ... 其他能力 }这样设计节点拿到Context就能直接调用ctx.getIoExecutor().submit(() - callApi())—— 不用自己new线程ctx.getCpuExecutor().invokeAll(tasks)—— 并行处理多个embeddingctx.retryScheduler.schedule(() - node.execute(ctx), 2, TimeUnit.SECONDS)—— 失败后2秒重试。注意Context必须是线程安全的所有data操作用ConcurrentHashMapstatusHistory也用并发Map。曾有个节点用context.getData().put(result, bigJson)结果另一个线程正在遍历data.keySet()触发ConcurrentModificationException。解决方案Context提供T T get(String key, ClassT type)和void put(String key, Object value)封装内部加锁或用computeIfAbsent。4. 实操过程从零开始手写流程引擎的完整步骤4.1 第一步定义核心接口15分钟创建core包写三个接口NodeExecutor所有节点的父接口public interface NodeExecutor { /** * 执行节点逻辑 * param context 执行上下文 * return 节点输出可为null */ Object execute(Context context) throws Exception; }StreamableNode流式节点扩展接口public interface StreamableNode extends NodeExecutor { /** * 订阅token流 * param tokenHandler token处理器 */ void subscribe(ConsumerString tokenHandler); }Workflow流程定义public class Workflow { private final String id; private final ListNodeDefinition nodes; // 节点定义列表 private final MapString, String edges; // nodeA - nodeB 的边映射 // getter/setter... }NodeDefinition包含nodeId,className,configJSON字符串这样节点实现类可从配置动态加载不用硬编码。4.2 第二步实现Context与状态管理30分钟Context类重点实现构造时初始化三个线程池ioExecutor用newCachedThreadPoolcpuExecutor用newWorkStealingPoolretryScheduler用newScheduledThreadPool(2)updateNodeStatus()方法先存入statusHistory再更新data里的nodeStatus_${nodeId}最后触发onStatusChange事件供监控用getOrCreateNodeStatus()懒加载避免空状态占用内存。状态枚举State定义为public enum State { PENDING, // 节点已入队未执行 EXECUTING, // 正在执行execute()但未产生输出 STREAMING, // 已开始流式输出但未结束 COMPLETED, // execute()正常返回或流式输出完成 FAILED, // execute()抛异常或流式中断 CANCELLED // 主动取消如用户中断 }注意STREAMING和COMPLETED不是互斥的——流式节点执行完execute()后状态是COMPLETED但期间会多次进入STREAMING。4.3 第三步编写引擎核心20分钟Engine类只做三件事run(Workflow workflow, Context context)主入口public void run(Workflow workflow, Context context) { // 1. 拓扑排序获取执行顺序 ListString executionOrder topologicalSort(workflow); // 2. 按序执行每个节点 for (String nodeId : executionOrder) { try { context.updateNodeStatus(nodeId, new NodeStatus(nodeId, State.EXECUTING, Instant.now())); NodeExecutor node createNodeInstance(workflow.getNode(nodeId)); Object result node.execute(context); context.updateNodeStatus(nodeId, new NodeStatus(nodeId, State.COMPLETED, Instant.now(), result)); } catch (Exception e) { context.updateNodeStatus(nodeId, new NodeStatus(nodeId, State.FAILED, Instant.now(), e)); throw e; // 或者记录后继续下一个节点 } } }createNodeInstance(NodeDefinition def)用Class.forName(def.getClassName()).getDeclaredConstructor().newInstance()动态加载节点topologicalSort()标准Kahn算法实现处理DAG环检测Agent流程严禁环检测到环直接抛IllegalArgumentException。实操心得topologicalSort必须返回ListString而非Set因为执行顺序严格依赖拓扑序。曾有同事用HashSet导致节点乱序执行LLM还没返回就去调用“格式化结果”节点结果空指针。4.4 第四步实现流式节点模板25分钟以“调用OpenAI API”为例写OpenAiNodepublic class OpenAiNode implements StreamableNode { private final OkHttpClient client new OkHttpClient.Builder() .connectTimeout(30, TimeUnit.SECONDS) .readTimeout(60, TimeUnit.SECONDS) .build(); Override public Object execute(Context context) throws Exception { // 1. 从context取输入 String prompt (String) context.get(userInput); // 2. 构建SSE请求 Request request new Request.Builder() .url(https://api.openai.com/v1/chat/completions) .post(buildRequestBody(prompt)) .header(Authorization, Bearer System.getenv(OPENAI_KEY)) .build(); // 3. 发起请求传入tokenHandler EventSource eventSource new EventSource(request, client, new TokenEventHandler(context)); // 自定义处理器 eventSource.start(); return streaming_started; // 占位符实际输出由tokenHandler处理 } Override public void subscribe(ConsumerString tokenHandler) { this.tokenHandler tokenHandler; // 存起来供EventSource回调 } // 内部类TokenEventHandler实现EventSource.Listener private class TokenEventHandler implements EventSource.Listener { Override public void onEvent(Event event) { if (message.equals(event.name())) { String token parseToken(event.data()); tokenHandler.accept(token); // 关键推给引擎 context.updateNodeStatus(openai, new NodeStatus( openai, State.STREAMING, Instant.now(), token.substring(0, Math.min(50, token.length())))); } } } }这里tokenHandler.accept(token)就是引擎与节点的契约——引擎拿到token后立刻更新状态并透传给网关层的SseEmitter。4.5 第五步网关层集成20分钟Spring Boot ControllerRestController public class WorkflowController { Autowired private Engine engine; PostMapping(/workflow/start) public ResponseEntityString start(RequestBody WorkflowRequest request) { String workflowId UUID.randomUUID().toString(); Context context new Context(); // 初始化上下文 context.put(userInput, request.getInput()); // 异步执行避免阻塞 CompletableFuture.runAsync(() - engine.run(request.getWorkflow(), context)); return ResponseEntity.ok(workflowId); } GetMapping(value /workflow/{id}/stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter stream(PathVariable String id) { SseEmitter emitter new SseEmitter(30000L); // 30秒超时 emitter.onTimeout(() - emitter.complete()); emitter.onError(throwable - emitter.complete()); // 将emitter的send方法包装成ConsumerString ConsumerString tokenHandler token - { try { emitter.send(SseEmitter.event().name(token).data(token)); } catch (IOException e) { emitter.complete(); } }; // 把tokenHandler注入到Context需改造Context支持动态注册 context.registerTokenHandler(id, tokenHandler); return emitter; } }关键点context.registerTokenHandler()让引擎在执行节点时能拿到正确的tokenHandler——因为一个Workflow可能有多个流式节点每个节点的token要推到不同前端连接。5. 常见问题与排查技巧实录那些文档里不会写的坑5.1 流式输出“卡顿”问题排查树当用户反馈“流式输出停了几秒才继续”按此顺序排查检查项命令/方法预期结果网关层缓冲curl -v http://localhost:8080/workflow/abc/stream响应头含Cache-Control: no-cache且每秒至少一个data:事件OkHttp连接池jstack查看线程栈搜EventSource线程状态为RUNNABLE无WAITINGLLM API限流查OpenAI Dashboard的Rate limit usage当前RPM未超配额如40K TPM对应200 RPM节点线程池饱和jstat -gc pid看GCT时间Full GC频率1次/小时GCT100ms前端EventSource重连浏览器Console搜eventsource无Failed to load resource错误重连间隔呈指数退避1s, 2s, 4s...我们遇到过最诡异的卡顿前端EventSource收到token后DOM更新慢。根源是textContent token触发了浏览器重排reflow。解决方案用document.createDocumentFragment()批量插入性能提升10倍。5.2 状态丢失的“幽灵bug”现象节点执行完statusHistory里找不到COMPLETED状态。原因通常是异常吞没节点execute()里try-catch了所有异常但没调用context.updateNodeStatus(..., State.FAILED)线程切换节点用了CompletableFuture异步执行但updateNodeStatus()在子线程调用而Context的statusHistory是ConcurrentHashMap没问题但data里的nodeStatus_xxx是普通Map导致状态没刷进去Spring代理干扰如果节点是Spring BeanAsync方法调用会走CGLIB代理this.execute()变成代理对象调用context参数可能被篡改。解决强制要求所有节点execute()方法必须是public且不被Spring代理加Scope(prototype)或用ApplicationContext.getBean()手动获取。5.3 内存泄漏的“静默杀手”Agent工作流最大的内存杀手不是大对象而是闭包引用。典型场景// 错误示范lambda捕获了整个Context node.subscribe(token - { context.put(lastToken, token); // Context被lambda持有无法GC updateUi(token); });正确做法// 正确只捕获必要字段 String nodeId context.getNodeId(); ConsumerString handler token - { // 用弱引用或局部变量 updateUi(nodeId, token); }; node.subscribe(handler);我们用jmap -histo pid | grep Context定期检查发现Context实例数持续增长立刻查lambda使用点。5.4 并发下的状态竞争当多个用户同时启动同一WorkflowContext是隔离的但Workflow定义是共享的。问题出在Workflow的edges是HashMap并发读写会死循环JDK7 HashMap扩容时Engine.topologicalSort()如果缓存了排序结果多线程访问同一Workflow对象会竞争。解决方案Workflow类所有集合字段用Collections.unmodifiableXXX()包装topologicalSort()每次执行都重新计算不缓存排序O(VE)V100时开销可忽略NodeDefinition的config字段用String而非Map避免JSON解析时的并发问题。5.5 “流式输出乱序”的真相用户看到token顺序是Hello, world,!但前端拼出来是 world!Hello。这不是网络问题而是OkHttp EventSource的onEvent()回调不是按序的它用ExecutorService分发事件线程调度导致乱序解决方案在TokenEventHandler.onEvent()里加序号节点输出时带上index:123前端按序号排序再拼接。我们实测不加序号1000次请求中有3%乱序加序号后乱序率为0。6. 进阶实战如何把这套引擎塞进现有项目6.1 与Spring Boot无缝集成不用改Spring Boot启动类只需写EngineConfiguration配置类Configuration public class EngineConfiguration { Bean Scope(prototype) // 每次getBean都是新实例 public Engine engine() { return new Engine(); } Bean public Context context() { return new Context(); // 默认构造 } }在Controller里Autowired private Engine engine;直接调用用EventListener监听Context的状态变更事件自动推送到Redis Pub/Sub供其他服务订阅。这样既享受Spring的DI便利又保持引擎的轻量独立。6.2 对接Dify/Coze工作流Dify导出的JSON结构是{ nodes: [{id:llm,type:llm,config:{model:qwen}}, {id:output,type:output}], edges: [{source:llm,target:output}] }写个DifyWorkflowConverternodes数组转成Workflow.nodestype映射到Java类名如llm→com.example.node.OpenAiNodeedges转成Workflow.edgesconfig字段转成NodeDefinition.config字符串。一行代码接入Workflow workflow new DifyWorkflowConverter().convert(difyJson); engine.run(workflow, context);6.3 性能压测与调优清单用JMeter模拟1000并发瓶颈通常在IO线程池ioExecutor默认newCachedThreadPool最大线程数无上限导致CPU 100%。改为newFixedThreadPool(50)LLM调用超时OpenAI默认timeout 60秒但用户等待10秒就会取消。在OpenAiNode.execute()里加TimeoutFuture包装SSE连接数Tomcat默认maxConnections2001000并发会拒绝连接。改server.tomcat.max-connections10000GC调优加JVM参数-XX:UseG1GC -XX:MaxGCPauseMillis200避免长时间STW打断流式输出。我们最终压测结果单机8C16G支撑3000并发SSE连接P95延迟800ms。6.4 安全加固要点Agent工作流的安全不是“防SQL注入”而是节点沙箱所有NodeExecutor实现类禁止Runtime.getRuntime().exec()、禁止Class.forName(com.sun.*)等敏感反射上下文隔离Context.data用Collections.unmodifiableMap()包装返回防止节点恶意修改流式输出过滤在tokenHandler里加正则过滤script标签防止XSS虽然前端应做但服务端双保险重试熔断retryScheduler调度重试时记录失败次数超过3次自动标记State.CANCELLED防止雪崩。这些不是可选项是上线前必须checklist。7. 最后分享一个真实教训别在节点里做“智能”上线前三天我们给“知识库检索”节点加了个“智能重试”逻辑如果第一次检索相关性0.5自动换关键词再试一次。结果发现重试导致平均延迟从1.2秒涨到3.8秒两次检索的token消耗翻倍成本激增更糟的是重试时Context被修改影响下游节点判断。最后砍掉所有“智能”回归朴素节点只做一件事调用API返回原始结果“是否重试”由引擎根据NodeStatus.error类型决定如TimeoutException重试IllegalArgumentException不重试“换关键词”交给上游Agent Planner不在工作流引擎里做决策。这就是Agent工作流的精髓引擎是高速公路不是导航仪。导航Planner决定走哪条路引擎Executor只负责把车开稳。把Planner逻辑塞进引擎就像在发动机里装GPS——车开不快还容易烧机油。现在回头看那套200行核心代码之所以能跑一年不重构正是因为守住了这条边界引擎只管“怎么执行”不管“执行什么”。if-else写法Low不是因为语法Low而是因为它把“业务逻辑”和“执行逻辑”搅在一起。而真正的专业是让每个模块各司其职像钟表齿轮一样严丝合缝地咬合——Java的强类型、线程模型、内存控制恰好是构建这种精密系统的最佳材料。