从单Agent到Multi-Agent:共享记忆层的设计与实现
目标读者:已有单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 冲突场景分类
- 写-写冲突:Agent A 和 Agent B 同时更新同一条记忆
- 读-写冲突:Agent A 正在读取时,Agent B 修改了该记忆
- 语义冲突: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 关键设计原则
- 分层存储:短期记忆走内存/Redis,长期记忆走向量库,不要混用
- 批量写入:向量数据库的写入永远批量做,单条写入是性能杀手
- 乐观优先:先用乐观锁解决 90% 的冲突,剩余 10% 用 LLM 仲裁
- 权限前置:在存储查询层面做过滤,不要只在应用层做
- 可观测性:所有记忆操作都要有 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% 的生产问题。在此基础上,再逐步引入语义合并、读写分离、审计日志等高级特性。
延伸阅读: