十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

智能体工作流引擎的设计与实现:从DAG调度到节点编排实践

智能体工作流引擎的设计与实现:从DAG调度到节点编排实践 做智能体应用一段时间的人基本都会撞上同一个问题Agent 在简单问答场景下很聪明一旦进入真实业务流程——比如收集工单信息→判断问题类型→调用接口查询→走审批流程→回执结果单靠模型自由发挥就完全不可控了。MaxKB4j 智能体平台里最核心的组件之一就是节点工作流引擎它把 LLM 调用、条件判断、工具执行、外部服务访问这些能力编排成一张有向无环图让整个智能体的行为从模糊不可预期变成确定可复现。这篇文章我会从技术原理和工程实现两个层面把这个引擎内部到底怎么设计的、执行链路怎么跑的、业务节点怎么扩展的完整讲一遍。适合正在做智能体平台、Agent 应用编排或者打算自研工作流引擎的后端工程师和架构师参考。1. 为什么需要自研工作流引擎从编排痛点说起1.1 单轮智能体演进到复杂编排的必然过程早期的智能体实现方式很简单一个Prompt一次LLM调用让模型自己决定调什么工具、按什么顺序调也就是现在常说的 ReAct 循环。这种模式在 Demo 阶段跑得很流畅因为场景足够窄工具数量少模型选错的概率低而且就算选错了改一版 Prompt 就能糊弄过去。一旦进入生产环境问题就冒出来了。我自己踩过几个非常典型的坑一是 Token 消耗不可控。一个工单流程如果完全靠模型自由决策它可能会反复重试、反复读同一个知识库文档、输出大段没有实际意义的中间推理单次任务的 Token 消耗能比确定性问题高出 3~5 倍。二是无法稳定复现。同一套 Prompt 在模型版本更新之后行为可能直接变化。比如原来能稳定调用某个 API升级之后模型开始自作主张改变参数顺序导致下游系统解析失败。三是外部系统状态难以对齐。业务系统是强状态的一个审批节点要么通过、要么驳回不存在也许通过的状态。但模型生成本质上是概率的它没办法保证每一次都输出合法的流程控制信号。工作流引擎解决的就是把流程确定性交给代码、把内容生成能力交给模型这个问题。在这个思路下智能体不再是一个大循环而是拆成一个个清晰的节点信息收集节点、意图判断节点、业务查询节点、人工确认节点、结果生成节点。节点之间的流转规则由引擎控制每个节点内部可以调用 LLM、可以调普通代码、可以访问外部服务。1.2 与现成开源方案对比后为什么还是选择自研提到智能体工作流业内最常被拿来比较的就是 Dify 这类开源智能体平台。Dify 的工作流编排界面做得非常成熟可视化拖拽、节点类型丰富、社区生态也很好。那为什么我们做 MaxKB4j 的时候还是要自己写一套引擎而不是直接把 Dify 嵌进来核心原因有三个。第一是部署形态的差异。Dify 是一个完整的平台前后端、数据库、任务队列一套全带上。但 MaxKB4j 的定位是作为一个 SDK/基础组件嵌入到我们自己已有的业务系统里而不是重新部署一套独立的平台。平台级方案和组件级方案在集成成本上完全不是一个量级。第二是节点扩展机制的约束。Dify 的工作流节点是平台内置好的你想加一个私有协议对接节点要么改 Dify 源码要么通过它暴露的插件机制但插件机制的边界是平台定的灵活度有限。我们的业务场景里有大量存量系统需要对接协议五花八门如果节点扩展权不在自己手里后面每个新对接都要走平台升级流程太被动。第三是执行控制粒度的差异。作为平台Dify 对单个工作流实例的控制比较粗粒度运行、停止、看日志。但我们希望引擎能提供更细的执行控制能力比如节点级别的超时设置、条件分支的动态解析、运行轨迹的分层记录这些需求在没有源码控制权的前提下很难做得彻底。简单总结一下我们的取舍逻辑Dify 适合我要快速搭建一套独立可用的智能体应用自研引擎适合我已经有一个业务系统需要把工作流能力作为模块嵌进去前者快后者灵活没有谁更好只有谁更匹配1.3 引擎需要具备的四个底层能力在动手设计之前我先把引擎必须要有的能力列了出来。这四个能力缺一个引擎都撑不起真实业务。第一个是图的描述与校验能力。工作流本质上就是一张有向无环图DAG节点是图的顶点连线是有向边。引擎必须能解析图定义、检查环的存在、识别孤立节点、判断边的连接是否合法。图都描述不清楚执行阶段必然出乱子。第二个是节点的注册与生命周期管理能力。每个节点不是硬编码在引擎里的而是通过注册机制动态加入。引擎只负责调度节点负责具体执行。注册机制做得够好新的业务节点才能低成本接入。第三个是数据流转与作用域隔离能力。上游节点的输出怎么传给下游节点、全局变量和局部变量怎么区分、多个并行分支之间数据会不会互相污染这些都要有明确的约定。这一步往往最容易被忽略但恰恰是后期问题最多的。第四个是运行状态的可观察性能力。工作流跑起来之后到底是哪个节点在执行、执行了多久、成功还是失败、失败在哪一步、当时的输入输出是什么全部要能查询。如果没有可观察性生产环境出了问题就只能靠猜。这四个能力对应到工程上分别就是图模型模块、节点注册中心、上下文数据总线、执行追踪模块。下面几节我会挨个展开讲。2. 引擎的宏观架构DAG 驱动的分层设计2.1 定义层、执行层、运行时的三层职责划分MaxKB4j 引擎在代码结构上分了三个清晰的层级每一层只干自己的事依赖方向是单向的定义层被执行层依赖执行层被运行时依赖反向不引用。定义层负责描述工作流长什么样。一个工作流定义就是一份 JSON 结构包含节点列表、连线列表、变量定义、全局配置。这一层是纯数据不含任何执行逻辑。好处是工作流定义可以被序列化到数据库、可以被管理后台编辑、可以被版本管理也能方便地导出和导入。执行层负责解析定义、构建内存图对象、进行调度。它拿到定义层的数据后会做三件事第一把 JSON 反序列化为对象模型第二校验图的合法性第三构建就绪队列和执行上下文。执行层不关心节点内部具体做了什么只关心节点的生命周期状态。运行时层负责具体执行节点代码、记录执行日志、触发事件回调、维护执行实例的状态。比如一个调用LLM节点定义层只知道它的类型是 llm、参数是模型名称和提示词模板执行层把它放进调度队列真正调 LLM API、处理返回结果的代码是运行时层的节点执行器。这个分层让每一层都能独立测试。我们实际开发中最常用的手段就是定义层用 JSON Schema 校验执行层用假节点做调度单测运行时层针对每个具体节点做集成测试。三层面向不同维度的验证问题定位会快很多。2.2 节点模型定义与执行分离的可扩展抽象节点的设计是整个引擎的灵魂。MaxKB4j 的节点模型把节点长什么样和节点干什么事拆成了两个对象NodeDefinition 和 NodeExecutor。NodeDefinition 是节点的静态描述包含基本信息。它被用来做工作流的可视化展示、参数表单渲染、JSON Schema 校验。NodeExecutor 是节点的运行时实现包含真正的执行逻辑。我用一个简化版的 Java 接口来说明这个设计// 节点的静态定义 public class NodeDefinition { private String type; // 节点类型如 llm、http、condition private String name; // 节点展示名称 private boolean isStart; // 是否为开始节点 private boolean isEnd; // 是否为结束节点 private ListFieldSchema inputs; // 输入参数的 JSON Schema } // 节点的运行时执行器 public interface NodeExecutor { NodeResult execute(NodeContext context) throws Exception; }NodeContext 里有引擎注入的各类依赖包括当前工作流实例 ID、节点实例 ID全局变量查询接口上游节点的输出结果日志输出接口事件发布器这样做最大的优势是扩展性极强。新增一个节点类型只需要实现 NodeExecutor 接口然后在注册中心登记类型名称引擎就能识别并调度它。业务系统对接什么新协议就加什么新节点完全不用动引擎主流程代码。2.3 图数据结构的落地邻接表与入度控制图的存储我斟酌过两种方式邻接矩阵和邻接表。工作流图的节点数一般不会很多几十个就算大型流程了所以理论上两者都行。但实际我选择了邻接表因为工作流图天然是稀疏的绝大多数顶点只连接了少数几个邻居用矩阵浪费空间不说遍历也不直观。定义层里的连线结构很简单{ edges: [ { source: node_start, target: node_collect_info }, { source: node_collect_info, target: node_condition } ] }引擎解析时会把边转成两类数据结构outputEdges记录每个节点的所有出边用于执行完一个节点后找下一批候选节点inDegree记录每个节点的入度数量用于调度判活有了入度信息调度逻辑就很清晰了入度为 0 的节点就是可以立即执行的节点。节点执行完成之后把它所有出边对应的下游节点入度减 1如果减到 0就把下游节点加入就绪队列。图校验方面最重要的就是检测环。实现上我用的是拓扑排序能完成拓扑排序的图是无环的反之就有环直接拒绝执行。这个校验必须在引擎启动加载工作流定义时就做不能等到运行时。试想一个环在工作流里意味着什么——A 节点等 BB 节点等 CC 节点又等 A整个执行会直接卡死而且这种问题在界面上还特别难发现。3. 节点生命周期与执行器机制3.1 节点的五状态状态机与容错转移每个节点实例在引擎里都会经历状态变化。MaxKB4j 的节点实例有五个状态PENDING、RUNNING、SUCCEEDED、FAILED、SKIPPED。节点刚被调度器选中但还没开始执行时是 PENDING。执行器开始跑 execute 方法的那一刻变成 RUNNING。执行方法正常返回变成 SUCCEEDED。抛异常或者超时变成 FAILED。条件分支判断为不满足的路径上的节点被标记为 SKIPPED不执行任何逻辑。状态转移规则如下当前状态可转移状态触发条件PENDINGRUNNING调度器分配线程开始执行RUNNINGSUCCEEDEDexecute 方法正常返回RUNNINGFAILED执行抛异常或触发超时PENDINGSKIPPED上游分支条件不满足FAILED-终态等待人工处理或工作流终止这个状态机看起来简单但工程上有个细节很关键FAILED 之后节点要不要自动重试。我做了可配置的策略因为并不是所有失败都值得重试。网络抖动类错误可以重试 2~3 次带退避间隔参数错误、权限错误这类确定性错误重试一万次也没用直接标记 FAILED 并触发工作流失败回调。3.2 上下文传递与变量作用域约定工作流引擎的上下文设计直接影响到节点的编写复杂度和排错难度。MaxKB4j 在上下文设计上遵循一个原则全局变量只读局部变量显式引用。全局变量放在一个 ImmutableMap 里引擎初始化工作流的时候注入进去。业务系统触发工作流时传入的初始参数、上一个节点的全部输出都会被固化在全局变量中。局部变量则是通过变量引用字符串来实现的。比如 A 节点输出了一个字段叫 orderId下游 B 节点想读取就配置${nodes.node_a.orderId}。引擎在执行 B 节点前会去解析这个引用从 A 节点的输出里取值再传入 B 节点的输入。如果引用路径不存在就直接报错绝不静默返回 null。为什么不做隐式的所有上游输出都注入所有下游节点这种设计因为副作用太大。一是并行分支场景下两个分支的节点输出谁先谁后没有确定性下游不知道拿谁的值二是大量数据堆积在上下文里单个节点处理不过来时内存会爆。显式引用虽然写起来啰嗦一点但可预测、可定位、可审计。3.3 分支、循环与并行三种控制流的实现控制流是工作流引擎和普通任务队列最大的区别。MaxKB4j 的引擎实现了三种控制流模式。条件分支是第一种。典型场景是如果用户输入的订单号以 10 开头走退款流程否则走人工客服。实现方式是分支节点的输出连接了多条出边每条边上配置一个条件表达式。节点执行完成后引擎对每条出边的表达式求值值为 true 的边被激活下游节点入度减 1值为 false 的边不激活下游节点直接标记为 SKIPPED并把这种情况也记录到日志里。并行是第二种。一条出边同时连接多个下游节点时这些节点理论上可以同时执行。实现上我用了一个线程池每个下游节点提交为一个独立任务。等所有并行节点都完成之后汇合节点再继续往下走。汇合时要注意的点是分支里某些节点可能是 SKIPPED 状态汇合逻辑应该只等待真正执行了的节点避免死等不存在的执行结果。循环是第三种。虽然主图是 DAG不允许出现环但真实业务里循环调用某个 API 直到拿到最终结果这种需求大量存在。我的做法是提供一个专门的循环节点它内部维护一个 MAX_ITERATIONS 上限在节点内部做循环逻辑而整个图上依然是一个节点。这样既保证了 DAG 约束不会被破坏又满足了实际需求。4. 核心实现从定义到运行的完整链路4.1 工作流定义的 JSON Schema 设计工作流定义是引擎的输入它的 Schema 直接决定了系统的表达能力和易用性。我提供一份精简后的核心结构{ workflowId: wf_001, title: 售后退款处理流程, variables: { input: { orderId: string, userId: string } }, nodes: [ { id: node_start, type: start, name: 开始 }, { id: node_cond, type: condition, name: 是否退款订单, config: { expression: payload.refundable true } }, { id: node_http, type: http, name: 调用退款接口, config: { method: POST, url: ${env.refund_service_url}/api/refund, body: { orderId: ${nodes.node_start.output.orderId} } } } ], edges: [ { source: node_start, target: node_cond }, { source: node_cond, target: node_http } ] }Schema 设计上有三个细节值得说一下。一是 config 字段用了 OpenAPI 风格的对象结构而不是传统的扁平 key-value这样每个节点的参数都带上了类型信息引擎可以统一做法在校验阶段检查类型对不对。二是变量表达式我用了${}模板语法。在执行器拿到 config 值的时候引擎会先把所有包含${}的字符串做一次插值解析再传给节点。这样节点执行器自己不需要关心变量从哪来它拿到的已经是最终值实现大大简化。三是这段 JSON 里没有直接体现分支不满足时走到哪因为分支路径是通过 edge 上的条件表达式表达的而不是靠节点里配置 target 字段。这种设计让图画起来更像实际流程图可读性好很多。4.2 调度引擎的核心实现思路调度是整个引擎最核心的部分。我贴一段简化后的调度器核心逻辑public class WorkflowScheduler { private final ExecutorService executor; private final MapString, WorkflowNode nodeMap; private final MapString, Integer inDegree new ConcurrentHashMap(); private final QueueString readyQueue new ConcurrentLinkedQueue(); public void schedule(WorkflowGraph graph) { // 1. 初始化入度表 for (String nodeId : graph.getNodeIds()) { inDegree.put(nodeId, graph.getInDegree(nodeId)); if (inDegree.get(nodeId) 0) { readyQueue.offer(nodeId); } } // 2. 循环消费就绪队列 while (!readyQueue.isEmpty()) { String nodeId readyQueue.poll(); WorkflowNode node nodeMap.get(nodeId); executor.submit(() - runNode(node)); } } private void runNode(WorkflowNode node) { try { NodeResult result node.getExecutor().execute(node.getContext()); markSuccess(node.getId(), result); } catch (Exception e) { markFailed(node.getId(), e); return; } // 3. 节点成功后更新下游入度 for (WorkflowEdge edge : node.getOutputEdges()) { int remain inDegree.computeIfPresent(edge.getTarget(), (k, v) - v - 1); if (remain 0) { readyQueue.offer(edge.getTarget()); } } } }这里有一个非常容易踩坑的点并行场景下的入度减少必须做原子操作。如果两个上游节点同时执行完成同时去更新同一个下游节点的入度不用原子类的话就会丢更新导致入度错误地停留在 1下游节点永远等不到调度。我在第一版实现里就踩过这个坑Debug 了整整一天才定位到后来全部换成了 computeIfPresent 这个原子方法。线程池的参数也值得单独说。LLM 调用、HTTP 请求这类节点全都是 IO 密集型的线程池核心线程数如果设置成 CPU 核数IO 等待时整个调度就会卡住。我最后设置为核心线程数 16、最大线程数 64、队列容量 512核心线程空闲回收时间 60 秒。这个参数是根据我们平均工作流节点数 15~20 个、单个节点平均耗时 1~3 秒测算出来的能够支撑约 10 个工作流同时并行执行而不互相阻塞。4.3 条件分支中的表达式求值表达式引擎的选择我斟酌了很久。最初考虑直接用现成的 Spring Expression LanguageSpEL功能足够强大但问题也出在强大上——SpEL 支持调用任意 Java 方法一旦工作流定义是可以被业务人员编辑的表达式就变成了一个远程代码执行漏洞。后来我改成了自研的轻量表达式解析器只支持有限的语法子集字面量字符串、数字、布尔值、null变量引用如payload.refundable、nodes.node_http.output.code比较运算、!、、、、逻辑运算、||、!可选的一元取反这个解析器的实现并不复杂核心就是分词 递归下降解析 求值。分词器把表达式拆成 Token解析器按优先级构建抽象语法树求值时从上下文取变量值做运算。实际落地时有个性能优化点一个工作流可能被触发成千上万次如果每次执行都重新解析一遍表达式开销完全没必要。我引入了表达式缓存把解析好的 AST 缓存在 ConcurrentHashMap 里key 是表达式字符串本身。实测下来表达式求值的平均耗时从最初的 3 毫秒降到了 0.1 毫秒以内。5. 扩展实践如何接入一个新的业务节点5.1 实际步骤从需求到实现一个发送短信节点光讲原理不够我拿一个实际例子走一遍完整接入流程。假设现在业务方要加一个节点发送短信验证码。第一步实现 NodeExecutor。短信服务是一个已有的内部 RPC我只需要把参数取出来传进去public class SmsSendExecutor implements NodeExecutor { private final SmsClient smsClient; Override public NodeResult execute(NodeContext context) { String phone context.getConfig(phone); String content context.getConfig(content); String bizId smsClient.send(phone, content); return NodeResult.success(Map.of(bizId, bizId)); } }第二步注册节点。我在配置文件里登记节点类型 sms_send同时提供输入参数 Schema{ type: sms_send, name: 发送短信, inputs: [ { field: phone, type: string, required: true }, { field: content, type: string, required: true } ] }注册完成之后工作流编辑器里就能看到这个节点业务人员可以直接拖拽使用。整个过程不需要改引擎任何代码这就是节点注册机制的价值。5.2 可视化执行轨迹可观测性的实际落地工作流引擎如果只输出一堆日志生产环境根本没法用。我专门做了执行轨迹记录模块每次执行都会留下三条层级的数据。流程级别的记录整个工作流的开始时间、结束时间、耗时、最终状态。节点级别的记录每个节点的开始时间、结束时间、状态、出参摘要。事件级别的记录节点内部打点输出的关键事件比如 LLM 调用的请求参数和返回结果。这套设计在执行链路里的效果非常直观。工作流出问题的时候我只要打开执行记录列表找到那一条失败记录放大看是哪个节点挂了再看它当时的入参出参基本几分钟内就能定位问题。如果哪天模型返回格式变了导致解析失败执行记录里能清楚地看到 LLM 节点的输出内容和解析错误日志不需要去后端系统里翻原始请求。5.3 超时、并发与资源保护的工程取舍最后一个实操主题是工程的自我保护。工作流引擎不像普通接口一个流程可能跑几分钟如果中间有节点是恶意慢调用整个链路都会卡住。我做的第一个保护是节点级超时。每个节点在执行前会计算一个截止时间执行器线程用 CompletableFuture 包装到了截止时间还没返回就直接取消并标记 FAILED。HTTP 节点默认超时 30 秒LLM 节点默认 60 秒这些都是根据实际服务的 P99 延迟来标定的。第二个保护是图规模限制。单个工作流的节点数量上限设为 100 个边的数量上限 200 条。业务方反馈过我的流程需要 150 个节点怎么办我的回复是拆成多个子流程用子流程节点串联。因为图一旦超过这个规模可视化界面排布和引擎调度的复杂度都会急剧上升收益完全覆盖不了成本。第三个保护是并发工作流实例的数量限制。每个工作流定义都有各自的并发上限默认是 20 个同时执行。超过上限的触发请求直接返回当前排队中由业务侧决定是排队等还是提示用户稍后再试。这个限流措施防止了某个高频业务把整个引擎线程池打满。就我个人的实际体验来说做工作流引擎难度最大的反而不是调度算法本身而是给能力边界做减法。节点的语义粒度定得太粗节点之间逻辑耦合严重定得太细一个简单流程要画十几二十个节点业务方根本不想用。最终我找到的相对平衡点是一个节点应该完成一个有明确产出、可独立验收的动作比如调用一个 API、生成一段文本、读一次知识库而不是把多个动作硬捏在一个节点里。这样画流程图的人在界面上看到的内容和实际执行时发生的事情才是一致的排查问题的时候也才能对着流程一路顺下来。
返回列表