共享记忆与智能体间通信
使用键值存储实现共享记忆层,让智能体从中读取和写入数据,从而支持异步协作,同时避免智能体之间的紧密耦合。
共享记忆与智能体间通信 是 CoddyKit 上的免费 AI Engineering Academy 课时。 这是第 4 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 AI Engineering Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 AI Engineering Academy 课程共包含 4 节课。
代理隔离带来的问题
在多代理系统中,每个代理都在自己的上下文中运行,不知道其他代理正在做什么或已经完成了什么。如果研究代理发现了一个关键事实,写作代理如何得知?如果编码代理遇到错误,监控代理如何知道?没有共享通信机制时,代理就会成为彼此孤立的信息孤岛,无法有效协作。
共享内存:Blackboard 模型
多代理通信的经典解决方案是Blackboard 模型:一个任何代理都可以读取或写入的共享数据存储(Blackboard)。代理发布自己的发现,读取其他代理的贡献,并通过共享状态进行隐式协调。这种模型解耦了代理之间的关系——代理不需要知道彼此是否存在,只需要了解共享内存的结构。
# Simple in-memory blackboard using a dictionary
from threading import Lock
class Blackboard:
def __init__(self):
self._data = {}
self._lock = Lock() # thread-safe for parallel agents
def write(self, key: str, value, agent_id: str):
with self._lock:
self._data[key] = {'value': value, 'written_by': agent_id}
print(f'[{agent_id}] wrote: {key}')
def read(self, key: str):
with self._lock:
return self._data.get(key, {}).get('value')
def keys(self):
with self._lock:
return list(self._data.keys())
blackboard = Blackboard()使用 Redis 实现持久化共享内存
对于代理以独立进程或服务运行的生产级多代理系统,内存字典并不够用。Redis 是共享代理内存最受欢迎的选择:它速度快,支持丰富的数据类型(字符串、哈希、列表、有序集合),内置用于过期的 TTL,并通过原子操作安全地处理并发读写。
import redis
import json
class RedisSharedMemory:
def __init__(self, prefix='agent:'):
self.redis = redis.Redis(host='localhost', port=6379, decode_responses=True)
self.prefix = prefix
def set(self, key: str, value, ttl_seconds=3600):
full_key = self.prefix + key
self.redis.setex(full_key, ttl_seconds, json.dumps(value))
def get(self, key: str):
full_key = self.prefix + key
raw = self.redis.get(full_key)
return json.loads(raw) if raw else None
def append_to_list(self, key: str, item):
full_key = self.prefix + key
self.redis.rpush(full_key, json.dumps(item))
def get_list(self, key: str):
full_key = self.prefix + key
return [json.loads(x) for x in self.redis.lrange(full_key, 0, -1)]
memory = RedisSharedMemory(prefix='research_project:')为共享内存设置命名空间
在复杂的多代理系统中,代理会写入许多不同类型的数据,扁平的键命名空间很快就会变得混乱。请使用分层命名空间来清晰地组织共享内存。常见模式是 project_id:agent_role:data_type,例如 proj_123:researcher:findings 或 proj_123:coder:error_log。这样可以轻松查询某个项目的所有数据,或某个特定代理的全部输出。
class NamespacedMemory:
def __init__(self, project_id: str, agent_id: str, redis_client):
self.base = f'{project_id}:{agent_id}'
self.redis = redis_client
def write_finding(self, topic: str, content: str):
key = f'{self.base}:findings:{topic}'
self.redis.set(key, content)
def read_all_findings(self, project_id: str):
# Read findings from ALL agents in this project
pattern = f'{project_id}:*:findings:*'
keys = self.redis.keys(pattern)
return {k: self.redis.get(k) for k in keys}
# Usage
researcher_memory = NamespacedMemory('proj_123', 'researcher', redis_client)
researcher_memory.write_finding('competitors', 'OpenAI, Anthropic, Google...')
writer_memory = NamespacedMemory('proj_123', 'writer', redis_client)
all_findings = writer_memory.read_all_findings('proj_123')结构化内存与非结构化内存
共享内存中的内容可以是非结构化的(供 LLM 读取的原始文本块),也可以是结构化的(带有类型化字段的 JSON/Python 对象)。结构化内存更为理想,因为它支持程序化查询、验证和合并。请始终为每个代理写入共享内存的数据定义模式并记录说明,同时根据该模式验证写入内容,以防止某个有缺陷的代理破坏内存存储。
from pydantic import BaseModel
from typing import Optional, list
from datetime import datetime
class ResearchFinding(BaseModel):
topic: str
summary: str
sources: list[str]
confidence: float # 0.0 to 1.0
written_by: str
timestamp: datetime
# Validated write - bad data is caught before it enters shared memory
def write_finding(memory, finding_dict: dict):
finding = ResearchFinding(**finding_dict) # validates on creation
memory.set(f'findings:{finding.topic}', finding.model_dump())
print(f'Validated finding written for topic: {finding.topic}')使用发布/订阅进行事件驱动通信
代理无需轮询共享内存来检查更新,而是可以使用发布/订阅(发布/订阅)通信,在任务完成时通知其他代理。代理 A 发布一个事件(“research_complete”),订阅了该事件的代理 B 被唤醒并开始处理。Redis 发布/订阅以及 RabbitMQ 或 Kafka 等消息队列都支持这种模式。
import redis
# Publisher (researcher agent)
def researcher_agent(topic, redis_client):
findings = do_research(topic)
redis_client.set(f'findings:{topic}', findings)
# Notify all subscribers that research is done
redis_client.publish('agent_events', f'research_complete:{topic}')
print(f'Research complete, published event for topic: {topic}')
# Subscriber (writer agent) - runs in separate process
def writer_agent_listener(redis_client):
pubsub = redis_client.pubsub()
pubsub.subscribe('agent_events')
for message in pubsub.listen():
if message['type'] == 'message':
event = message['data']
if event.startswith('research_complete:'):
topic = event.split(':')[1]
findings = redis_client.get(f'findings:{topic}')
write_draft(findings) # start writing immediatelyLangGraph 中的共享内存
在 LangGraph 中,代理之间共享的内存本身就是图状态对象。每个节点都会从同一个类型化状态字典中读取数据并写入数据。LangGraph 会自动处理读写协调。对于更复杂的场景,您还可以通过依赖注入,将外部内存客户端(Redis、数据库)注入每个节点函数。
from langgraph.graph import StateGraph
from typing import TypedDict
class SharedState(TypedDict):
# All shared data lives here - every node can read any field
query: str
research_findings: str # written by researcher, read by writer
written_draft: str # written by writer, read by reviewer
review_notes: str # written by reviewer, read by writer (loop)
final_output: str # written by synthesizer
# Researcher writes to 'research_findings'
def researcher(state: SharedState) -> dict:
findings = search_and_summarize(state['query'])
return {'research_findings': findings} # partial state update
# Writer reads 'research_findings', writes 'written_draft'
def writer(state: SharedState) -> dict:
draft = write_from_findings(state['research_findings']) # reads researcher output
return {'written_draft': draft}内存冲突与一致性
当多个代理并发写入共享内存时,可能会发生写入冲突。两个代理可能会相互覆盖对方的工作,或者在读取和写入之间读到过时的数据。您可以采用以下方式处理:乐观锁(写入前检查版本)、在 Redis 中执行原子比较并交换操作,或者通过协调器代理串行化写入,让它成为关键内存字段的唯一写入者。
# Optimistic locking with Redis
def safe_write(redis_client, key, new_value, expected_version):
with redis_client.pipeline() as pipe:
try:
pipe.watch(key + ':version') # watch for concurrent modification
current_version = int(pipe.get(key + ':version') or 0)
if current_version != expected_version:
raise ValueError(f'Version conflict: expected {expected_version}, got {current_version}')
pipe.multi() # start transaction
pipe.set(key, new_value)
pipe.set(key + ':version', current_version + 1)
pipe.execute() # atomic commit
print('Write successful')
except redis.WatchError:
print('Conflict detected, retry write')内存 TTL 与清理
共享代理内存会随着时间不断累积,如果不加管理,可能无限增长。请始终为内存条目设置生存时间(TTL),使其自动过期。对于项目范围的内存,请在项目完成时清理所有条目。TTL 的取值应与工作流的预期持续时间相匹配:临时数据使用较短的 TTL(几分钟),可能会重复使用的结果使用较长的 TTL(几小时到几天)。
def cleanup_project_memory(redis_client, project_id: str):
pattern = f'{project_id}:*'
keys = redis_client.keys(pattern)
if keys:
redis_client.delete(*keys)
print(f'Cleaned up {len(keys)} memory entries for project {project_id}')
# Set TTL when writing
def write_with_ttl(redis_client, key, value, ttl_hours=2):
redis_client.setex(
key,
ttl_hours * 3600, # convert to seconds
json.dumps(value)
)
# Register cleanup callback when workflow completes
def on_workflow_complete(project_id):
cleanup_project_memory(redis_client, project_id)
print(f'Workflow {project_id} complete, memory cleaned up')作为代理历史记录的内存
共享内存不仅可以存储代理的操作历史,还可以记录代理的输出。记录哪个代理在何时、出于什么原因执行了什么操作,可以形成审计轨迹,这对于调试故障、了解最终输出的生成过程,以及恢复中断的工作流都非常有价值。此操作日志相当于基于内存实现的 LangSmith 跟踪。
import time
from dataclasses import dataclass
@dataclass
class AgentAction:
agent_id: str
action_type: str # 'research', 'write', 'review', 'tool_call'
input_summary: str
output_summary: str
timestamp: float
success: bool
def log_action(memory, action: AgentAction):
key = f'action_log:{action.agent_id}:{action.timestamp}'
memory.set(key, vars(action))
# Usage in an agent
def researcher_with_logging(state, memory):
start = time.time()
findings = do_research(state['query'])
log_action(memory, AgentAction(
agent_id='researcher',
action_type='research',
input_summary=state['query'][:100],
output_summary=findings[:100],
timestamp=start,
success=True
))
return findings选择内存架构
合适的共享内存架构取决于您的部署模式。对于单进程 LangGraph 工作流,图状态已经足够。对于多进程或分布式代理,请使用 Redis。对于需要在重启后保持数据的长期运行项目,请使用带有适当索引的关系型数据库。对于服务之间的事件驱动协调,请在您选择的存储之上添加发布/订阅机制。
快速检查
请测试您对本课所讲共享内存和代理间通信的理解。
课程回顾
在本课中,您学到了:黑板模型使用共享数据存储,任何代理都可以读取和写入,从而在不让代理彼此直接耦合的情况下实现隐式协调;在分布式多代理系统中,Redis是持久化共享内存的首选;发布/订阅支持事件驱动通信,使代理能够在依赖项完成后立即做出响应。接下来,我们将学习代理任务的代码执行循环。
常见问题解答
「共享记忆与智能体间通信」课时是免费的吗?
是的 — 「共享记忆与智能体间通信」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 AI Engineering Academy 课程的其余内容,请升级到 CoddyKit PRO。 AI Engineering Academy 课程共包含 4 节课。
「共享记忆与智能体间通信」这节课中我会学到什么?
使用键值存储实现共享记忆层,让智能体从中读取和写入数据,从而支持异步协作,同时避免智能体之间的紧密耦合。 你通过在浏览器中直接运行的动手代码来练习 AI Engineering Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 AI Engineering Academy 需要有经验吗?
无需任何先前经验。CoddyKit 上的 AI Engineering Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 4 节课,共 4 节。
「共享记忆与智能体间通信」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 AI Engineering Academy 课中编写并运行代码吗?
能。每节 AI Engineering Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。
此课程中的所有课时
- 单智能体为何会遇到瓶颈
- 编排器-子智能体模式
- 使用 LangGraph 构建多智能体流程
- 共享记忆与智能体间通信