从单Agent到Multi-Agent:共享记忆层的设计与实现

agent高级
AI Engineer Roadmap2026年08月03日

目标读者:已有单Agent开发经验,正在探索Multi-Agent架构扩展的开发者。


一、引言:从"一个大脑"到"一群大脑"

如果你已经用 LangChain、LlamaIndex 或自研框架搭建过单Agent系统,你一定熟悉这个循环:

感知(Observation)→ 检索(Retrieval)→ 推理(Reasoning)→ 行动(Action)→ 记忆(Memory)

单Agent的记忆层设计相对简单:一个向量数据库、一个对话历史缓冲区、一个工具调用记录,Agent自己读写即可。

但当你尝试把系统扩展为 Multi-Agent 架构时,问题开始变得复杂:

  • Agent A 刚完成的分析结论,Agent B 如何获取?
  • 两个Agent同时更新同一份知识,会不会互相覆盖?
  • 敏感信息(如用户隐私)应该让哪些Agent可见?
  • 短期任务上下文和长期领域知识,应该放在同一个存储里吗?

Multi-Agent 不是单Agent的简单堆叠,记忆层的共享、隔离与同步才是架构的关键瓶颈。 本文将从工程实践角度,拆解共享记忆层的设计要点与实现方案。


二、核心挑战:为什么记忆层是瓶颈?

在单Agent中,记忆是"私有财产"。在Multi-Agent中,记忆变成了"公共资源",由此带来三类核心挑战:

挑战类型具体问题影响
一致性Agent A 写入新结论,Agent B 读取到旧版本决策基于过时信息,产生逻辑矛盾
并发性多个Agent同时读写向量数据库数据竞争、向量索引损坏、查询结果不一致
安全性所有Agent共享同一套记忆,无权限边界敏感数据泄露、越权访问

这些问题不会在你跑 Demo 时暴露,但一旦进入生产环境、Agent 数量超过 3 个、QPS 超过 10,它们就会成为系统的阿喀琉斯之踵。


三、记忆分层模型:短期工作记忆 vs 长期知识库

解决上述问题的第一步,是把"记忆"这个概念拆细。我们借鉴认知科学的双存储模型,将Agent记忆分为两层:

3.1 短期工作记忆(Short-Term Working Memory, STWM)

特征:

  • 生命周期:随单次任务或单次会话存在,任务结束后清理或归档
  • 内容:当前对话上下文、多轮协商的中间结论、正在进行的工具调用链
  • 访问模式:高频读写、低延迟要求、强一致性要求
  • 存储介质:内存(Redis/Memcached)、进程内字典、或本地 SQLite

设计要点:

STWM 结构示例:
{
  "session_id": "task_20240803_001",
  "agents_involved": ["planner", "coder", "reviewer"],
  "context_window": [
    {"role": "planner", "content": "目标:实现用户登录接口", "timestamp": "..."},
    {"role": "coder", "content": "已生成 FastAPI 代码...", "timestamp": "..."}
  ],
  "shared_state": {
    "current_file": "auth.py",
    "pending_review": true
  },
  "ttl": 3600  // 1小时后自动清理
}

3.2 长期知识库(Long-Term Knowledge Base, LTKB)

特征:

  • 生命周期:持久化存储,跨任务、跨会话、跨Agent共享
  • 内容:领域知识文档、历史任务总结、用户偏好画像、代码规范
  • 访问模式:低频写入、高频读取、最终一致性可接受
  • 存储介质:向量数据库(Milvus/Pinecone/Weaviate)、图数据库(Neo4j)、对象存储

设计要点:

  • 采用 Write-Ahead Log (WAL) 机制保证写入可靠性
  • 向量索引与原始数据分离,支持异步重建
  • 按领域/项目/权限做 Namespace 隔离

3.3 两层之间的数据流动

┌─────────────────────────────────────────────────────────────┐
│                      Multi-Agent 系统                        │
│  ┌─────────────┐    ┌─────────────┐    ┌─────────────┐     │
│  │   Agent A   │◄──►│   Agent B   │◄──►│   Agent C   │     │
│  └──────┬──────┘    └──────┬──────┘    └──────┬──────┘     │
│         │                  │                  │             │
│         └──────────────────┼──────────────────┘             │
│                            ▼                                │
│              ┌─────────────────────────┐                    │
│              │   短期工作记忆 (STWM)    │  ◄── 实时同步      │
│              │      Redis/Memory       │                    │
│              └───────────┬─────────────┘                    │
│                          │ 任务结束/归档                     │
│                          ▼                                  │
│              ┌─────────────────────────┐                    │
│              │   长期知识库 (LTKB)      │  ◄── 异步写入      │
│              │   向量数据库 + 图数据库   │                    │
│              └─────────────────────────┘                    │
└─────────────────────────────────────────────────────────────┘

