编程 YAML 编译成一张 LangGraph:Supervisor + Worker 多 Agent 平台的 7 条已知限制

2026-09-24 00:05:53

YAML 编译成一张 LangGraph:Supervisor + Worker 多 Agent 平台的 7 条已知限制

项目地址:github.com/diguike/book-langchain-agent,作者站点:inferloop.dev

目标是一个简化版的 Dify / n8n:用户用 YAML 描述工作流,平台在运行时把它编译成一张 LangGraph,然后按 Supervisor + Worker 模式执行。只做三件事:

  1. YAML 定义工作流——节点用什么模型、跑什么工具、怎么衔接,全部声明式描述;
  2. 运行时编译——读 YAML,生成一张 StateGraph.compile(),不改代码;
  3. Supervisor + Worker 拓扑——一个调度 Agent + 多个执行 Agent,调度结果用 Command({ goto, update }) 显式传递。

不做可视化拖拽编辑器、节点级权限、计费、多租户隔离。

架构与模型选型

YAML 工作流定义
-> GraphCompiler 运行时编译
-> CompiledGraph
-> Supervisor 节点(Claude Haiku 4.5)
-> handoff 到 Worker(research / write / review)
-> 结果回流 Supervisor
-> done 到 END

模型分工按任务复杂度切:Supervisor 用 Haiku 4.5,它只做「下一步谁来干」的分类决策;Worker 默认 Sonnet 4.6 做业务执行,YAML 里可以覆盖;需要重推理的 Worker 用 Opus 4.7,由 YAML 显式声明。

YAML schema

id: content-pipeline
name: 内容生产流水线
workers:
- id: researcher
model: anthropic:claude-sonnet-4-6
systemPrompt: |
你是调研员。根据用户给的主题,列 3-5 个关键事实和引用源。
tools: [web_search]
- id: writer
model: anthropic:claude-sonnet-4-6
systemPrompt: |
你是写手。基于调研结果写一篇 300-500 字的中文短文。
tools: []
- id: reviewer
model: anthropic:claude-opus-4-7
systemPrompt: |
你是审稿编辑。读文章,指出 3 个具体可改进的点。
tools: []
supervisor:
model: anthropic:claude-haiku-4-5
policy: |
一般流程:先 researcher,再 writer,再 reviewer。
若用户已提供调研材料则跳过 researcher。
所有 worker 都执行完后,回复 done。
maxSteps: 6   # 防止死循环

编译器 src/compiler.ts

State 用 Annotation.Root 扩展 MessagesAnnotation,额外加两个 channel:workerOutputs 用 reducer 合并对象,step 是累加计数器。

const State = Annotation.Root({
...MessagesAnnotation.spec,
workerOutputs: Annotation>({
reducer: (a, b) => ({ ...a, ...b }),
default: () => ({}),
}),
step: Annotation({
reducer: (a, b) => a + b,
default: () => 0,
}),
});

compileWorkflow 分五步:

1. 把每个 worker 包成可调用 Agent,存成 Map,用 createAgent

2. Supervisor 决策 schema 用 zodnext 是枚举(所有 worker id 加 "done"),instruction 是字符串,走 functionCalling

const DecisionSchema = z.object({
next: z.enum([...workerIds, "done"]),
instruction: z.string(),
});

const decide = supervisorModel.withStructuredOutput(DecisionSchema, {
method: "functionCalling",
});

3. supervisorNode:先看 state.step >= maxSteps,超了直接 goto END;否则把 formatContext 的结果喂给决策模型,next === "done" 时合并所有 workerOutputs 作为最终回复并 goto END,否则返回 Command

const supervisorNode = async (state: typeof State.State) => {
if (state.step >= config.supervisor.maxSteps) {
return new Command({
goto: END,
update: { messages: [new AIMessage("[maxSteps reached]")] },
});
}

const decision = await decide.invoke([
new SystemMessage(config.supervisor.policy),
new HumanMessage(formatContext(state)),
]);

if (decision.next === "done") {
return new Command({
goto: END,
update: {
messages: [new AIMessage(mergeOutputs(state.workerOutputs))],
},
});
}

return new Command({
goto: decision.next,
update: { step: 1, messages: [new AIMessage(decision.instruction)] },
});
};

4. makeWorkerNode:从最后一条 assistant 消息里取出指令,agent.invoke,然后用 Command 把产物回流到 supervisor:

const makeWorkerNode =
(w: WorkerConfig, agent: Agent) => async (state: typeof State.State) => {
const instruction = state.messages.at(-1)?.content as string;
const result = await agent.invoke({ messages: [new HumanMessage(instruction)] });

return new Command({
goto: "supervisor",
update: {
workerOutputs: { [w.id]: String(result.messages.at(-1)?.content ?? "") },
messages: [new AIMessage(`[${w.id}] 已完成`)] ,
},
});
};

5. 拼图:supervisor 的 ends 是所有 worker id 加 END,每个 worker 的 ends 只有 "supervisor",入口边从 START 指向 supervisor,最后带 MemorySaver 编译:

const builder = new StateGraph(State).addNode("supervisor", supervisorNode, {
ends: [...workerIds, END],
});

for (const w of config.workers) {
builder.addNode(w.id, makeWorkerNode(w, agents.get(w.id)!), {
ends: ["supervisor"],
});
}

builder.addEdge(START, "supervisor");
return builder.compile({ checkpointer: new MemorySaver() });

formatContext 对每个 worker 的输出只取前 200 字拼进 supervisor 的决策上下文。完整输出超出这个长度后,对「下一步选谁」的边际收益很低,截断省 token:

const formatContext = (state: typeof State.State) =>
Object.entries(state.workerOutputs)
.map(([id, out]) => `## ${id}\n${out.slice(0, 200)}`)
.join("\n\n");

类型与 YAML 校验

src/types.ts 里两个结构:WorkflowConfig { id, name, description?, workers: WorkerConfig[], supervisor: SupervisorConfig }WorkerConfig { id, model, systemPrompt, tools }

src/load.tsvalidate() 校验:id / workers / supervisor 必填,worker id 不重复,每个 worker 有 model 与 systemPrompt,supervisor.maxSteps >= 1

export function validate(cfg: WorkflowConfig) {
if (!cfg.id || !cfg.workers?.length || !cfg.supervisor) {
throw new Error("id / workers / supervisor 必填");
}
const seen = new Set();
for (const w of cfg.workers) {
if (seen.has(w.id)) throw new Error(`worker id 重复: ${w.id}`);
seen.add(w.id);
if (!w.model || !w.systemPrompt) {
throw new Error(`worker ${w.id} 缺少 model 或 systemPrompt`);
}
}
if (!(cfg.supervisor.maxSteps >= 1)) {
throw new Error("supervisor.maxSteps 必须 >= 1");
}
}

工具注册表 src/tool-registry.ts

工具不写在 YAML 里,YAML 只写 id。真实工具注册到运行时 registry,解析时按 id 取,未注册直接抛错:

const registry = new Map();

export const register = (id: string, tool: StructuredToolInterface) =>
registry.set(id, tool);

export function resolveTools(ids: string[]) {
return ids.map((id) => {
const tool = registry.get(id);
if (!tool) throw new Error(`tool 未注册: ${id}`);
return tool;
});
}

register("web_search", tavilySearch); // 内置示例,真实接 Tavily / Bing API

HTTP 服务:Hono + SSE

编译结果按 workflowId 缓存。POST /workflows/:id/runstreamSSEgraph.streamstreamMode: "updates" 逐节点推送:

app.post("/workflows/:id/run", async (c) => {
const graph = await getCompiled(c.req.param("id"));
const { input } = await c.req.json();

return streamSSE(c, async (stream) => {
const events = await graph.stream(
{ messages: [new HumanMessage(input)] },
{
configurable: { thread_id: crypto.randomUUID() },
streamMode: "updates",
},
);

for await (const chunk of events) {
await stream.writeSSE({ event: "step", data: JSON.stringify(chunk) });
}
});
});

POST /workflows/:id/reload 清掉缓存,供开发调试。

依赖:@langchain/langgraph ^1.0.0langchain ^1.4.0hono ^4.6.0yamlzodengines.node >= 20;dev 脚本用 tsx watch.env.example 里有 ANTHROPIC_API_KEYLANGSMITH_TRACING / API_KEY / PROJECT

{
"type": "module",
"engines": { "node": ">=20" },
"scripts": { "dev": "tsx watch src/server.ts" },
"dependencies": {
"@langchain/langgraph": "^1.0.0",
"langchain": "^1.4.0",
"hono": "^4.6.0",
"yaml": "^2.5.0",
"zod": "^3.23.0"
}
}

跑起来

npm install
# 把上面的 YAML 存到 workflows/content-pipeline.yaml
npm run dev

curl -N -X POST http://localhost:3000/workflows/content-pipeline/run \
-d '{"input":"写一篇关于 RAG 在企业搜索中演进的短文"}'

SSE 会按顺序吐出 supervisor -> researcher -> supervisor -> writer -> ... -> done

运行时节奏与四个关键点

控制权始终在 Supervisor 和 Worker 之间往返。Supervisor 每被唤醒一次,就用结构化输出决策下一个谁来干,用 Command({ goto }) 跳到对应 Worker;Worker 执行完再 goto 回 supervisor,直到判定 done 或撞上 maxSteps。Worker 之间从不直接连边。

  1. Supervisor 用 Command({ goto, update }) 显式跳转——比 addConditionalEdges 灵活,状态更新和路由一起返回。
  2. Worker 完成后回到 Supervisor——不要 Worker 之间直接连边,否则退化成静态流水线,YAML 就没意义了。
  3. maxSteps 是硬约束——动态调度的代价是可能死循环,Supervisor 自己不会停,就需要外部计数。
  4. Worker 是独立 Agent——每个 Worker 持有自己的工具集和 prompt,互不污染。

已知限制(这一版没做的事)

  1. 只支持 Anthropic providerparseModel 没分发到 OpenAI / Google,接其他 provider 加一个 switch 即可。
  2. 没有并发分支:Supervisor 一次只跳一个 worker。要做 researcher 和 fact_checker 并发,得让 Supervisor 返回多 goto,或在 graph 里加 Send API。
  3. handoff 上下文截断粗暴formatContext 取每个 worker 输出前 200 字,长链下信息丢失。生产环境应传 hash + 完整内容存 State 另一个 channel。
  4. Worker 内部出错没有重试 / 降级:抛异常会让整张图崩,要包 try/catch + withFallbacks 才能上线。
  5. 没接 LangSmith 评估:YAML 改了之后没自动跑回归,下一步要加 evaluate(targets, { data })
  6. YAML 没有 schema 验证validate() 只校验关键字段,应该用 Zod 定义完整 schema + 版本号字段做向前兼容。
  7. MemorySaver 不能多进程:上线必须换 PostgresSaver。

小结

平台层最值得展示的不是「做了多少功能」,而是把「什么变」和「什么不变」切干净:不变的是编译器、Supervisor 决策框架、State Schema、HTTP 服务;可变的是 YAML 里的 worker 列表、模型、prompt、工具引用。加业务只需写新工具注册到 registry、加新 YAML,不用碰编译器,也不用碰 LangGraph 拓扑构造。

推荐文章

程序员茄子在线接单