0Pricing
AI Engineering Academy · 课时

共享记忆与智能体间通信

使用键值存储实现共享记忆层,让智能体从中读取和写入数据,从而支持异步协作,同时避免智能体之间的紧密耦合。

共享记忆与智能体间通信 是 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 immediately

LangGraph 中的共享内存

在 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 反馈 — 无需本地设置。

此课程中的所有课时

  1. 单智能体为何会遇到瓶颈
  2. 编排器-子智能体模式
  3. 使用 LangGraph 构建多智能体流程
  4. 共享记忆与智能体间通信
← 返回 AI Engineering Academy