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 产出) └─────────────────────────┘
|
在深度自主任务(如 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/