📷 [图片 token=Wy2CbrxSOoMSeLxMPYGc9co4nSd(未能下载,见飞书原文)]
前言
关键代码:app/agent/aiops/ 目录下的 state.py、planner.py、executor.py、replanner.py,以及 app/services/aiops_service.py。
📷 [图片 token=BxHebByEioqLYQxFy82chBk2nnh(未能下载,见飞书原文)]
流程梳理
运维 Agent 的核心目标是 规划 → 执行 → 评估 → 调整。整体流程就是三个节点:
Planner:拆解排查步骤,生成执行计划
Executor:从计划中取出第一个步骤,调用工具执行
Replanner:评估执行结果,决定继续、调整计划还是生成最终报告
三个节点通过 LangGraph StateGraph 串联,共享同一份 PlanExecuteState 状态对象在整个流程中传递。
📷 [图片 token=S0bFbHxhmoplyuxMjrJcPpXhn0d(未能下载,见飞书原文)]
实战
状态定义
整个 Plan-Execute-Replan 流程的数据流通过 PlanExecuteState 承载,字段设计非常简洁:
class PlanExecuteState(TypedDict):
input: str # 用户输入的任务描述
plan: List[str] # 待执行的步骤列表
past_steps: Annotated[List[tuple], operator.add] # 已执行的步骤历史(追加式更新)
response: str # 最终报告/响应
past_steps 使用 Annotated[List[tuple], operator.add] 声明,LangGraph 会将每次节点返回的 past_steps 自动追加到列表中,而不是覆盖,无需手动维护历史。
Planner 节点
Planner 负责制定执行计划,输出一个结构化的步骤列表。核心流程:
先调用
retrieve_knowledge查询知识库,寻找历史经验文档获取所有可用工具(本地工具 + MCP 工具),格式化为文字描述
将工具列表和经验文档注入 prompt,调用 LLM 生成结构化计划
async def planner(state: PlanExecuteState) -> Dict[str, Any]:
input_text = state.get("input", "")
# 1. 查询内部知识库,寻找相关经验
context_str = await retrieve_knowledge.ainvoke({"query": input_text})
experience_docs = context_str if context_str and context_str.strip() else ""
# 2. 获取所有可用工具(本地 + MCP)
local_tools = [get_current_time, retrieve_knowledge]
mcp_tools = await (await get_mcp_client_with_retry()).get_tools()
all_tools = local_tools + mcp_tools
# 3. 调用 LLM 生成结构化计划
llm = ChatQwen(model=config.rag_model, api_key=config.dashscope_api_key, temperature=0)
planner_chain = planner_prompt | llm.with_structured_output(Plan)
plan_result = await planner_chain.ainvoke({
"messages": [("user", input_text)],
"tools_description": format_tools_description(all_tools),
"experience_context": experience_docs
})
return {"plan": plan_result.steps}
计划的输出格式用 Pydantic Plan 模型约束,通过 llm.with_structured_output(Plan) 保证 LLM 的输出可以直接解析为步骤列表:
class Plan(BaseModel):
steps: List[str] = Field(
description="完成任务所需的不同步骤,按顺序执行,每一步建立在前一步的基础上。"
)
Planner Prompt
Planner 的系统提示词要求模型将任务分解为逻辑独立的步骤,每步指明使用哪个工具及所需参数。如果查到了经验文档,也会作为参考注入:
planner_prompt = ChatPromptTemplate.from_messages([
("system", """
作为一个专家级别的规划者,你需要将复杂的任务分解为可执行的步骤。
可用工具列表(用于制定计划时参考):
{tools_description}
注意:你的职责是制定计划,实际的工具调用由 Executor 负责执行。
{experience_context}
对于给定的任务,请创建一个简单的、逐步的计划:
- 将任务分解为逻辑上独立的步骤
- 每个步骤明确使用哪些工具(如果需要),最好能同时提供工具所需参数
- 步骤之间应有清晰的依赖关系
- 如果有相关经验文档,请参考其中的方法和步骤制定计划
"""),
("placeholder", "{messages}"),
])
Executor 节点
Executor 每次只执行计划中的第一个步骤,使用 LangGraph 的 ToolNode 自动处理工具调用,执行完后将该步骤从 plan 中移除,并将执行结果追加到 past_steps:
async def executor(state: PlanExecuteState) -> Dict[str, Any]:
plan = state.get("plan", [])
task = plan[0] # 只取第一个步骤
# 绑定工具的 LLM
all_tools = local_tools + mcp_tools
llm_with_tools = llm.bind_tools(all_tools)
tool_node = ToolNode(all_tools)
messages = [
SystemMessage(content="你是一个能力强大的助手,负责执行具体的任务步骤。..."),
HumanMessage(content=f"请执行以下任务: {task}")
]
# 第一步:LLM 决定是否需要工具调用
llm_response = await llm_with_tools.ainvoke(messages)
# 第二步:如果有工具调用,使用 ToolNode 自动执行
if hasattr(llm_response, "tool_calls") and llm_response.tool_calls:
messages.append(llm_response)
tool_messages = await tool_node.ainvoke({"messages": messages})
messages.extend(tool_messages["messages"])
# 第三步:将工具结果返回给 LLM 生成最终答案
final_response = await llm_with_tools.ainvoke(messages)
result = final_response.content
else:
result = llm_response.content
return {
"plan": plan[1:], # 移除已执行的第一个步骤
"past_steps": [(task, result)], # 追加执行历史
}
Replanner 节点
Replanner 根据原始任务、已执行步骤和剩余计划做出三选一的决策:
| 决策 | 含义 | 触发条件 |
|---|---|---|
respond | 信息充足,立即生成最终报告 | 最高优先级,已执行 ≥ 3 步且有关键信息 |
continue | 当前计划合理,继续执行 | 剩余步骤确实必要 |
replan | 调整计划,替换剩余步骤 | 最低优先级,计划有重大偏差时才使用 |
决策同样用 Pydantic 模型约束输出:
class Act(BaseModel):
action: str = Field(description="下一步行动: 'continue' | 'replan' | 'respond'")
new_steps: List[str] = Field(default_factory=list, description="replan 时的新步骤列表")
Replanner 核心逻辑(含安全限制):
async def replanner(state: PlanExecuteState) -> Dict[str, Any]:
past_steps = state.get("past_steps", [])
plan = state.get("plan", [])
# 安全限制:已执行步骤超过 8 步,强制生成响应,防止无限循环
if len(past_steps) >= 8:
return await _generate_response(state, llm)
if plan:
# 还有剩余步骤,让 LLM 做决策
act = await replanner_chain.ainvoke({
"messages": [
("user", f"原始任务: {input_text}"),
("user", f"已执行的步骤:\n{steps_summary}"),
("user", f"剩余计划: {', '.join(plan)}"),
],
"tools_description": tools_description
})
if act.action == "respond":
return await _generate_response(state, llm)
elif act.action == "replan":
# 安全限制:新步骤数不能超过当前剩余步骤数
new_steps = act.new_steps[:len(plan)]
return {"plan": new_steps}
else: # continue
return {} # 不修改状态,继续执行下一步
else:
# 计划已执行完毕,直接生成最终响应
return await _generate_response(state, llm)
当决定 respond 时,_generate_response 会整理所有执行历史,生成结构化的 Markdown 报告:
async def _generate_response(state, llm) -> Dict[str, Any]:
execution_history = "\n\n".join([
f"### 步骤: {step}\n**结果:**\n{result}"
for step, result in past_steps
])
response_gen = response_prompt | llm.with_structured_output(Response)
response_obj = await response_gen.ainvoke({
"messages": [
("user", f"原始任务: {input_text}"),
("user", f"执行历史:\n{execution_history}"),
("user", "请基于以上信息生成全面的最终响应"),
]
})
return {"response": response_obj.response}
构建 LangGraph 工作流
三个节点通过 StateGraph 连接,replanner 之后根据状态中是否存在 response 进行条件路由:
def _build_graph(self):
workflow = StateGraph(PlanExecuteState)
# 添加三个节点
workflow.add_node("planner", planner)
workflow.add_node("executor", executor)
workflow.add_node("replanner", replanner)
# 固定边:planner -> executor -> replanner
workflow.set_entry_point("planner")
workflow.add_edge("planner", "executor")
workflow.add_edge("executor", "replanner")
# 条件边:replanner 根据状态决定结束还是继续执行
def should_continue(state: PlanExecuteState) -> str:
if state.get("response"):
return END # 已生成最终响应,结束流程
if state.get("plan"):
return "executor" # 还有步骤,继续执行
return END
workflow.add_conditional_edges("replanner", should_continue, {
"executor": "executor",
END: END
})
return workflow.compile(checkpointer=MemorySaver())
工作流结构如下:
[用户任务] → planner → executor → replanner
↑ |
| continue/replan
└───────────┘
|
respond/空计划
↓
[END]
执行工作流(流式输出)
AIOpsService.execute 使用 graph.astream(stream_mode="updates") 流式执行,每个节点完成后立即产生事件,可以实时推送给前端展示执行进度:
async def execute(self, user_input: str, session_id: str) -> AsyncGenerator:
initial_state: PlanExecuteState = {
"input": user_input,
"plan": [],
"past_steps": [],
"response": ""
}
async for event in self.graph.astream(
input=initial_state,
config={"configurable": {"thread_id": session_id}},
stream_mode="updates" # 每个节点更新后立即推送
):
for node_name, node_output in event.items():
if node_name == "planner":
yield self._format_planner_event(node_output) # type: plan
elif node_name == "executor":
yield self._format_executor_event(node_output) # type: step_complete
elif node_name == "replanner":
yield self._format_replanner_event(node_output) # type: report / status
# 所有节点完成,推送最终响应
final_state = self.graph.get_state({"configurable": {"thread_id": session_id}})
final_response = final_state.values.get("response", "") if final_state else ""
yield {
"type": "complete",
"stage": "complete",
"message": "任务执行完成",
"response": final_response
}
流式事件类型说明:
type | stage | 含义 |
|---|---|---|
plan | plan_created | Planner 生成了执行计划 |
step_complete | step_executed | Executor 执行完一个步骤 |
report | final_report | Replanner 生成了最终报告 |
status | 各节点名 | 节点运行中的状态通知 |
complete | complete | 整个工作流结束 |
error | error | 执行出错 |
总结
通过上面的分析,我们已经了解了 Planner、Executor、Replanner 的作用和相关 prompt。代码的核心是 LangGraph 的 StateGraph 管理节点间的状态流转,PlanExecuteState 在整个流程中传递,三个节点各司其职:
Planner 查询知识库获取经验,生成结构化步骤列表
Executor 取出第一步,通过
ToolNode自动执行工具调用,返回结果并移除该步Replanner 评估已执行结果,决定是继续、调整还是收敛生成最终报告
框架帮我们处理了节点间的数据传递、条件路由和流式输出,核心还是要搞懂设计原理:Plan-Execute-Replan 本质上是一个带反馈的闭环调度器,Replanner 的收敛策略决定了整个流程的质量与效率。
📷 [图片 token=GIrDbzR7WouUWMxHlD3cQGyYnpf(未能下载,见飞书原文)]