首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >Agent Loop + Graph Engineer 实战:从循环式智能体到图编排引擎的完整代码实现

Agent Loop + Graph Engineer 实战:从循环式智能体到图编排引擎的完整代码实现

原创
作者头像
资源大佬 jzit-top
发布于 2026-09-22 16:03:31
发布于 2026-09-22 16:03:31
1490
举报

摘要:Agent Loop(智能体循环)是当前 LLM Agent 最核心的执行范式,而 Graph Engineer(图工程)则代表下一代智能体编排方向——用有向图显式建模状态、节点、边与条件分支。本文从专业视角系统讲解两者的原理、边界与融合方式,并给出一套可运行的 Python 实现:包含 ReAct 循环、反思循环、工具调用、状态机、DAG 编排、条件边、并行节点、检查点与恢复、人机协作、子图嵌套、可观测性与测试。


目录

  1. 为什么需要 Agent Loop
  2. 为什么需要 Graph Engineer
  3. Agent Loop 与 Graph 的关系
  4. Agent Loop 核心原理与实现
  5. Graph Engineer 核心原理与实现
  6. 关键机制:状态、节点、边、检查点
  7. 代码实战:从 Loop 到 Graph 的完整系统
  8. 模块一:状态与消息模型
  9. 模块二:Agent Loop 引擎
  10. 模块三:工具系统
  11. 模块四:Graph 编排引擎
  12. 模块五:条件边与并行节点
  13. 模块六:检查点与恢复
  14. 模块七:人机协作(HITL)
  15. 模块八:子图与多智能体
  16. 模块九:可观测性与审计
  17. 模块十:API 与测试
  18. 生产化清单
  19. 常见陷阱
  20. 总结

一、为什么需要 Agent Loop

最早的 LLM 应用是单轮问答:

代码语言:javascript
复制
用户输入 -> LLM -> 输出

但它有三个根本问题:

  1. 不能调用外部工具;
  2. 不能多步推理;
  3. 出错无法自我修正。

ReAct(Reasoning + Acting)范式解决了这个问题:

代码语言:javascript
复制
Thought -> Action -> Observation -> Thought -> ... -> Final Answer

这个循环就是 Agent Loop。

Agent Loop 的价值:

能力

单轮问答

Agent Loop

多步任务

❌

✅

工具调用

❌

✅

自我修正

❌

✅

长任务

❌

✅

可观测

弱

强

一个最小 Agent Loop:

代码语言:javascript
复制
while not done:
    thought = llm.think(state)
    if thought.is_final:
        break
    action = thought.action
    observation = tools.run(action)
    state.append(observation)

二、为什么需要 Graph Engineer

Agent Loop 简单,但有明显上限:

  • 状态不可见:循环里发生了什么,很难追踪;
  • 无法并行:所有步骤串行;
  • 无法分支:复杂业务逻辑难以表达;
  • 无法恢复:进程崩溃就得重来;
  • 难以调试:黑盒循环,出错难定位;
  • 难以协作:多智能体协调困难。

Graph Engineer 的思路是:把 Agent 执行过程显式建模为图。

代码语言:javascript
复制
节点(Node)= 一个执行单元(LLM 调用、工具调用、函数)
边(Edge)  = 控制流(顺序、条件、并行)
状态(State)= 在节点间流动的共享数据
检查点(Checkpoint)= 每一步的状态快照

Graph 的优势:

能力

Agent Loop

Graph Engineer

状态可见

弱

强

并行

❌

✅

条件分支

隐式

显式

断点恢复

❌

✅

人机协作

难

原生支持

多智能体

难

子图

可视化

难

天然

一句话:Loop 是运行时,Graph 是编排层。


三、Agent Loop 与 Graph 的关系

两者不是替代关系,而是互补:

代码语言:javascript
复制
┌─────────────────────────────────────────┐
│              Graph 编排层                │
│  ┌──────────┐  ┌──────────┐  ┌────────┐ │
│  │  Node A  │->│  Node B  │->│ Node C │ │
│  │ (Loop)   │  │ (Tool)   │  │ (LLM)  │ │
│  └──────────┘  └──────────┘  └────────┘ │
└─────────────────────────────────────────┘
  • Graph:负责宏观流程、状态流转、分支、并行、恢复;
  • Loop:负责单个节点内部的推理与工具调用循环。

工程上常见组合:

  • 简单任务:纯 Loop;
  • 复杂任务:Graph 编排多个 Loop 节点;
  • 多智能体:Graph 的子图 + 多 Loop。

四、Agent Loop 核心原理与实现

4.1 基本循环

