作 者:吴佳浩(Alben)
公众号:全栈架构师笔记
系列专栏:《企业级 Agent 实战指南————Multi-Agent 架构设计:从单体 ReAct 到群智协同》· 第 03 篇
导读
点对点通信是网状的,Agent 一多就会演变成 的通信风暴;共享黑板(Blackboard)是星型的,把复杂度重新降回 。
很多团队以为共享内存就是挂一个 Redis 让所有 Agent 随便读写;但在高并发分布式协作中,脏读、脏写与脑裂会让整个系统产生不可逆的幻觉扩散。
没有版本控制的共享黑板是并发灾难,基于 CAS 乐观锁与分区隔离的黑板架构才是 Multi-Agent 协同的基石。
在多智能体系统从玩具走向工业级落地的过程中,当协同的 Agent 数量从 2 个增加到 5 个、10 个时,点对点消息传递(P2P Messaging)会立刻引发三大系统灾难:
| 绝境现象 | 具体翻车表现 | 架构根因 |
|---|---|---|
| 1. 通信网络风暴 | 5 个 Agent 互相同步状态,产生 | 通信拓扑呈网状 扩散, |
| (Message Storm) | 海量重复冗余的上下文转发 | 缺乏集中状态收敛中枢 |
| 2. 状态脏写覆写 | Coder A 和 Coder B 同时修改全局 | 缺乏并发控制与版本比对机制, |
| (Dirty Overwrites) | 架构配置,后提交的无情覆盖前者 | 发生典型的分布式写冲突 |
| 3. 认知撕裂与脑裂 | 各 Agent 手头的状态版本不一致, | 缺乏唯一的“单一信任事实源” |
| (State Split-Brain) | 架构师看到的版本落后于一线开发 | (Single Source of Truth) |
要解决多 Agent 并发协作中的混乱问题,经典人工智能与分布式系统中的最高范式就是:分级共享黑板模式(Partitioned Blackboard Pattern)。
一、共享黑板模式的核心架构:单一事实源与订阅分发
黑板模式的核心思想非常纯粹:所有 Agent 不直接互相通信,而是共同围绕一块“带版本控制与事件通知的共享黑板”进行读写。
| (全局状态与任务契约中枢) |
|---|
| Architect Agent |
| (写入系统规范) |
- 🔸 分区隔离原则:黑板分为
Global(全局规范)、Working(协同工作区)与Private(私有暂存区)。私有探索不污染全局; - 🔸 单一信任源(Single Source of Truth):所有 Agent 的推理基准完全对齐于黑板的最新版本;
- 🔸 事件驱动感知:当黑板特定分区发生状态变更时,通过 Redis Pub/Sub 或 WebSocket 广播通知下游关注该事件的 Agent。
一句话总结这一章的核心观点:
黑板模式把复杂网状沟通降维为星型结构,确保所有 Agent 拥有完全一致的全局心智。
二、并发控制深水区:基于 CAS 的乐观锁机制
当多个子代理(如 3 个并发 Coder Agent)同时尝试修改黑板上的同一份配置文件时,系统必须引入乐观并发控制(OCC / Compare-And-Swap):
乐观并发控制三原则
- 🔸 提交必须携带预期版本号:所有写操作必须声明
expected_version; - 🔸 冲突立即拒绝,杜绝脏写:一旦版本不匹配,原子化拒绝写入并返回最新快照;
- 🔸 语义合并重试:被拒绝的 Agent 必须重新拉取最新版本,结合自身修改意图进行语义合并后再次提交。
一句话总结这一章的核心观点:
乐观并发控制保证了数据一致性,语义合并保证了业务连续性。
三、生产级代码实战:带 CAS 乐观锁与事件广播的共享黑板
以下为基于 Python 3.11+ 构建的企业级共享黑板核心实现,完整包含分区存储、CAS 乐观并发控制与状态变更广播:
<span>"""
shared_blackboard_engine.py
生产级带 CAS 并发控制的共享黑板引擎
包含:
- 分区隔离
- CAS 乐观锁
- 并发冲突检测
- 黑板事件通知
"""</span>
<span>import</span> asyncio
<span>from</span> datetime <span>import</span> datetime
<span>from</span> typing <span>import</span> <span>Any</span>, <span>Callable</span>, <span>Dict</span>, <span>List</span>, <span>Optional</span>
<span>from</span> pydantic <span>import</span> BaseModel, Field
<span># ============================================================</span>
<span># Blackboard Entry</span>
<span># ============================================================</span>
<span>class</span> <span>BlackboardEntry</span>(<span>BaseModel</span>):
<span>"""黑板中的数据项"""</span>
key: <span>str</span>
value: <span>Any</span>
<span># CAS 版本号</span>
version: <span>int</span> = <span>1</span>
<span># 最近修改者</span>
last_modified_by: <span>str</span>
<span># 更新时间</span>
updated_at: datetime = Field(
default_factory=datetime.utcnow
)
<span># ============================================================</span>
<span># CAS 返回结果</span>
<span># ============================================================</span>
<span>class</span> <span>CASMutationResult</span>(<span>BaseModel</span>):
<span>"""CAS 写入结果"""</span>
success: <span>bool</span>
current_version: <span>int</span>
message: <span>str</span>
entry: <span>Optional</span>[BlackboardEntry] = <span>None</span>
<span># ============================================================</span>
<span># Shared Blackboard</span>
<span># ============================================================</span>
<span>class</span> <span>SharedBlackboardHub</span>:
<span>"""
企业级共享黑板
包含三个逻辑分区:
1. Global Partition
全局知识区
2. Working Partition
Agent 协作区(CAS)
3. Scratchpad
Agent 私有临时空间
"""</span>
<span>def</span> <span>__init__</span>(<span>self</span>):
<span># 全局知识</span>
self._global_partition: <span>Dict</span>[
<span>str</span>,
BlackboardEntry,
] = {}
<span># 协作区(CAS)</span>
self._working_partition: <span>Dict</span>[
<span>str</span>,
BlackboardEntry,
] = {}
<span># Agent 私有工作区</span>
self._agent_scratchpads: <span>Dict</span>[
<span>str</span>,
<span>Dict</span>[<span>str</span>, <span>Any</span>],
] = {}
<span># 黑板事件监听器</span>
self._subscribers: <span>List</span>[
<span>Callable</span>[[<span>str</span>, BlackboardEntry], <span>None</span>]
] = []
<span># CAS 原子锁</span>
self._lock = asyncio.Lock()
<span># --------------------------------------------------------</span>
<span># Event Subscription</span>
<span># --------------------------------------------------------</span>
<span>def</span> <span>subscribe_events</span>(<span>
self,
callback: <span>Callable</span>[[<span>str</span>, BlackboardEntry], <span>None</span>],
</span>) -> <span>None</span>:
<span>"""注册黑板事件监听器"""</span>
self._subscribers.append(callback)
<span># --------------------------------------------------------</span>
<span># Read</span>
<span># --------------------------------------------------------</span>
<span>def</span> <span>read_working_entry</span>(<span>
self,
key: <span>str</span>,
</span>) -> <span>Optional</span>[BlackboardEntry]:
<span>"""读取协作区数据"""</span>
<span>return</span> self._working_partition.get(key)
<span># --------------------------------------------------------</span>
<span># CAS Write</span>
<span># --------------------------------------------------------</span>
<span>async</span> <span>def</span> <span>cas_write_working_entry</span>(<span>
self,
key: <span>str</span>,
new_value: <span>Any</span>,
expected_version: <span>int</span>,
agent_role: <span>str</span>,
</span>) -> CASMutationResult:
<span>"""
CAS 原子写入
流程:
Read
│
▼
Version Check
│
┌─────┴─────┐
│ │
Pass Conflict
│ │
▼ ▼
Update Reject
"""</span>
<span>async</span> <span>with</span> self._lock:
existing = self._working_partition.get(key)
<span># ------------------------------------------------</span>
<span># 新建数据</span>
<span># ------------------------------------------------</span>
<span>if</span> existing <span>is</span> <span>None</span>:
<span>if</span> expected_version != <span>0</span>:
<span>return</span> CASMutationResult(
success=<span>False</span>,
current_version=<span>0</span>,
message=(
<span>f"Key '<span>{key}</span>' "</span>
<span>"does not exist. "</span>
<span>"expected_version must be 0."</span>
),
)
new_entry = BlackboardEntry(
key=key,
value=new_value,
version=<span>1</span>,
last_modified_by=agent_role,
)
self._working_partition[key] = new_entry
self._notify_subscribers(
key,
new_entry,
)
<span>return</span> CASMutationResult(
success=<span>True</span>,
current_version=<span>1</span>,
message=<span>"Key created successfully."</span>,
entry=new_entry,
)
<span># ------------------------------------------------</span>
<span># CAS Version Check</span>
<span># ------------------------------------------------</span>
<span>if</span> existing.version != expected_version:
<span>return</span> CASMutationResult(
success=<span>False</span>,
current_version=existing.version,
message=(
<span>"CAS Conflict: "</span>
<span>f"Current=<span>{existing.version}</span>, "</span>
<span>f"Expected=<span>{expected_version}</span>"</span>
),
entry=existing,
)
<span># ------------------------------------------------</span>
<span># 更新数据</span>
<span># ------------------------------------------------</span>
existing.value = new_value
existing.version += <span>1</span>
existing.last_modified_by = agent_role
existing.updated_at = datetime.utcnow()
self._notify_subscribers(
key,
existing,
)
<span>return</span> CASMutationResult(
success=<span>True</span>,
current_version=existing.version,
message=<span>"Version updated successfully."</span>,
entry=existing,
)
<span># --------------------------------------------------------</span>
<span># Event Notify</span>
<span># --------------------------------------------------------</span>
<span>def</span> <span>_notify_subscribers</span>(<span>
self,
key: <span>str</span>,
entry: BlackboardEntry,
</span>) -> <span>None</span>:
<span>"""广播黑板事件"""</span>
<span>for</span> subscriber <span>in</span> self._subscribers:
<span>try</span>:
subscriber(key, entry)
<span>except</span> Exception <span>as</span> e:
<span>print</span>(
<span>"Blackboard subscriber error:"</span>,
e,
)
<span># --------------------------------------------------------</span>
<span># Scratchpad</span>
<span># --------------------------------------------------------</span>
<span>def</span> <span>write_scratchpad</span>(<span>
self,
agent_role: <span>str</span>,
key: <span>str</span>,
value: <span>Any</span>,
</span>) -> <span>None</span>:
<span>"""
写入 Agent 私有 Scratchpad。
Scratchpad 为 Agent 独占,
无需加锁。
"""</span>
<span>if</span> agent_role <span>not</span> <span>in</span> self._agent_scratchpads:
self._agent_scratchpads[agent_role] = {}
self._agent_scratchpads[agent_role][key] = value
<span>def</span> <span>read_scratchpad</span>(<span>
self,
agent_role: <span>str</span>,
key: <span>str</span>,
</span>) -> <span>Optional</span>[<span>Any</span>]:
<span>"""读取 Agent 私有 Scratchpad"""</span>
<span>return</span> (
self._agent_scratchpads
.get(agent_role, {})
.get(key)
)
本篇总结
- 🔸 点对点通信是网状灾难,共享黑板是星型解耦;
- 🔸 分区设计是核心:全局规范只读、协作工作区 CAS 保护、私有暂存区自由探索;
- 🔸 CAS 乐观并发控制是保证多 Agent 并发写入不脏写的唯一数学保障;
- 🔸 事件驱动广播让下游 Agent 在状态更新时实现毫秒级自动感知与响应。
在多 Agent 规模化运行中,当 Agent 之间产生死锁、互相反驳或消耗海量 Token 时,如何进行治理?
筒子们本篇为《企业级 Agent 实战指南》· 第三章的第 3 篇,后续续会更新完整的agent的开发的全部过程,如果你对Agent开发感兴趣不妨关注一下本合集。
在下一篇中,我们将深入解构:《多智能体系统的通信风暴与死锁治理:生产级降级与容灾方案》!
把 Redis 当共享内存随便读写,是并发灾难的开始。本文给出分区隔离+CAS+事件广播的落地范式,适合正从两三个 Agent 走向十个以上、被脏写和脑裂困扰的团队借鉴。