关键设计决策:

  • STWM 是"热数据",所有Agent通过 Pub/Sub 或共享内存实时同步
  • LTKB 是"冷数据",由专门的"记忆管理员Agent"或后台任务定期归档、去重、向量化
  • Agent 启动时从 LTKB 加载领域知识,运行中只读写 STWM

四、Agent间记忆冲突解决

当多个Agent同时操作共享记忆时,冲突不可避免。我们需要一套冲突检测与解决机制。

4.1 冲突场景分类

  1. 写-写冲突:Agent A 和 Agent B 同时更新同一条记忆
  2. 读-写冲突:Agent A 正在读取时,Agent B 修改了该记忆
  3. 语义冲突:Agent A 写入"使用 MySQL",Agent B 写入"使用 PostgreSQL"——技术上不冲突,但业务逻辑矛盾

4.2 乐观锁 + 向量版本控制

对于结构化记忆(如 JSON 状态),采用乐观锁:

class SharedMemory:
    def __init__(self, redis_client):
        self.redis = redis_client
    
    def update_state(self, session_id: str, agent_id: str, 
                     key: str, value: Any, expected_version: int) -> bool:
        """
        乐观锁更新:只有版本号匹配时才写入
        """
        lock_key = f"memory:{session_id}:lock"
        version_key = f"memory:{session_id}:version"
        
        # 获取分布式锁(防止并发写)
        with self.redis.lock(lock_key, timeout=5):
            current_version = self.redis.get(version_key) or 0
            if int(current_version) != expected_version:
                # 版本冲突,返回失败,由上层重试或合并
                return False
            
            # 写入新值并递增版本
            state_key = f"memory:{session_id}:state"
            self.redis.hset(state_key, key, json.dumps(value))
            self.redis.incr(version_key)
            
            # 广播变更事件给其他Agent
            self.redis.publish(f"channel:{session_id}", json.dumps({
                "agent": agent_id,
                "key": key,
                "value": value,
                "version": expected_version + 1
            }))
            return True

4.3 语义冲突解决:三向合并(Three-Way Merge)

对于向量记忆(如文本片段),版本号无法解决语义冲突。我们引入 三向合并 策略:

from typing import List, Dict
import hashlib

class SemanticMemoryMerger:
    def __init__(self, llm_client):
        self.llm = llm_client
    
    def three_way_merge(self, 
                        base: str,           # 冲突前的原始版本
                        agent_a_version: str, # Agent A 的修改
                        agent_b_version: str  # Agent B 的修改
                       ) -> str:
        """
        当两个Agent对同一段记忆产生语义冲突时,
        由 LLM 作为"仲裁者"进行智能合并。
        """
        prompt = f"""你是一名记忆冲突解决专家。以下是一段共享记忆的三个版本:

【原始版本】
{base}

【Agent A 的修改】
{agent_a_version}

【Agent B 的修改】
{agent_b_version}

请分析两者的修改意图,合并为一个最准确、无矛盾的版本。
如果存在不可调和的矛盾,优先保留更具体、更可验证的信息。
直接输出合并后的文本,不要解释。"""

        merged = self.llm.complete(prompt)
        return merged.strip()
    
    def detect_conflict(self, entries: List[Dict]) -> List[Dict]:
        """
        基于向量相似度和时间戳检测潜在语义冲突
        """
        conflicts = []
        for i, entry_a in enumerate(entries):
            for entry_b in entries[i+1:]:
                # 1. 检查是否涉及同一主题(高向量相似度)
                similarity = self.cosine_similarity(
                    entry_a['embedding'], 
                    entry_b['embedding']
                )
                # 2. 检查时间窗口(1小时内)
                time_diff = abs(entry_a['timestamp'] - entry_b['timestamp'])
                
                if similarity > 0.85 and time_diff < 3600:
                    # 3. 检查内容是否矛盾(用LLM判断)
                    if self.is_contradictory(entry_a['content'], entry_b['content']):
                        conflicts.append({
                            'entry_a': entry_a,
                            'entry_b': entry_b,
                            'similarity': similarity
                        })
        return conflicts

工程建议:

  • 对于高频冲突场景(如 3 个以上Agent同时编辑),引入 "记忆仲裁Agent" 专门负责合并
  • 对于低频冲突,由冲突检测器自动触发合并,失败时人工介入

五、向量存储的并发读写

向量数据库(如 Milvus、Pinecone、Weaviate)的并发模型与传统关系型数据库不同,需要特别关注。