代码语言:javascript
复制
┌─────────────────────────────────────────┐
│                                         │
│   ┌───────────┐                         │
│   │  Observe  │◄────────────┐           │
│   └─────┬─────┘             │           │
│         ▼                   │           │
│   ┌───────────┐             │           │
│   │  Think    │             │           │
│   └─────┬─────┘             │           │
│         ▼                   │           │
│   ┌───────────┐             │           │
│   │   Act     │             │           │
│   └─────┬─────┘             │           │
│         ▼                   │           │
│   ┌───────────┐             │           │
│   │  Check    │───── 未完成 ┘           │
│   └─────┬─────┘                         │
│         ▼ 完成                          │
│      Final Answer                       │
└─────────────────────────────────────────┘

4.2 变体

变体

说明

适用场景

ReAct

Thought + Action + Observation

通用工具调用

Reflexion

增加自我反思节点

需要纠错

Plan-and-Execute

先规划再执行

长任务

Tree-of-Thought

多路径探索

复杂推理

Self-Consistency

多次采样投票

提高稳定性

4.3 循环控制要点

  • 最大步数:防止死循环;
  • 超时:防止长时间挂起;
  • 成本上限:防止 Token 爆炸;
  • 重复检测:相同动作重复时打断;
  • 终止条件:显式定义完成标准。

五、Graph Engineer 核心原理与实现

5.1 图的基本构成

代码语言:javascript
复制
State(状态):节点间共享的数据
Node(节点):执行单元
Edge(边):控制流
Condition(条件):决定走哪条边
Checkpoint(检查点):状态快照

5.2 图的执行模型

代码语言:javascript
复制
1. 初始化 State
2. 从 START 节点开始
3. 执行当前节点,更新 State
4. 根据边与条件选择下一个节点
5. 重复直到到达 END
6. 每一步可保存检查点

5.3 节点类型

类型

说明

LLM Node

调用大模型

Tool Node

调用工具

Function Node

纯函数处理

Condition Node

条件判断

Parallel Node

并行分支

Subgraph Node

嵌套子图

Human Node

等待人工输入

5.4 边的类型

类型

说明

普通边

固定跳转

条件边

根据 State 选择

并行边

同时触发多个节点

汇聚边

等待多个节点完成


六、关键机制:状态、节点、边、检查点

6.1 状态

状态是一个可序列化的字典,包含:

  • 消息历史;
  • 工具结果;
  • 中间变量;
  • 错误信息;
  • 元数据(用户、会话、成本)。

6.2 节点

每个节点是一个函数:

代码语言:javascript
复制
def node(state: State) -> State:
    ...
    return updated_state

6.3 边

代码语言:javascript
复制
graph.add_edge("node_a", "node_b")
graph.add_conditional_edge("node_b", router_fn, {"yes": "node_c", "no": "node_d"})

6.4 检查点

每步保存 State 快照,用于:

  • 断点恢复;
  • 时间旅行调试;
  • 人工审批后继续;
  • 审计。

七、代码实战:从 Loop 到 Graph 的完整系统

7.1 技术栈

  • Python 3.12
  • FastAPI(API 层)
  • Pydantic(Schema)
  • SQLModel(持久化)
  • SQLite(示例)/ PostgreSQL(生产)
  • Pytest

7.2 项目结

代码语言:javascript
复制
agent-graph-engine/
├── app/
│   ├── __init__.py
│   ├── main.py
│   ├── state.py
│   ├── loop.py
│   ├── tools.py
│   ├── graph.py
│   ├── nodes.py
│   ├── checkpoint.py
│   ├── hitl.py
│   ├── subgraph.py
│   ├── observability.py
│   └── llm.py
├── tests/
│   └── test_agent_graph.py
├── requirements.txt
└── Dockerfile

7.3 依赖

代码语言:javascript
复制
pip install fastapi uvicorn pydantic sqlmodel pytest httpx

requirements.txt:

代码语言:javascript
复制
fastapi
uvicorn[standard]
pydantic
sqlmodel
pytest
httpx

八、模块一:状态与消息模型

app/state.py:

代码语言:javascript
复制
from dataclasses import dataclass, field
from typing import Any, Dict, List, Optional
from datetime import datetime
import uuid


@dataclass
class Message:
    role: str                    # system / user / assistant / tool
    content: str
    tool_name: Optional[str] = None
    tool_call_id: Optional[str] = None
    timestamp: str = field(default_factory=lambda: datetime.utcnow().isoformat())


@dataclass
class AgentState:
    """
    图执行的共享状态。所有节点读写同一个 State。
    """
    run_id: str = field(default_factory=lambda: str(uuid.uuid4()))
    user_id: int = 0
    goal: str = ""
    messages: List[Message] = field(default_factory=list)
    variables: Dict[str, Any] = field(default_factory=dict)
    tool_results: List[Dict[str, Any]] = field(default_factory=list)
    errors: List[str] = field(default_factory=list)
    step: int = 0
    max_steps: int = 20
    status: str = "running"      # running / done / failed / waiting_human
    next_node: Optional[str] = None
    metadata: Dict[str, Any] = field(default_factory=dict)

    def add_message(self, role: str, content: str, **kwargs):
        self.messages.append(Message(role=role, content=content, **kwargs))

    def set(self, key: str, value: Any):
        self.variables[key] = value

    def get(self, key: str, default: Any = None) -> Any:
        return self.variables.get(key, default)

    def to_dict(self) -> dict:
        return {
            "run_id": self.run_id,
            "user_id": self.user_id,
            "goal": self.goal,
            "messages": [m.__dict__ for m in self.messages],
            "variables": self.variables,
            "tool_results": self.tool_results,
            "errors": self.errors,
            "step": self.step,
            "status": self.status,
            "next_node": self.next_node,
            "metadata": self.metadata,
        }

    @classmethod
    def from_dict(cls, data: dict) -> "AgentState":
        state = cls(
            run_id=data["run_id"],
            user_id=data.get("user_id", 0),
            goal=data.get("goal", ""),
            variables=data.get("variables", {}),
            tool_results=data.get("tool_results", []),
            errors=data.get("errors", []),
            step=data.get("step", 0),
            status=data.get("status", "running"),
            next_node=data.get("next_node"),
            metadata=data.get("metadata", {}),
        )
        state.messages = [Message(**m) for m in data.get("messages", [])]
        return state

九、模块二:Agent Loop 引擎

app/llm.py:

代码语言:javascript
复制
from dataclasses import dataclass
from typing import Optional


@dataclass
class LLMDecision:
    kind: str                       # "tool" / "final" / "think"
    content: str = ""
    tool_name: Optional[str] = None
    arguments: Optional[dict] = None


def mock_llm(
    state,
    tool_schemas: list,
) -> LLMDecision:
    """
    Mock LLM:规则驱动,便于离线运行。
    生产环境替换为 OpenAI / Claude function calling。
    """
    goal = state.goal.lower()
    last_msgs = state.messages[-5:]

    # 看是否已经有工具结果,避免重复调用
    has_tool_result = any(m.role == "tool" for m in state.messages)

    if "计算" in state.goal and not has_tool_result:
        return LLMDecision(
            kind="tool",
            tool_name="calculator",
            arguments={"expression": "2+3*4"},
        )

    if "天气" in state.goal and not has_tool_result:
        return LLMDecision(
            kind="tool",
            tool_name="weather",
            arguments={"city": "上海"},
        )

    if "搜索" in state.goal and not has_tool_result:
        return LLMDecision(
            kind="tool",
            tool_name="search",
            arguments={"query": state.goal},
        )

    if has_tool_result:
        return LLMDecision(
            kind="final",
            content=f"任务完成。基于工具结果,答案是:{state.tool_results[-1]}",
        )

    return LLMDecision(
        kind="final",
        content=f"收到目标:{state.goal}。当前没有合适工具,直接回答。",
    )

app/loop.py:

代码语言:javascript
复制
from app.state import AgentState
from app.llm import mock_llm
from app.tools import ToolRegistry


class AgentLoop:
    """
    单节点内的 ReAct 循环。
    """

    def __init__(
        self,
        tools: ToolRegistry,
        max_steps: int = 10,
        max_repeats: int = 2,
    ):
        self.tools = tools
        self.max_steps = max_steps
        self.max_repeats = max_repeats

    def run(self, state: AgentState) -> AgentState:
        action_history = []

        for step in range(self.max_steps):
            state.step += 1

            if state.step > state.max_steps:
                state.status = "failed"
                state.errors.append("超过最大步数")
                return state

            decision = mock_llm(state, self.tools.to_schema())

            if decision.kind == "final":
                state.add_message("assistant", decision.content)
                state.status = "done"
                return state

            if decision.kind == "tool":
                # 重复检测
                key = f"{decision.tool_name}:{decision.arguments}"
                action_history.append(key)
                if action_history.count(key) > self.max_repeats:
                    state.errors.append(f"重复动作:{key}")
                    state.add_message("assistant", "检测到重复动作,停止循环。")
                    state.status = "failed"
                    return state

                state.add_message(
                    "assistant",
                    f"调用工具 {decision.tool_name}",
                    tool_name=decision.tool_name,
                )

                result = self.tools.execute(
                    decision.tool_name,
                    decision.arguments or {},
                    state,
                )
                state.tool_results.append(result)
                state.add_message(
                    "tool",
                    str(result),
                    tool_name=decision.tool_name,
                )
                continue

        state.status = "failed"
        state.errors.append("循环未收敛")
        return state

十、模块三:工具系统

app/tools.py:

代码语言:javascript
复制
from dataclasses import dataclass, field
from typing import Any, Callable, Dict, List


@dataclass
class ToolSpec:
    name: str
    description: str
    parameters: Dict[str, Any]
    handler: Callable[..., Any]
    risk_level: str = "low"


class ToolRegistry:
    def __init__(self):
        self._tools: Dict[str, ToolSpec] = {}

    def register(self, spec: ToolSpec):
        self._tools[spec.name] = spec

    def get(self, name: str) -> ToolSpec | None:
        return self._tools.get(name)

    def to_schema(self) -> List[dict]:
        return [
            {
                "name": t.name,
                "description": t.description,
                "parameters": t.parameters,
            }
            for t in self._tools.values()
        ]

    def execute(self, name: str, args: dict, state) -> dict:
        spec = self.get(name)
        if not spec:
            return {"error": f"未知工具:{name}"}
        try:
            result = spec.handler(**args)
            return {"tool": name, "ok": True, "result": result}
        except Exception as e:
            return {"tool": name, "ok": False, "error": str(e)}


def tool_calculator(expression: str) -> dict:
    allowed = set("0123456789+-*/(). ")
    if not set(expression) <= allowed:
        raise ValueError("表达式包含非法字符")
    value = eval(expression, {"__builtins__": {}}, {})
    return {"expression": expression, "value": value}


def tool_weather(city: str) -> dict:
    return {"city": city, "temperature": 24, "condition": "多云"}


def tool_search(query: str) -> dict:
    return {
        "query": query,
        "results": [
            {"title": "示例结果 1", "url": "https://example.com/1"},
            {"title": "示例结果 2", "url": "https://example.com/2"},
        ],
    }


def build_default_registry() -> ToolRegistry:
    reg = ToolRegistry()
    reg.register(ToolSpec(
        name="calculator",
        description="计算数学表达式",
        parameters={
            "type": "object",
            "properties": {"expression": {"type": "string"}},
            "required": ["expression"],
        },
        handler=tool_calculator,
        risk_level="low",
    ))
    reg.register(ToolSpec(
        name="weather",
        description="查询城市天气",
        parameters={
            "type": "object",
            "properties": {"city": {"type": "string"}},
            "required": ["city"],
        },
        handler=tool_weather,
    ))
    reg.register(ToolSpec(
        name="search",
        description="搜索信息",
        parameters={
            "type": "object",
            "properties": {"query": {"type": "string"}},
            "required": ["query"],
        },
        handler=tool_search,
    ))
    return reg

十一、模块四:Graph 编排引擎

app/graph.py:

代码语言:javascript
复制
from typing import Callable, Dict, List, Optional, Any
from app.state import AgentState


NodeFn = Callable[[AgentState], AgentState]
RouterFn = Callable[[AgentState], str]


class Node:
    def __init__(self, name: str, fn: NodeFn, kind: str = "function"):
        self.name = name
        self.fn = fn
        self.kind = kind

    def run(self, state: AgentState) -> AgentState:
        return self.fn(state)


class Graph:
    """
    极简图编排引擎。

    - 支持普通边
    - 支持条件边
    - 支持并行节点
    - 支持检查点(通过 checkpoint_store)
    - 支持循环(通过条件边回到上游)
    """

    def __init__(self, name: str = "graph"):
        self.name = name
        self.nodes: Dict[str, Node] = {}
        self.edges: Dict[str, List[str]] = {}
        self.conditional_edges: Dict[str, tuple] = {}
        self.entry: Optional[str] = None
        self.finish: str = "END"
        self.checkpoint_store = None
        self.parallel_targets: Dict[str, List[str]] = {}

    def add_node(self, name: str, fn: NodeFn, kind: str = "function"):
        self.nodes[name] = Node(name, fn, kind)
        return self

    def set_entry(self, name: str):
        self.entry = name
        return self

    def add_edge(self, src: str, dst: str):
        self.edges.setdefault(src, []).append(dst)
        return self

    def add_conditional_edge(
        self,
        src: str,
        router: RouterFn,
        mapping: Dict[str, str],
    ):
        self.conditional_edges[src] = (router, mapping)
        return self

    def add_parallel(self, src: str, targets: List[str], join: str):
        """
        并行:执行 targets 中所有节点,然后汇聚到 join。
        """
        self.parallel_targets[src] = targets
        for t in targets:
            self.add_edge(t, join)
        return self

    def _save_checkpoint(self, state: AgentState, node: str):
        if self.checkpoint_store:
            self.checkpoint_store.save(state, node)

    def run(
        self,
        state: AgentState,
        resume_from: Optional[str] = None,
    ) -> AgentState:
        current = resume_from or self.entry
        if not current:
            raise ValueError("图未设置入口节点")

        visited = []

        while current != self.finish:
            node = self.nodes.get(current)
            if not node:
                state.errors.append(f"节点不存在:{current}")
                state.status = "failed"
                return state

            visited.append(current)

            try:
                state = node.run(state)
            except Exception as e:
                state.errors.append(f"节点 {current} 异常:{e}")
                state.status = "failed"
                self._save_checkpoint(state, current)
                return state

            state.next_node = None
            self._save_checkpoint(state, current)

            if state.status in ("failed", "waiting_human"):
                return state

            # 并行分支
            if current in self.parallel_targets:
                targets = self.parallel_targets[current]
                for t in targets:
                    t_node = self.nodes.get(t)
                    if t_node:
                        state = t_node.run(state)
                        self._save_checkpoint(state, t)
                # 并行后走普通边
                nexts = self.edges.get(current, [])
                current = nexts[0] if nexts else self.finish
                continue

            # 条件边优先
            if current in self.conditional_edges:
                router, mapping = self.conditional_edges[current]
                key = router(state)
                nxt = mapping.get(key)
                if not nxt:
                    state.errors.append(f"条件边未匹配:{current} -> {key}")
                    state.status = "failed"
                    return state
                current = nxt
                continue

            # 普通边
            nexts = self.edges.get(current, [])
            if not nexts:
                current = self.finish
            else:
                current = nexts[0]

        if state.status == "running":
            state.status = "done"
        return state

    def mermaid(self) -> str:
        """
        输出 Mermaid 图,便于可视化。
        """
        lines = ["graph TD"]
        for src, dsts in self.edges.items():
            for d in dsts:
                lines.append(f"  {src} --> {d}")
        for src, (_, mapping) in self.conditional_edges.items():
            for label, dst in mapping.items():
                lines.append(f"  {src} -- {label} --> {dst}")
        return "\n".join(lines)

十二、模块五:条件边与并行节点

app/nodes.py:

代码语言:javascript
复制
from app.state import AgentState
from app.loop import AgentLoop


def node_classify(state: AgentState) -> AgentState:
    goal = state.goal
    if any(k in goal for k in ["计算", "天气", "搜索"]):
        state.set("task_type", "tool")
    else:
        state.set("task_type", "chat")
    state.add_message("system", f"任务分类:{state.get('task_type')}")
    return state


def node_plan(state: AgentState) -> AgentState:
    state.set("plan", [
        "识别目标",
        "调用工具",
        "汇总结果",
    ])
    state.add_message("system", "已生成计划")
    return state


def node_execute_loop(tools) -> callable:
    def _run(state: AgentState) -> AgentState:
        loop = AgentLoop(tools)
        return loop.run(state)
    return _run


def node_chat(state: AgentState) -> AgentState:
    state.add_message("assistant", f"普通对话回复:{state.goal}")
    state.status = "done"
    return state


def node_review(state: AgentState) -> AgentState:
    if state.errors:
        state.set("review", "has_errors")
    else:
        state.set("review", "ok")
    return state


def node_parallel_a(state: AgentState) -> AgentState:
    state.set("branch_a", "A 完成")
    return state


def node_parallel_b(state: AgentState) -> AgentState:
    state.set("branch_b", "B 完成")
    return state


def node_join(state: AgentState) -> AgentState:
    state.set("join_result", f"{state.get('branch_a')} + {state.get('branch_b')}")
    state.add_message("system", "并行分支已汇聚")
    return state


def route_by_task(state: AgentState) -> str:
    return state.get("task_type", "chat")


def route_by_review(state: AgentState) -> str:
    return state.get("review", "ok")

十三、模块六:检查点与恢复

app/checkpoint.py:

代码语言:javascript
复制
import json
from datetime import datetime
from typing import Optional
from sqlmodel import SQLModel, Field, Session, select

from app.state import AgentState
from app.db import engine


class Checkpoint(SQLModel, table=True):
    id: Optional[int] = Field(default=None, primary_key=True)
    run_id: str = Field(index=True)
    node: str
    state_json: str
    created_at: datetime = Field(default_factory=datetime.utcnow)


class SQLCheckpointStore:
    def save(self, state: AgentState, node: str):
        with Session(engine) as db:
            cp = Checkpoint(
                run_id=state.run_id,
                node=node,
                state_json=json.dumps(state.to_dict(), ensure_ascii=False, default=str),
            )
            db.add(cp)
            db.commit()

    def latest(self, run_id: str) -> Optional[dict]:
        with Session(engine) as db:
            row = db.exec(
                select(Checkpoint)
                .where(Checkpoint.run_id == run_id)
                .order_by(Checkpoint.id.desc())
            ).first()
            if not row:
                return None
            return {
                "node": row.node,
                "state": json.loads(row.state_json),
            }

    def history(self, run_id: str) -> list:
        with Session(engine) as db:
            rows = db.exec(
                select(Checkpoint)
                .where(Checkpoint.run_id == run_id)
                .order_by(Checkpoint.id.asc())
            ).all()
            return [
                {"node": r.node, "state": json.loads(r.state_json), "at": str(r.created_at)}
                for r in rows
            ]

app/db.py:

代码语言:javascript
复制
from sqlmodel import SQLModel, create_engine, Session

DATABASE_URL = "sqlite:///./agent_graph.db"

engine = create_engine(
    DATABASE_URL,
    connect_args={"check_same_thread": False},
)


def init_db():
    SQLModel.metadata.create_all(engine)


def get_session():
    with Session(engine) as session:
        yield session

恢复执行:

代码语言:javascript
复制
def resume_run(run_id: str, graph, store: SQLCheckpointStore) -> AgentState:
    latest = store.latest(run_id)
    if not latest:
        raise ValueError("没有找到检查点")

    state = AgentState.from_dict(latest["state"])
    state.status = "running"
    return graph.run(state, resume_from=latest["node"])

十四、模块七:人机协作(HITL)

app/hitl.py:

代码语言:javascript
复制
from app.state import AgentState


def node_high_risk_action(state: AgentState) -> AgentState:
    """
    高风险动作节点:需要人工确认。
    """
    state.set("pending_action", {
        "type": "delete_resource",
        "target": "prod-db-01",
        "reason": "用户请求清理",
    })
    state.status = "waiting_human"
    state.add_message("system", "等待人工确认高风险操作")
    return state


def approve(state: AgentState, approved: bool, comment: str = "") -> AgentState:
    if state.status != "waiting_human":
        return state

    state.set("human_decision", {
        "approved": approved,
        "comment": comment,
    })

    if approved:
        state.add_message("system", "人工已批准,继续执行")
        state.status = "running"
    else:
        state.add_message("system", "人工已拒绝,终止流程")
        state.status = "failed"
        state.errors.append("人工拒绝")

    return state

十五、模块八:子图与多智能体

app/subgraph.py:

代码语言:javascript
复制
from app.graph import Graph
from app.state import AgentState


def as_node(subgraph: Graph, name: str):
    """
    把子图封装成节点,实现嵌套编排。
    """
    def _run(state: AgentState) -> AgentState:
        return subgraph.run(state)
    return name, _run


def build_research_subgraph(tools) -> Graph:
    from app.nodes import node_execute_loop

    sub = Graph(name="research")
    sub.add_node("search", node_execute_loop(tools))
    sub.add_node("summarize", _summarize)
    sub.set_entry("search")
    sub.add_edge("search", "summarize")
    sub.add_edge("summarize", "END")
    return sub


def _summarize(state: AgentState) -> AgentState:
    state.set("summary", "已汇总搜索结果")
    return state

多智能体示例:

代码语言:javascript
复制
def build_multi_agent_graph(tools) -> Graph:
    from app.nodes import (
        node_classify, node_plan, node_execute_loop,
        node_chat, node_review, route_by_task, route_by_review,
    )

    g = Graph(name="multi_agent")
    g.add_node("classify", node_classify)
    g.add_node("plan", node_plan)
    g.add_node("execute", node_execute_loop(tools))
    g.add_node("chat", node_chat)
    g.add_node("review", node_review)

    g.set_entry("classify")
    g.add_edge("classify", "plan")
    g.add_conditional_edge("plan", route_by_task, {
        "tool": "execute",
        "chat": "chat",
    })
    g.add_edge("execute", "review")
    g.add_conditional_edge("review", route_by_review, {
        "ok": "END",
        "has_errors": "plan",   # 出错回到 plan,形成循环
    })
    return g

十六、模块九:可观测性与审计

app/observability.py:

代码语言:javascript
复制
import json
from datetime import datetime
from typing import Any, Dict, List
from sqlmodel import SQLModel, Field, Session, select

from app.db import engine


class TraceEvent(SQLModel, table=True):
    id: int | None = Field(default=None, primary_key=True)
    run_id: str = Field(index=True)
    node: str
    event_type: str           # node_start / node_end / tool_call / error
    payload: str
    created_at: datetime = Field(default_factory=datetime.utcnow)


class Tracer:
    def __init__(self, run_id: str):
        self.run_id = run_id

    def log(self, node: str, event_type: str, payload: Dict[str, Any]):
        with Session(engine) as db:
            db.add(TraceEvent(
                run_id=self.run_id,
                node=node,
                event_type=event_type,
                payload=json.dumps(payload, ensure_ascii=False, default=str),
            ))
            db.commit()

    def timeline(self) -> List[dict]:
        with Session(engine) as db:
            rows = db.exec(
                select(TraceEvent)
                .where(TraceEvent.run_id == self.run_id)
                .order_by(TraceEvent.id.asc())
            ).all()
            return [
                {
                    "node": r.node,
                    "type": r.event_type,
                    "payload": json.loads(r.payload),
                    "at": str(r.created_at),
                }
                for r in rows
            ]

带追踪的节点包装:

代码语言:javascript
复制
def traced(name: str, fn, tracer: Tracer):
    def _run(state: AgentState) -> AgentState:
        tracer.log(name, "node_start", {"step": state.step})
        try:
            state = fn(state)
            tracer.log(name, "node_end", {"status": state.status})
        except Exception as e:
            tracer.log(name, "error", {"error": str(e)})
            raise
        return state
    return _run

十七、模块十:API 与测试

app/main.py:

代码语言:javascript
复制
from fastapi import FastAPI, Depends, HTTPException
from pydantic import BaseModel
from sqlmodel import Session

from app.db import init_db, get_session
from app.state import AgentState
from app.tools import build_default_registry
from app.checkpoint import SQLCheckpointStore
from app.nodes import build_multi_agent_graph
from app.observability import Tracer
from app.hitl import approve

app = FastAPI(title="Agent Loop + Graph Engineer")
store = SQLCheckpointStore()
tools = build_default_registry()


@app.on_event("startup")
def on_startup():
    init_db()


class RunRequest(BaseModel):
    goal: str
    user_id: int = 1


class ApproveRequest(BaseModel):
    run_id: str
    approved: bool
    comment: str = ""


@app.get("/health")
def health():
    return {"status": "ok"}


@app.post("/run")
def run_graph(payload: RunRequest):
    graph = build_multi_agent_graph(tools)
    graph.checkpoint_store = store

    state = AgentState(user_id=payload.user_id, goal=payload.goal)
    tracer = Tracer(state.run_id)

    tracer.log("graph", "start", {"goal": payload.goal})
    state = graph.run(state)
    tracer.log("graph", "end", {"status": state.status})

    return {
        "run_id": state.run_id,
        "status": state.status,
        "messages": [m.__dict__ for m in state.messages],
        "tool_results": state.tool_results,
        "variables": state.variables,
        "errors": state.errors,
        "timeline": tracer.timeline(),
    }


@app.get("/runs/{run_id}/checkpoints")
def get_checkpoints(run_id: str):
    return store.history(run_id)


@app.post("/runs/resume")
def resume(run_id: str):
    latest = store.latest(run_id)
    if not latest:
        raise HTTPException(status_code=404, detail="No checkpoint")

    graph = build_multi_agent_graph(tools)
    graph.checkpoint_store = store

    state = AgentState.from_dict(latest["state"])
    state.status = "running"
    state = graph.run(state, resume_from=latest["node"])
    return {"run_id": state.run_id, "status": state.status}


@app.post("/runs/approve")
def approve_run(payload: ApproveRequest):
    latest = store.latest(payload.run_id)
    if not latest:
        raise HTTPException(status_code=404, detail="No checkpoint")

    state = AgentState.from_dict(latest["state"])
    state = approve(state, payload.approved, payload.comment)

    if state.status == "running":
        graph = build_multi_agent_graph(tools)
        graph.checkpoint_store = store
        state = graph.run(state, resume_from=latest["node"])

    return {"run_id": state.run_id, "status": state.status}


@app.get("/graph/mermaid")
def graph_mermaid():
    graph = build_multi_agent_graph(tools)
    return {"mermaid": graph.mermaid()}

tests/test_agent_graph.p

代码语言:javascript
复制
import os
import pytest
from fastapi.testclient import TestClient

os.environ["DATABASE_URL"] = "sqlite:///./test_agent_graph.db"

from app.main import app  # noqa: E402
from app.graph import Graph  # noqa: E402
from app.state import AgentState  # noqa: E402
from app.tools import build_default_registry  # noqa: E402
from app.loop import AgentLoop  # noqa: E402

client = TestClient(app)


def test_health():
    assert client.get("/health").json()["status"] == "ok"


def test_agent_loop_calculator():
    tools = build_default_registry()
    loop = AgentLoop(tools)
    state = AgentState(goal="帮我计算 2+3*4")
    state = loop.run(state)
    assert state.status == "done"
    assert any("value" in str(r) for r in state.tool_results)


def test_agent_loop_weather():
    tools = build_default_registry()
    loop = AgentLoop(tools)
    state = AgentState(goal="查一下上海天气")
    state = loop.run(state)
    assert state.status == "done"
    assert any("上海" in str(r) for r in state.tool_results)


def test_graph_run_tool_task():
    r = client.post("/run", json={"goal": "帮我计算 2+3*4"})
    assert r.status_code == 200
    data = r.json()
    assert data["status"] == "done"
    assert len(data["tool_results"]) >= 1


def test_graph_run_chat_task():
    r = client.post("/run", json={"goal": "你好,介绍一下你自己"})
    assert r.status_code == 200
    data = r.json()
    assert data["status"] == "done"


def test_checkpoints_saved():
    r = client.post("/run", json={"goal": "查一下上海天气"})
    run_id = r.json()["run_id"]

    r = client.get(f"/runs/{run_id}/checkpoints")
    assert r.status_code == 200
    assert len(r.json()) >= 2


def test_resume_run():
    r = client.post("/run", json={"goal": "帮我计算 2+3*4"})
    run_id = r.json()["run_id"]

    r = client.post("/runs/resume", params={"run_id": run_id})
    assert r.status_code == 200
    assert r.json()["status"] in ("done", "failed")


def test_mermaid_output():
    r = client.get("/graph/mermaid")
    assert r.status_code == 200
    assert "graph TD" in r.json()["mermaid"]


def test_conditional_routing():
    from app.nodes import node_classify, route_by_task

    state = AgentState(goal="帮我计算 1+1")
    state = node_classify(state)
    assert route_by_task(state) == "tool"

    state2 = AgentState(goal="随便聊聊天")
    state2 = node_classify(state2)
    assert route_by_task(state2) == "chat"

运行:

代码语言:javascript
复制
pytest -v

十八、生产化清单

  • □ 接入真实 LLM(OpenAI / Claude / 本地)
  • □ 支持 function calling 而非规则模拟
  • □ 状态持久化到 PostgreSQL
  • □ 检查点写入独立表并加索引
  • □ 支持时间旅行调试(任意检查点恢复)
  • □ 节点级超时与重试
  • □ 循环最大步数与成本上限
  • □ 并行节点使用线程池 / asyncio
  • □ 汇聚节点处理部分失败
  • □ HITL 接入工单 / 飞书 / 钉钉
  • □ 子图支持独立版本管理
  • □ 多智能体通信协议(消息总线)
  • □ 全链路追踪(OpenTelemetry)
  • □ 成本与 Token 统计
  • □ 内容安全过滤
  • □ 权限与数据边界
  • □ 灰度发布与回滚
  • □ 图可视化(前端渲染 Mermaid)
  • □ 单元测试 + 集成测试 + 回归测试
  • □ 失败案例归档与复盘

十九、常见陷阱

  1. 用 Loop 硬写复杂流程:状态不可见,调试地狱;
  2. 用 Graph 做简单任务:过度设计,维护成本高;
  3. 状态无版本:升级后旧检查点无法恢复;
  4. 检查点不持久化:进程一挂全丢;
  5. 并行节点无汇聚:数据竞争;
  6. 循环无上限:Token 爆炸;
  7. HITL 无超时:流程永久挂起;
  8. 子图无边界:状态污染;
  9. 缺少追踪:出错无法定位;
  10. 不做测试:图一改就崩。

二十、总结

Agent Loop 与 Graph Engineer 的关系:

代码语言:javascript
复制
Loop  = 微观执行(一个节点内的推理循环)
Graph = 宏观编排(多个节点间的状态流转)

核心公式:

代码语言:javascript
复制
Agent System = Graph(State, Nodes, Edges, Checkpoints, Loop)

三条工程原则:

  1. 状态显式:所有中间结果可序列化、可恢复;
  2. 控制流显式:分支、并行、循环都在图里;
  3. 执行可观测:每一步都有检查点和追踪。

适用场景:

场景

推荐

简单工具调用

纯 Agent Loop

多步业务流程

Graph 编排

高风险操作

Graph + HITL

多智能体协作

Graph + 子图

长任务

Graph + 检查点

需要调试

Graph + 追踪

当你能把业务逻辑画成一张图,并让每个节点内部跑一个可控的 Agent Loop 时,你就真正掌握了下一代智能体系统的工程方法。


附:完整项目结构

代码语言:javascript
复制
agent-graph-engine/
├── app/
│   ├── __init__.py
│   ├── main.py
│   ├── db.py
│   ├── state.py
│   ├── loop.py
│   ├── tools.py
│   ├── graph.py
│   ├── nodes.py
│   ├── checkpoint.py
│   ├── hitl.py
│   ├── subgraph.py
│   ├── observability.py
│   └── llm.py
├── tests/
│   └── test_agent_graph.py
├── requirements.txt
└── Dockerfile

运行顺序:

代码语言:javascript
复制
pip install -r requirements.txt
pytest -v
uvicorn app.main:app --reload

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 目录
  • 一、为什么需要 Agent Loop
  • 二、为什么需要 Graph Engineer
  • 三、Agent Loop 与 Graph 的关系
  • 四、Agent Loop 核心原理与实现
    • 4.1 基本循环
    • 4.2 变体
    • 4.3 循环控制要点
  • 五、Graph Engineer 核心原理与实现
    • 5.1 图的基本构成
    • 5.2 图的执行模型
    • 5.3 节点类型
    • 5.4 边的类型
  • 六、关键机制:状态、节点、边、检查点
    • 6.1 状态
    • 6.2 节点
    • 6.3 边
    • 6.4 检查点
  • 七、代码实战:从 Loop 到 Graph 的完整系统
    • 7.1 技术栈
    • 7.2 项目结
    • 7.3 依赖
  • 八、模块一:状态与消息模型
  • 九、模块二:Agent Loop 引擎
  • 十、模块三:工具系统
  • 十一、模块四:Graph 编排引擎
  • 十二、模块五:条件边与并行节点
  • 十三、模块六:检查点与恢复
  • 十四、模块七:人机协作(HITL)
  • 十五、模块八:子图与多智能体
  • 十六、模块九:可观测性与审计
  • 十七、模块十:API 与测试
  • 十八、生产化清单
  • 十九、常见陷阱
  • 二十、总结
  • 附:完整项目结构
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档