拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Higress api-workflow 插件实战:用 DAG 配置编排 AI API 工作流

Higress api-workflow 插件实战用 DAG 配置编排 AI API 工作流【免费下载链接】higress AI Gateway | AI Native API Gateway项目地址: https://gitcode.com/GitHub_Trending/hi/higressapi-workflow 是 Higress 提供的一个可编排的 Wasm 插件它把一次请求处理抽象成一张由配置定义的 DAG有向无环图支持并行调用多个 API、按条件分支、汇聚结果并构造下一个请求非常适合在 AI Gateway 场景中编排 LLM、Embedding、Rerank、图像/音频等多模型服务的调用链路。读完本文你将掌握该插件完整的配置语法、控制流与数据流机制并能独立写出一个可运行的多节点工作流。功能概览为什么需要API 工作流在 AI 网关场景下一个业务请求往往不是一次简单转发而是需要串联多个模型服务先用 Embedding 服务把用户问题向量化再查询向量库最后把结果交给 LLM 生成回答。如果把这些调用逻辑写在业务代码里每新增一个环节都要改代码、重新发布而把调哪些接口、先后顺序、如何分支、数据怎么流动抽象成一张 DAG 配置就能让网关在不改代码的情况下完成多 API 编排。api-workflow插件正是为此设计根据配置定义生成 DAG 并执行工作流功能说明。它把工作流抽象成 DAG 配置文件配合控制流决定哪些步骤执行和数据流决定数据如何拼接进请求实现流程编排与请求构造。从源码看插件在初始化时通过wrapper.SetCtx注册了配置解析函数与请求处理函数main.go工作流在onHttpRequestBody阶段被触发整体上是一个基于 Proxy-Wasm 的 Go 插件位于plugins/wasm-go/extensions/api-workflow目录。下图展示了插件从接收到外部请求、解析 DAG 配置、执行工作流、调用各类 APILLM、rerank、embedding、audio、Tool 等到最后根据target是end返回最后一个操作的返回值还是continue放行到下一个插件结束的完整过程配置结构总览插件配置由两个顶层字段组成workflow必填DAG 的定义与env选填环境变量。名称数据类型填写要求默认值描述备注workflowobject必填-DAG 的定义-envobject选填-一些环境变量-env对象的配置字段名称数据类型填写要求默认值描述备注timeoutint选填5000每次请求的过期时间单位是毫秒(ms)max_depthint选填100工作流最大迭代次数-在源码中这两个默认值被定义为常量DefaultMaxDepth 100与DefaultTimeout 5000main.goparseConfig解析时若读到 0 则回退到默认值main.go。max_depth同时充当两个作用一是递归深度上限防止工作流死循环recursive中depth config.Env.MaxDepth直接报错见 main.go二是每次 httpCall 的超时基数实际超时为max_depth * timeout见 main.go。workflow对象的配置字段名称数据类型填写要求默认值描述备注nodesarray of node object选填-DAG 的定义的节点-edgesarray of edge object必填-DAG 的定义的边-edge对象的配置字段名称数据类型填写要求默认值描述sourcestring必填-上一步的操作必须是定义的 node 的 name或者初始化工作流的 starttargetstring必填-当前的操作必须是定义的 node 的 name或者结束工作流的关键字 end / continueconditionalstring选填-这一步是否执行的判断条件node对象的配置字段名称数据类型填写要求默认值描述备注namestring必填-node 名称全局唯一service_namestring必填-higress 配置的服务名称-service_portint选填80higress 配置的服务端口-service_domainstring选填-higress 配置的服务 domain-service_pathstring必填-请求的 path-service_headersarray of header object选填-请求的头-service_body_replace_keysarray of bodyReplaceKeyPair object选填-请求 body 模板替换键值对用来构造请求如果为空则直接使用 service_body_tmpl 请求service_body_tmplstring选填-请求的 body 模板-service_methodstring必填-请求的方法GETPOSTheader对象的配置字段名称数据类型填写要求默认值描述备注keystring必填-头文件的 key-valuestring必填-头文件的 value-bodyReplaceKeyPair对象的配置字段名称数据类型填写要求默认值描述备注fromstring必填-描述数据从哪获得-tostring必填-描述数据最后放到那-对应地PluginConfig、Env、Workflow、Edge、Node、BodyReplaceKeyPair、ServiceHeader等结构体定义在 workflow/workflow.goparseConfig会逐项校验必填字段如source、target、name、service_name、service_method为空都会直接报错见 main.go。值得注意的两个细节service_port为 0 时若service_name以.static结尾则默认使用 80static 类型服务否则报错service_path为空时默认取/。工作流骨架用边Edge描述 DAG 编排edges描述操作如何编排每条边表示从 source 到 target 的一次执行关系。以下是一个典型的边定义样例edges: - source: start target: A - source: start target: B - source: start target: C - source: A target: D - source: B target: D - source: C target: D - source: D target: end conditional: gt {{D||check}} 0.9 - source: D target: E conditional: lt {{D||check}} 0.9 - source: E target: end这张图表达的是start并行派出 A、B、C 三个节点三者都完成后汇聚到 DD 根据自身返回值check字段的大小决定走end 0.9还是先执行 E 0.9再结束。在插件内部target有三个特殊语义workflow/workflow.go 中的IsEnd/IsContinue当target为节点name执行该节点的操作当target为end直接返回 source 的结果结束工作流recursive中调用proxywasm.SendHttpResponse把最终 body 回给客户端见 main.go当target为continue结束工作流将请求放行到下一个插件调用proxywasm.ResumeHttpRequest()见 main.go。多节点并行汇聚fan-in依赖入度计数机制initWorkflowExecStatus在配置解析阶段为每个节点统计其入边数量即依赖的前驱数执行时每个前驱完成后将目标节点计数减 1只有计数归零才真正执行该节点main.go 与 main.go。这就是 D 节点能等到 ABC 数据都就位的原因。WorkflowExecStatus以map[string]int形式存放在请求上下文中key 为workflowExecStatus见 main.go。控制流用 conditional 表达式决定分支走向当某条edge的conditional定义不为空时插件执行到该边会根据表达式判断是否执行这一步判断为否时跳过该分支。表达式支持使用参数参数用{{xxx}}标注具体定义见下文数据流支持的比较操作符如下操作符语义示例eq arg1 arg2arg1 arg2 时为 true不只是数字支持 stringeq 1 1lt arg1 arg2arg1 arg2 时为 truelt 1.1 2le arg1 arg2arg1 arg2 时为 truele 1 2gt arg1 arg2arg1 arg2 时为 truegt 2 1ge arg1 arg2arg1 arg2 时为 truege 2 1and arg1 arg2arg1 arg2and (eq 1 1) (lt 2 3)or arg1 arg2arg1 || arg2or (eq 1 2) (lt 1 3)contain arg1 arg2arg1 包含 arg2 时为 truecontain helloworld world支持and、or的嵌套例如and (eq 1 1) (or (contain hello hi) (lt 1 2))。在源码中这些操作符定义在 utils/conditional.goExecConditionalStr通过正则先提取最内层括号内的原子表达式递归求值再用结果替换回外层表达式最终对三元组操作符、arg1、arg2分派执行utils/conditional.go。对应的单元测试覆盖了数字/字符串/浮点比较、布尔与或、contain、非法输入以及多层嵌套等大量用例utils/conditional_test.go例如or (eq 1 2) (and (eq 1 1) (gt 2 3))应得到false。条件判断的执行链路是Edge.IsPass先把conditional中的模板变量{{str1||str2}}替换为实际数据WrapperDataByTmplStr再调用ExecConditional求值注意返回值取反!ok表示不通过则跳过语义上 conditional 为真才继续执行workflow/workflow.go。数据流请求 body 的模板构造与上下文传递进入插件的数据request body会按如下机制流转每个节点执行前会根据该节点配置的service_body_tmpl构造模板 JSON与service_body_replace_keys替换键值对构造请求 body执行后的结果以节点 name 为 key存入请求上下文ctx.SetContext(edge.Target, responseBody)见 main.go供后续节点取用。目前只支持 JSON 格式的数据。模板与变量{{str1||str2}}在工作流配置文件中edge.conditional支持模板和变量方便根据数据流的数据构建判断表达式。变量使用{{str1||str2}}包裹用||分隔str1代表使用哪个 node 的输出数据str2代表如何取数据过滤表达式基于GJSON PATH语法提取字符串all代表全都要。例如conditional: lt {{D||check}} 0.9若 node D 的返回值是{check: 0.99}则解析后的表达式为lt 0.99 0.9即判断 0.99 是否小于 0.9。service_body_tmpl与service_body_replace_keys这组配置用来构造请求 bodyservice_body_tmpl是模板 JSONservice_body_replace_keys描述如何填充模板 JSON一个 object 数组from数据从哪里来使用str1||str2字符串str1代表使用哪个 node 的执行返回数据str2代表如何取数据表达式基于 GJSON PATH 语法提取字符串to数据放到哪表达式基于 GJSON PATH 语法描述填充位置内部使用 sjson 拼接 JSON填充进service_body_tmpl模板 JSON 里。当service_body_replace_keys为空时代表直接发送service_body_tmpl。示例service_body_tmpl: embeddings: result: msg: sk: sk-xxxxxx service_body_replace_keys: - to: embeddings.result from: A||output.embeddings.0.embedding - to: msg from: B||all假设 A 节点的输出是{embeddings:{output:{embeddings:[{embedding:[0.014398524595686043],text_index:0}]},usage:{total_tokens:12},request_id:2a5229bc-53d9-91ca-bce2-00ae5e01a1d3}}B 节点的输出是[higress项目主仓库的github地址是什么]那么根据模板和替换键值对构造出的 request body 为{embeddings:{result:[0.014398524595686043,......]},msg:[higress项目主仓库的github地址是什么],sk:sk-xxxxxx}从源码看WrapperDataByTmplStrAndKeys会解析from中的||从上下文取出对应节点的原始 body再用 GJSON 按路径取值最后用 sjson 的SetRaw把取到的值原样写入to指定的位置workflow/workflow.go当from的第二段为all时整段数据整体填入。若路径不存在会返回明确的错误信息最终该节点以 500 结束工作流。节点定义封装 httpCall 的操作单元node是工作流中具体执行的单元它封装了 httpCall提供 HTTP 访问能力以获取各种 API 的能力request body 支持自主构建。以下是一个调用官方 text-embedding-v2 模型的节点样例nodes: - name: A service_domain: dashscope.aliyuncs.com service_name: dashscope service_port: 443 service_path: /api/v1/services/embeddings/text-embedding/text-embedding service_method: POST service_body_tmpl: model: text-embedding-v2 input: texts: parameters: text_type: query service_body_replace_keys: - from: start||messages.#(roleuser)#.content to: input.texts service_headers: - key: Authorization value: Bearer sk-b98f462xxxxxxxx - key: Content-Type value: application/json这里展示了几处关键用法start||messages.#(roleuser)#.content表示从进入插件的原始请求 body即 start 节点的输出中用 GJSON 查询语法按roleuser过滤出 messages 数组里用户消息的 content填入模板的input.textsservice_domain指定直连域名service_name指定 Higress 中配置的服务名。在wrapperNodeTask中节点信息会被封装为 FQDN ClusterHost 为service_domainFQDN 为service_namePort 为service_port再依次封装 body、Method、Path 与 Headersworkflow/workflow.go。完整实战并行采集 → 汇聚判断 → 条件落库下面给出 README 中的完整示例从三个节点 A、B、C 获取信息等数据都就位后再执行 D并根据 D 的输出判断是否需要执行 E 还是直接结束。其 DAG 结构如下start 的返回值请求插件的 body{ model:qwen-7b-chat-xft, frequency_penalty:0, max_tokens:800, stream:false, messages: [{role:user,content:higress项目主仓库的github地址是什么}], presence_penalty:0,temperature:0.7,top_p:0.95 }A 的返回值是{ output:{ embeddings: [ {text_index: 0, embedding: [-0.006929283495992422,-0.005336422007530928]}, {text_index: 1, embedding: [-0.006929283495992422,-0.005336422007530928]}, {text_index: 2, embedding: [-0.006929283495992422,-0.005336422007530928]}, {text_index: 3, embedding: [-0.006929283495992422,-0.005336422007530928]} ] }, usage:{total_tokens:12}, request_id:d89c06fb-46a1-47b6-acb9-bfb17f814969 }B 的返回值是{llm:this is b}C 的返回值是{get: this is c}D 的返回值是{check: 0.99, llm:{}}E 的返回值是{save: ok, date:{}}。这个工作流的完整配置文件如下env: max_depth: 100 timeout: 3000 workflow: edges: - source: start target: A - source: start target: B - source: start target: C - source: A target: D - source: B target: D - source: C target: D - source: D target: end conditional: lt {{D||check}} 0.9 - source: D target: E conditional: gt {{D||check}} 0.9 - source: E target: end nodes: - name: A service_domain: dashscope.aliyuncs.com service_name: dashscope service_port: 443 service_path: /api/v1/services/embeddings/text-embedding/text-embedding service_method: POST service_body_tmpl: model: text-embedding-v2 input: texts: parameters: text_type: query service_body_replace_keys: - from: start||messages.#(roleuser)#.content to: input.texts service_headers: - key: Authorization value: Bearer sk-b98f462xxxxxxxx - key: Content-Type value: application/json - name: B service_body_tmpl: embeddings: default msg: default request body sk: sk-xxxxxx service_body_replace_keys: service_headers: - key: AK value: ak-xxxxxxxxxxxxxxxxxxxx - key: Content-Type value: application/json service_method: POST service_name: whoai.static service_path: /llm service_port: 80 - name: C service_method: GET service_name: whoai.static service_path: /get service_port: 80 - name: D service_headers: service_method: POST service_name: whoai.static service_path: /check_cache service_port: 80 service_body_tmpl: A_result: B_result: C_result: service_body_replace_keys: - from: A||output.embeddings.0.embedding.0 to: A_result - from: B||llm to: B_result - from: C||get to: C_result - name: E service_method: POST service_name: whoai.static service_path: /save_cache service_port: 80 service_body_tmpl: save: service_body_replace_keys: - from: D||llm to: save执行请求curl -v 127.0.0.1:8080 -H Accept: application/json, text/event-stream -H Content-Type: application/json --data-raw {model:qwen-7b-chat-xft,frequency_penalty:0,max_tokens:800,stream:false,messages:[{role:user,content:higress项目主仓库的github地址是什么}],presence_penalty:0,temperature:0.7,top_p:0.95}执行后的简略 debug 日志可以看到工作流等到前置的 ABC 流程执行完毕后根据返回值构建了 D 的 body{A_result:0.007155838584362588,B_result:this is b,C_result:this is c}执行 D 后根据 D 的返回值{check: 0.99, llm:{}}进行条件判断最终继续执行了 Egt 0.99 0.9然后结束流程[api-workflow] workflow exec task,source is start,target is A, body is {input:{texts:[higress项目主仓库的github地址是什么]},model:text-embedding-v2,parameters:{text_type:query}},header is [[Authorization Bearer sk-b98f4628125xxxxxxxxxxxxxxxx] [Content-Type application/json]] [api-workflow] workflow exec task,source is start,target is B, body is {embeddings:default,msg:default request body,sk:sk-xxxxxx},header is [[AK ak-xxxxxxxxxxxxxxxxxxxx] [Content-Type application/json]] [api-workflow] workflow exec task,source is start,target is C, body is ,header is [] [api-workflow] source is B,target is D,status is map[A:0 B:0 C:0 D:2 E:1] [api-workflow] source is C,target is D,status is map[A:0 B:0 C:0 D:1 E:1] [api-workflow] source is A,target is D,status is map[A:0 B:0 C:0 D:0 E:1] [api-workflow] workflow exec task,source is A,target is D, body is,header is [] [api-workflow] source is D,target is end,workflow is pass [api-workflow] source is D,target is E,status is map[A:0 B:0 C:0 D:0 E:0] [api-workflow] workflow exec task,source is D,target is E, body is {save:{\A_result\:0.007155838584362588,\B_result\:\this is b\,\C_result\:\this is c\}},header is [] [api-workflow] source is E,target is end,workflow is end从日志中可以直观看到入度计数机制source is B,target is D,status is map[A:0 B:0 C:0 D:2 E:1]中 D 的计数从 2 递减到 0只有当 A、B、C 三个前驱全部完成D 才真正被触发执行。这与你配置的边关系一一对应也是排查编排问题时最重要的调试线索。测试与验证源码中的行为证据插件自带的测试用例可以直接验证上述机制main_test.go基本工作流start → A → end验证onHttpRequestBody返回ActionPause等待外部 HTTP 调用响应成功后返回 200条件分支A返回{score: 0.8}时命中gt {{A||score}} 0.5分支直接到end并行执行A、B、C 三个请求分别返回后D 作为汇聚节点接收三者结果最终返回 200验证 fan-in 计数逻辑continue 场景执行完 A 后 target 为continue最终GetHttpStreamAction()返回ActionContinue即请求被放行到下一个插件。条件表达式求值本身则有 utils/conditional_test.go 提供逐操作符的断言与嵌套用例可作为编写 conditional 时的参考语法清单。快速上手与部署提示api-workflow 是一个标准的 Higress Wasm Go 插件插件目录下提供了 Dockerfile 用于构建 Wasm 镜像构建与发布流程可参考插件整体说明plugins/wasm-go/README.md以及仓库根目录的 Makefile 中的相关 target。部署时通过 Higress 的 WasmPlugin 资源将插件挂载到目标域名/路由上并以 JSON 格式提供上文所述的workflow、env配置配置参考样例见 samples/wasmplugin 目录下的 yaml 写法。需要特别说明的是插件在请求体阶段ProcessRequestBody执行因此默认只对带 body 的 POST 类请求生效工作流执行过程中使用ActionPause挂起请求等待异步 HTTP 回调这是它能够编排多节点串并行调用的底层原理见 main.go。小结api-workflow 插件用一张 DAG 配置把多 API 编排从业务代码中解放出来edges定义编排拓扑conditional提供基于表达式与 GJSON 变量的分支控制service_body_tmplservice_body_replace_keys实现节点间 JSON 数据流的请求构造配合入度计数实现并行汇聚、end/continue两种收尾方式足以覆盖 AI 网关中诸如向量化 → 检索 → 生成、并行多模型打分 → 聚合判断 → 落库等典型编排诉求。需要深度定制时可直接阅读 workflow/workflow.go 与 main.go 了解执行细节并用 main_test.go 中的测试模式快速验证自己的配置。【免费下载链接】higress AI Gateway | AI Native API Gateway项目地址: https://gitcode.com/GitHub_Trending/hi/higress创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门