5.1 向量数据库的并发特性

数据库写入并发读取并发索引更新策略
Milvus段级锁MVCC后台异步构建
Pinecone全局队列高并发实时更新
Weaviate对象级锁MVCC批量索引

核心问题:向量索引的重建是计算密集型操作,大量并发写入会导致:

  • 查询延迟飙升(索引碎片化)
  • 写入排队(锁竞争)
  • 向量距离计算结果不一致(索引版本不同)

5.2 写入缓冲 + 批量刷新模式

import asyncio
from collections import deque
from typing import List, Dict
import numpy as np

class BufferedVectorStore:
    def __init__(self, vector_db, batch_size: int = 100, 
                 flush_interval: float = 5.0):
        self.db = vector_db
        self.batch_size = batch_size
        self.flush_interval = flush_interval
        self.buffer = deque()
        self.lock = asyncio.Lock()
        self._start_flush_loop()
    
    async def add(self, agent_id: str, texts: List[str], 
                  embeddings: List[List[float]], metadata: List[Dict]):
        """
        Agent 写入时先进入内存缓冲,而非直接写向量库
        """
        async with self.lock:
            for text, emb, meta in zip(texts, embeddings, metadata):
                self.buffer.append({
                    'agent_id': agent_id,
                    'text': text,
                    'embedding': emb,
                    'metadata': meta,
                    'timestamp': asyncio.get_event_loop().time()
                })
            
            # 达到批量阈值立即刷新
            if len(self.buffer) >= self.batch_size:
                await self._flush()
    
    async def _flush(self):
        """批量写入向量数据库,减少索引重建次数"""
        if not self.buffer:
            return
        
        async with self.lock:
            batch = list(self.buffer)
            self.buffer.clear()
        
        # 按 namespace 分组写入(隔离不同Agent的数据)
        grouped = {}
        for item in batch:
            ns = item['metadata'].get('namespace', 'default')
            grouped.setdefault(ns, []).append(item)
        
        for namespace, items in grouped.items():
            texts = [i['text'] for i in items]
            embeddings = [i['embedding'] for i in items]
            metadatas = [i['metadata'] for i in items]
            
            # 批量写入,单次触发索引更新
            await self.db.upsert(
                vectors=embeddings,
                texts=texts,
                metadatas=metadatas,
                namespace=namespace
            )
    
    def _start_flush_loop(self):
        """定时刷新,避免数据滞留"""
        async def loop():
            while True:
                await asyncio.sleep(self.flush_interval)
                await self._flush()
        asyncio.create_task(loop())
    
    async def search(self, agent_id: str, query_embedding: List[float], 
                     namespace: str = None, top_k: int = 5) -> List[Dict]:
        """
        查询时合并缓冲区和持久化数据的结果
        """
        # 1. 查询持久化向量库
        db_results = await self.db.search(
            vector=query_embedding,
            namespace=namespace,
            top_k=top_k * 2  # 多取一些,后续合并
        )
        
        # 2. 查询内存缓冲区(未刷新的最新数据)
        buffer_results = []
        async with self.lock:
            for item in self.buffer:
                if namespace and item['metadata'].get('namespace') != namespace:
                    continue
                # 计算与查询向量的相似度
                sim = self._cosine_sim(query_embedding, item['embedding'])
                buffer_results.append({
                    'text': item['text'],
                    'metadata': item['metadata'],
                    'score': sim
                })
        
        # 3. 合并去重,按相似度排序
        all_results = db_results + buffer_results
        all_results.sort(key=lambda x: x['score'], reverse=True)
        
        # 去重(基于文本哈希)
        seen = set()
        unique = []
        for r in all_results:
            h = hashlib.md5(r['text'].encode()).hexdigest()
            if h not in seen:
                seen.add(h)
                unique.append(r)
                if len(unique) >= top_k:
                    break
        
        return unique

关键收益:

  • 将随机小写入聚合成批量写入,减少索引重建次数 10 倍以上
  • 查询时合并缓冲区,保证"写入即可见"的语义
  • 按 Namespace 隔离,避免不同Agent的数据互相污染索引

5.3 读写分离架构

对于高并发场景,建议采用读写分离:

┌─────────────┐     ┌─────────────┐     ┌─────────────┐
│   Agent A   │     │   Agent B   │     │   Agent C   │
└──────┬──────┘     └──────┬──────┘     └──────┬──────┘
       │                   │                   │
       └─────────┬─────────┴─────────┬─────────┘
                 ▼                   ▼
        ┌──────────────┐    ┌──────────────┐
        │  写入缓冲队列  │    │   读取副本    │
        │  (Kafka/Redis)│    │  (只读向量库) │
        └──────┬───────┘    └──────┬───────┘
               │                   │
               ▼                   ▼
        ┌──────────────┐    ┌──────────────┐
        │  批量写入器   │    │   查询服务    │
        │  (单线程)     │    │  (多线程)     │
        └──────┬───────┘    └──────────────┘
               │
               ▼
        ┌──────────────┐
        │  主向量数据库  │
        │  (Milvus等)  │
        └──────────────┘

六、记忆权限控制

在 Multi-Agent 系统中,并非所有记忆都应该对所有Agent开放。权限控制是生产环境的必选项。

6.1 权限模型设计

我们采用 RBAC + ABAC 混合模型:

  • RBAC(基于角色):定义Agent角色(如 Planner、Coder、Reviewer),角色决定基础权限
  • ABAC(基于属性):根据记忆内容的敏感级别、所属项目、创建者等属性动态判定
from enum import Enum
from typing import Set, List

class Permission(Enum):
    READ = "read"
    WRITE = "write"
    DELETE = "delete"
    ADMIN = "admin"

class AgentRole(Enum):
    PLANNER = "planner"
    CODER = "coder"
    REVIEWER = "reviewer"
    ORCHESTRATOR = "orchestrator"

class SensitivityLevel(Enum):
    PUBLIC = 0      # 所有Agent可见
    INTERNAL = 1    # 同项目Agent可见
    CONFIDENTIAL = 2 # 仅特定角色可见
    SECRET = 3      # 仅创建者和 orchestrator 可见

class MemoryACL:
    def __init__(self):
        # 角色基础权限
        self.role_permissions = {
            AgentRole.PLANNER: {Permission.READ, Permission.WRITE},
            AgentRole.CODER: {Permission.READ, Permission.WRITE},
            AgentRole.REVIEWER: {Permission.READ},
            AgentRole.ORCHESTRATOR: {Permission.READ, Permission.WRITE, 
                                     Permission.DELETE, Permission.ADMIN}
        }
    
    def check_access(self, agent_id: str, agent_role: AgentRole,
                     memory_entry: Dict, requested_perm: Permission) -> bool:
        """
        检查Agent是否有权访问某条记忆
        """
        # 1. 角色权限检查
        if requested_perm not in self.role_permissions.get(agent_role, set()):
            return False
        
        # 2. 敏感级别检查
        sensitivity = memory_entry.get('sensitivity', SensitivityLevel.PUBLIC)
        owner = memory_entry.get('created_by')
        
        if sensitivity == SensitivityLevel.SECRET:
            # 仅创建者和 orchestrator 可访问
            if agent_id != owner and agent_role != AgentRole.ORCHESTRATOR:
                return False
        
        elif sensitivity == SensitivityLevel.CONFIDENTIAL:
            # 检查Agent是否在允许的列表中
            allowed = memory_entry.get('allowed_agents', [])
            if agent_id not in allowed and agent_role != AgentRole.ORCHESTRATOR:
                return False
        
        elif sensitivity == SensitivityLevel.INTERNAL:
            # 检查是否同项目
            agent_project = self._get_agent_project(agent_id)
            memory_project = memory_entry.get('project_id')
            if agent_project != memory_project:
                return False
        
        # PUBLIC 无需额外检查
        
        return True
    
    def _get_agent_project(self, agent_id: str) -> str:
        # 从注册中心获取Agent所属项目
        pass

6.2 存储层权限隔离

权限控制不能只停留在应用层,必须在存储层落实:

方案一:Namespace 隔离(推荐)

# 每个项目/权限组一个 Namespace
namespace_map = {
    "project_alpha": "ns_project_alpha",
    "project_beta": "ns_project_beta", 
    "confidential": "ns_confidential"
}

# 查询时自动路由到对应 Namespace
def get_namespace(agent_id: str, sensitivity: SensitivityLevel) -> str:
    project = get_agent_project(agent_id)
    if sensitivity == SensitivityLevel.CONFIDENTIAL:
        return f"ns_confidential_{project}"
    return f"ns_{project}"

方案二:元数据过滤

在向量数据库不支持 Namespace 时,使用元数据过滤:

# 查询时自动注入权限过滤条件
def build_filter(agent_id: str, agent_role: AgentRole) -> Dict:
    return {
        "$or": [
            {"sensitivity": "PUBLIC"},
            {"sensitivity": "INTERNAL", "project_id": get_agent_project(agent_id)},
            {"allowed_agents": {"$contains": agent_id}},
            {"created_by": agent_id}
        ]
    }

