共享黑板模式(Blackboard)实战:多 Agent 如何并发协作而不冲突?

文章来源声明: 原文作者:吴佳浩Alben; 来源站点:掘金; 原文链接:https://juejin.cn/post/7688603830762782735; 本文基于上述来源整理/加工,觅优补充点评,仅供技术学习交流。版权归原作者所有。
觅优短评

把 Redis 当共享内存随便读写,是并发灾难的开始。本文给出分区隔离+CAS+事件广播的落地范式,适合正从两三个 Agent 走向十个以上、被脏写和脑裂困扰的团队借鉴。

共享黑板模式(Blackboard)实战:多 Agent 如何并发协作而不冲突? ----------------------------------------

请添加图片描述

作 者:吴佳浩(Alben)
公众号:全栈架构师笔记
系列专栏:《企业级 Agent 实战指南————Multi-Agent 架构设计:从单体 ReAct 到群智协同》· 第 03 篇


导读
点对点通信是网状的,Agent 一多就会演变成 O(N2)O(N^2) 的通信风暴;共享黑板(Blackboard)是星型的,把复杂度重新降回 O(N)O(N)
很多团队以为共享内存就是挂一个 Redis 让所有 Agent 随便读写;但在高并发分布式协作中,脏读、脏写与脑裂会让整个系统产生不可逆的幻觉扩散。
没有版本控制的共享黑板是并发灾难,基于 CAS 乐观锁与分区隔离的黑板架构才是 Multi-Agent 协同的基石。


在多智能体系统从玩具走向工业级落地的过程中,当协同的 Agent 数量从 2 个增加到 5 个、10 个时,点对点消息传递(P2P Messaging)会立刻引发三大系统灾难:

绝境现象具体翻车表现架构根因
1. 通信网络风暴5 个 Agent 互相同步状态,产生通信拓扑呈网状 O(N2)O(N^2) 扩散,
(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
(写入系统规范)

mermaid diagram

  • 🔸 分区隔离原则:黑板分为 Global(全局规范)、Working(协同工作区)与 Private(私有暂存区)。私有探索不污染全局;
  • 🔸 单一信任源(Single Source of Truth):所有 Agent 的推理基准完全对齐于黑板的最新版本;
  • 🔸 事件驱动感知:当黑板特定分区发生状态变更时,通过 Redis Pub/Sub 或 WebSocket 广播通知下游关注该事件的 Agent。

一句话总结这一章的核心观点:
黑板模式把复杂网状沟通降维为星型结构,确保所有 Agent 拥有完全一致的全局心智。


二、并发控制深水区:基于 CAS 的乐观锁机制

当多个子代理(如 3 个并发 Coder Agent)同时尝试修改黑板上的同一份配置文件时,系统必须引入乐观并发控制(OCC / Compare-And-Swap)

mermaid diagram

乐观并发控制三原则

  • 🔸 提交必须携带预期版本号:所有写操作必须声明 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开发感兴趣不妨关注一下本合集。

在下一篇中,我们将深入解构:《多智能体系统的通信风暴与死锁治理:生产级降级与容灾方案》