0%

Agent Orchestration - 智能体编排

1. Eino 如何实现多智能体

1
2
3
4
5
6
7
8
9
                           ┌────────────────────────────────────────────────────────┐
│ Eino 多智能体协同体系 (Multi-Agent) │
└──────────────────────────┬─────────────────────────────┘

┌────────────────────────────────────────────┼────────────────────────────────────────────┐
▼ ▼ ▼
【范式 1: Host Multi-Agent】 【范式 2: Agent as a Tool】 【范式 3: State Graph 协同】
• 机制: 意图识别分发 + Summarizer 聚合 • 机制: 主控 Agent 将 Sub-Agent 包装为 Tool • 机制: 强类型全局状态机 + 条件边流转
• 适用: 专家分流 (客服/风控/导购) • 适用: 深度自主规划、跨领域能力委派 • 适用: 复杂闭环业务、自反思、Saga 事务

1.1 Host-Specialist 模式(基于 Flow 集成)

Host Multi-Agent 是 Eino 官方封装的高内聚开箱模式。
Host 负责理解用户 Query 并做意图路由,分发给一个或多个 Specialist Agent,最后由可选的 Summarizer 模块聚合成最终流式输出。

架构核心机制:

  • 单/多 Specialist 动态激活:Host 可判定单路由转发,也可同时 Fan-Out 派发给多个 Specialist。
  • 流式平滑拼接:底层自动处理多个 Specialist 的 StreamChunk 合并,未配置 Summarizer 时默认拼接文本,配置后由 Summarizer 模型总结后下发。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
                    ┌─────────────────────────┐
│ User Query │
└────────────┬────────────┘

┌─────────────────────────┐
│ Host Agent (Router) │ ──> (识别意图并路由)
└──────┬───────────┬──────┘
│ │
┌───────────────┘ └───────────────┐
▼ ▼
┌───────────────────────┐ ┌───────────────────────┐
│ Specialist A (售后退款)│ │ Specialist B (商品推荐)│
└──────────┬────────────┘ └──────────┬────────────┘
│ │
└───────────────────┬───────────────────────┘

┌─────────────────────────┐
│ Summarizer (聚合响应) │ ──> (合并多 Agent 产出)
└─────────────────────────┘

1.2 Agent-as-a-Tool / 层级委派(基于 ADK)

在深度自主任务(如 DeepAgent / 复杂分析)中,主控 Agent 拥有全局目标,遇到具体执行动作时,通过调用封装为 tool.BaseTool 的子 Agent。

架构核心机理:

  • 上下文截断与降噪:主 Agent 的上下文不应被子 Agent 的长推理链(ReAct Loop)污染。子 Agent 在独立的 Context 空间内运行,执行完成后仅将最终摘要/结构化结果作为 Tool Output 返回给主 Agent。
  • 控制权转移(Hand-off):在 Eino ADK 中,通过设置 Agent 的转移链(Transfer/SubAgents),支持主 Agent 在识别到特定专业阶段后,彻底将交互控制权移交给子 Agent 接管,待子任务结束后再回弹。

基于 State Graph 的强类型多 Agent 协同

在复杂度最高、一致性要求最严的电商/金融场景中,应直接采用 Eino-Compose Graph 构建基于全局黑板(Blackboard State)的状态机协同。

1
2
3
4
5
6
7
8
9
10
11
12
(Start) ──> [Node: Intent_Router]

┌───────────────┴───────────────┐
▼ ▼
[Node: Sales_Agent] [Node: Support_Agent]
│ │
└───────────────┬───────────────┘

[Node: Medical_Critic]

├── (Score < 0.8 && HopCount < 3) ──> (回到 Agent 反思重写)
└── (Score >= 0.8) ──> [Node: Stream_Responder] ──> (End)

状态流转示意图

1
2
3
4
5
6
7
8
9
10
11
12
13
// 添加评价节点的条件路由边
graph.AddBranch("Medical_Critic", func(ctx context.Context, state *SessionState) (string, error) {
state.HopCount++
// 死循环熔断
if state.HopCount > 3 {
return "Stream_Responder", nil
}
// 质检不合格,反思重试
if state.CriticScore < 0.8 {
return "Sales_Agent", nil
}
return "Stream_Responder", nil
})

2. 多 Agent 协同四大工程防线

2.1 跨 Agent 上下文隔离和 token 剪枝

痛点:若每个 Agent 都能看到全部交互与所有工具原始报文,多轮多 Agent 交互后 Token 迅速暴涨。

Eino 架构解法:在 Graph 节点间引入 State Reducer / Masking 机制。

  • Sales_Agent 仅消费 UserQuery + UserProfile;
  • Critic_Agent 仅消费 DraftAnswer + Medical_Rules;
  • 状态持久化存储在 Redis,但注入大模型推理上下文的切片进行按需动态投影。

2.2 全链路 Composable Streaming(流式透传与背压)

痛点:多 Agent 串行或嵌套调用时,中间节点的阻塞会导致前端长时间白屏,TTFT 严重劣化。

Eino 架构解法:

  • 利用 schema.StreamReader[T] 抽象,下游节点无需等待上游完全生成,而是以 Chunk 为单位进行流式订阅;
  • 在路由节点使用 Tee 分流,一路流式输出给前端做打字机渲染,一路流式喂给后置安全审计拦截器并行校验。

2.3 中断与人工接管(Human-in-the-loop / Checkpointing)

Eino ADK 中断恢复机制:

  • 涉及敏感交易或高危操作(如发放高额优惠券、退款审批)时,Agent 执行节点调用 adk.Interrupt() 暂停当前执行图,将完整的 State 序列化持久化至存储;
  • 外部运营/用户在审批系统确认后,通过 adk.Resume(taskID, payload) 从中断点精准恢复执行,避免从头重新推理。

adk.Interrupt() 会触发下游的 checkPointStore 进行数据存储

2.4 切面观测与 Parent-Child Span 治理(Aspect Tracing)

Eino Aspect 机制:利用无侵入拦截器,在多 Agent 协同流转中自动维护 W3C TraceContext:

  • Root Span (Graph Run)
  • Child Span (Host Router)
  • Child Span (Specialist Sub-Agent)
  • Grandchild Span (Tool Calling & LLM Inference)

精准归因每个 Agent 的延迟占比、Token 消耗及错误堆栈,直接对接内部 OpenTelemetry 大盘。

3. 数据中断与恢复流程

  • 状态模式演进与版本漂移
    • 历史数据结构不兼容,用兼容的数据结构或序列化组件
    • 在 checkpoint 上打标版本号
  • 任务超时与清理
    • 配置超时时间
    • 定时任务扫描标记过期
  • 防止唤醒脑裂
    • 通过 Redis 分布式锁
    • 更新的时候带上状态校验 CAS
  • 敏感数据处理
    • 敏感数据加密处理
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
[1. Agent 节点执行] ──> 触发 adk.Interrupt(payload)


[2. Checkpointer 拦截] ──> 抽取当前 Blackboard State

├── (State > 64KB) ──> 上传 S3 并获取 Object_URI

├── 写入 MySQL (Status = SUSPENDED, 记录挂起节点与审批摘要)
└── 缓存至 Redis (Key: task_id, TTL: 24h)


[3. 释放 Worker 资源] ──> Worker 协程销毁退出,向审批工作台 / Webhook 发送待办事件

[ 人工审批/修改参数 ]


[4. 触发 Resume API] ──> POST /api/v1/tasks/{task_id}/resume


[5. 图引擎状态重建] ──> 加载分布式锁 ──> 读取最新 Checkpoint ──> 反序列化 State

├── 注入 Human Feedback (覆盖/合并修改字段)
├── 更新 Checkpoint 状态为 RESUMED
└── 从 suspended node 的下游条件边继续触发图流转

4. 多智能体架构示意

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
+-----------------------------------------------------------------------------------+
| 1. 接入层 (API Gateway) |
| 电商 App / 商家后台 / 客服工作台 (多端接入) | SSE/WebSocket 流式传输 | 鉴权 & 限流 |
+-----------------------------------------------------------------------------------+

+----------------------------------------▼------------------------------------------+
| 2. 编排与调度中枢 (Agent Orchestration & Engine) |
| ┌──────────────────┐ ┌─────────────────────────┐ ┌──────────────────────────┐ |
| │ 意图识别/Router │ │ 状态机 (LangGraph/DAG) │ │ Planning & Reflection │ |
| └─────────┬────────┘ └───────────┬─────────────┘ └────────────┬─────────────┘ |
| │ │ │ |
| ┌─────────▼───────────────────────▼─────────────────────────────▼─────────────┐ |
| │ Multi-Agent 协作层: 客服 Agent / 售后 Agent / 商家 Copilot / 推荐 Agent │ |
| └─────────────────────────────────────────────────────────────────────────────┘ |
+-----------------------------------------------------------------------------------+
│ │
+-------------------▼------------------+ +-------------------▼-----------------+
| 3. 记忆与上下文引擎 (Memory Engine) | | 4. 能力与连接层 (Skills & MCP) |
| • 短时窗口管理 (Sliding Window/Pruning)| | • MCP Client / Server 标准接入协议 |
| • 长期记忆 (Vector DB / User Profile) | | • 业务 API: 订单/物流/退款/商品接口 |
| • 会话状态持久化 (Redis / MySQL) | | • RAG Pipeline: 向量检索 + 重排/分块 |
+--------------------------------------+ +-------------------------------------+
│ │
+-------------------▼---------------------------------------------▼-----------------+
| 5. 模型与推理网关 (Model & Infra Gateway) |
| 模型路由 (Claude / GPT / 自研开源模型) | Prompt 管理 | 语义缓存 & Prompt Cache |
+-----------------------------------------------------------------------------------+

+----------------------------------------▼------------------------------------------+
| 6. 评测、安全与可观测 (Eval & Observability) |
| 安全护栏 (Guardrails/注入防御) | Trace 链路追踪 (LangSmith/OTel) | 自动化评测 (Eval) |
+-----------------------------------------------------------------------------------+

5. 智能体记忆共享

多 Agent 共享记忆治理思路

  • 架构解耦:采用 CQRS 模式,写走 Append-Only 日志流保吞吐,读走‘Redis 热缓存 + 向量/图谱物化视图’的双轨读,兼顾写吞吐与写后读强一致;
  • 认知一致性:建立‘Slot 碰撞 → NLI 小模型 → LLM 仲裁’三级漏斗,以毫秒级确定性计算化解语义冲突;
  • 零信任治理:基于 ABAC 策略做向量检索算子下推,避免先搜后滤带来的召回截断与数据越界,并通过 Taint Analysis 阻断记忆投毒;
  • 成本与演进:参考 LSM 思想做双时态生命周期管理与分级 Compaction,将海量碎片 Trace 蒸馏为高阶全局认知,实现 Agent 系统的经验复利。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
[ Multi-Agents: Perception & Action ]
│ ▲
│ 1. Append (毫秒级 ACK) │ 4. Hybrid Read (双轨读)
▼ │
┌──────────────────────────────┐ ┌────────────────────────────────────────────────────────┐
│ Ingestion Gateway (WAL Log) │ │ Query Engine │
│ (Kafka / Raft Append-Only) │ │ ├─ Hot Path: Redis (Session/Working Memory) [强一致] │
└──────────────┬───────────────┘ │ └─ Cold Path: Hybrid Vector + Graph [最终一致] │
│ └───────────────────────────▲────────────────────────────┘
▼ 2. Async Dispatch │
┌─────────────────────────────────────────────────────────┐ │ 3. Materialized View Update
│ Background Async Reconciliation & Compaction Pipeline │──────┘
│ ├─ Fast Path: Key-Slot Hash & Vector Clock (确定性覆盖) │
│ ├─ Mid Path: NLI 轻量模型 (逻辑冲突/蕴含检测) │
│ ├─ Slow Path: LLM Semantic Merge (复杂因果仲裁) │
│ ├─ ABAC Pushdown Indexer (构建带权限掩码的物理索引) │
│ └─ Knowledge Compaction Engine (LSM-style 碎片聚合) │
└─────────────────────────────────────────────────────────┘

6. 参考文档

  • https://arthurchiao.art/blog/built-multi-agent-research-system-zh/