# Milvus 示例
results = milvus_client.search(
    collection_name="shared_memory",
    data=[query_embedding],
    filter=build_filter(agent_id, agent_role),
    limit=10
)

6.3 审计日志

所有记忆访问都应记录审计日志,用于事后追溯:

class MemoryAuditLogger:
    def log_access(self, agent_id: str, memory_id: str, 
                   action: str, success: bool, reason: str = None):
        log_entry = {
            "timestamp": datetime.utcnow().isoformat(),
            "agent_id": agent_id,
            "memory_id": memory_id,
            "action": action,
            "success": success,
            "reason": reason,
            "trace_id": get_current_trace_id()  # 关联分布式追踪
        }
        # 写入审计日志队列(如 Kafka)
        self.audit_queue.send(log_entry)

七、完整架构示例

综合以上设计,一个生产级的共享记忆层架构如下:

┌─────────────────────────────────────────────────────────────────┐
│                        Multi-Agent 协调层                        │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐  ┌──────────┐        │
│  │ Planner  │  │  Coder   │  │ Reviewer │  │  Tools   │        │
│  └────┬─────┘  └────┬─────┘  └────┬─────┘  └────┬─────┘        │
│       │             │             │             │               │
│       └─────────────┴──────┬──────┴─────────────┘               │
│                            ▼                                    │
│              ┌─────────────────────────┐                        │
│              │    记忆访问网关 (Gateway)  │  ◄── 权限校验、审计    │
│              │    - ACL 检查              │                        │
│              │    - 冲突检测              │                        │
│              │    - 路由分发              │                        │
│              └───────────┬─────────────┘                        │
│                          │                                      │
│         ┌────────────────┼────────────────┐                     │
│         ▼                ▼                ▼                     │
│  ┌──────────────┐ ┌──────────────┐ ┌──────────────┐            │
│  │  短期工作记忆  │ │  长期知识库   │ │   审计日志    │            │
│  │   (Redis)    │ │ (Milvus+RDB) │ │  (ClickHouse)│            │
│  │              │ │              │ │              │            │
│  │ • Pub/Sub   │ │ • 向量索引    │ │ • 访问记录    │            │
│  │ • 状态锁    │ │ • 知识图谱    │ │ • 异常告警    │            │
│  │ • TTL 清理  │ │ • 全文检索    │ │ • 合规报表    │            │
│  └──────────────┘ └──────────────┘ └──────────────┘            │
│                                                                 │
│  ┌─────────────────────────────────────────────────────────┐   │
│  │              后台服务                                    │   │
│  │  • 记忆归档器:STWM → LTKB 定期迁移                      │   │
│  │  • 索引优化器:向量索引重建、去重、压缩                   │   │
│  │  • 冲突仲裁器:语义合并、人工介入队列                     │   │
│  └─────────────────────────────────────────────────────────┘   │
└─────────────────────────────────────────────────────────────────┘

八、最佳实践与总结

8.1 关键设计原则

  1. 分层存储:短期记忆走内存/Redis,长期记忆走向量库,不要混用
  2. 批量写入:向量数据库的写入永远批量做,单条写入是性能杀手
  3. 乐观优先:先用乐观锁解决 90% 的冲突,剩余 10% 用 LLM 仲裁
  4. 权限前置:在存储查询层面做过滤,不要只在应用层做
  5. 可观测性:所有记忆操作都要有 Trace ID,方便定位"Agent 为什么做了这个决策"

8.2 避坑指南

坑后果解决方案
所有Agent共享一个向量 Collection权限无法隔离、索引互相污染按项目/权限分 Namespace
直接让Agent写向量库并发冲突、索引碎片化加写入缓冲层,批量刷新
忽略记忆TTL内存爆炸、查询变慢STWM 必须设 TTL,LTKB 定期归档
不做版本控制无法追溯、无法回滚所有记忆变更带版本号

8.3 总结

从单Agent到Multi-Agent,最大的认知转变是:记忆不再是Agent的私有属性,而是系统的共享基础设施。

一个好的共享记忆层应该像"分布式大脑的外层皮层":

  • 让Agent们实时感知彼此的进展(STWM)
  • 让知识持久沉淀并跨任务复用(LTKB)
  • 在冲突时智能仲裁而非简单粗暴覆盖
  • 对敏感信息精细管控而非一刀切隔离

如果你正在从单Agent向Multi-Agent演进,建议先实现 STWM + 乐观锁 + Namespace 隔离 这三个基础能力,它们能解决 80% 的生产问题。在此基础上,再逐步引入语义合并、读写分离、审计日志等高级特性。


延伸阅读: