第三部分:多智能体协作演示与性能统计 这一部分创建实际场景,对前两部分进行演示。 代码建立了三个智能体: planner:负责生成装配计划。 robot:负责执行部件装配。 inspector:负责质量检查。 运行过程如下: 规划智能体开始生成装配计划,进度为 20%。 机器人开始装配,进度为 45%。 检查智能体等待质量检查。 机器人只更新进度到 80%,任务名称不会被重复发送。 规划智能体完成任务。 检查智能体开始工作,进度达到 60%。 三个智能体分别处理收到的消息,形成一致的共享状态视图。 程序输出最终状态、消息数量和总传输字节数。 核心作用:证明多个智能体可以通过低开销消息同步状态,并直观展示通信成本。负责人:任德锦
“”” 面向多智能体协作的低开销通信与状态共享系统
运行方式:python multi_agent_communication.py 本示例仅依赖 Python 标准库,并明确分为三个部分。 “””
from future import annotations
import asyncio import json import time from dataclasses import asdict, dataclass, field from typing import Any
@dataclass(slots=True) class AgentState: “””智能体的可共享状态。”””
agent_id: str status: str = "idle" task: str = "" progress: int = 0 updated_at: float = field(default_factory=time.time) def update(self, **changes: Any) -> dict[str, Any]: """更新本地状态,并只返回发生变化的字段(增量数据)。""" delta: dict[str, Any] = {} for key, value in changes.items(): if not hasattr(self, key): raise ValueError(f"未知状态字段:{key}") if getattr(self, key) != value: setattr(self, key, value) delta[key] = value if delta: self.updated_at = time.time() delta["updated_at"] = self.updated_at return delta
@dataclass(slots=True) class Message: “””轻量级智能体消息。”””
sender: str topic: str sequence: int payload: dict[str, Any] def encode(self) -> bytes: # 紧凑 JSON 去除多余空格,降低通信字节数。 return json.dumps( asdict(self), ensure_ascii=False, separators=(",", ":") ).encode("utf-8") @classmethod def decode(cls, data: bytes) -> "Message": return cls(**json.loads(data.decode("utf-8")))
class CommunicationBus: “””进程内异步发布/订阅总线,可替换为 MQTT、Redis 或网络传输层。”””
def __init__(self) -> None: self._subscribers: dict[str, list[asyncio.Queue[bytes]]] = {} self.message_count = 0 self.transferred_bytes = 0 def subscribe(self, topic: str) -> asyncio.Queue[bytes]: queue: asyncio.Queue[bytes] = asyncio.Queue() self._subscribers.setdefault(topic, []).append(queue) return queue async def publish(self, message: Message) -> None: data = message.encode() queues = self._subscribers.get(message.topic, []) for queue in queues: await queue.put(data) self.message_count += 1 self.transferred_bytes += len(data) * len(queues)
class CollaborativeAgent: “””能够发布增量状态、接收状态并维护共享视图的智能体。”””
TOPIC = "agent/state" def __init__(self, agent_id: str, bus: CommunicationBus) -> None: self.state = AgentState(agent_id=agent_id) self.bus = bus self.inbox = bus.subscribe(self.TOPIC) self.shared_states: dict[str, dict[str, Any]] = {} self._sequence = 0 self._last_seen: dict[str, int] = {} async def change_state(self, **changes: Any) -> None: delta = self.state.update(**changes) if not delta: return # 状态未变化时不发送消息,避免无效通信。 self._sequence += 1 await self.bus.publish( Message( sender=self.state.agent_id, topic=self.TOPIC, sequence=self._sequence, payload=delta, ) ) async def receive_once(self) -> None: message = Message.decode(await self.inbox.get()) # 利用序列号忽略重复或过期消息。 if message.sequence <= self._last_seen.get(message.sender, 0): return self._last_seen[message.sender] = message.sequence state = self.shared_states.setdefault(message.sender, {}) state.update(message.payload) async def receive_available(self) -> None: while not self.inbox.empty(): await self.receive_once()
本项目面向多智能体协作场景,研究一种低开销通信、状态传递与共享记忆机制,提高智能体之间的信息交互效率。通过优化通信策略和状态同步方法,实现多个智能体在复杂环境下的协同决策,为智能制造、智能机器人及人工智能应用提供支持。
版权所有:中国计算机学会技术支持:开源发展技术委员会 京ICP备13000930号-9 京公网安备 11010802047560号
“”” 面向多智能体协作的低开销通信与状态共享系统
运行方式:python multi_agent_communication.py 本示例仅依赖 Python 标准库,并明确分为三个部分。 “””
from future import annotations
import asyncio import json import time from dataclasses import asdict, dataclass, field from typing import Any
============================================================
第一部分:状态与消息模型
============================================================
@dataclass(slots=True) class AgentState: “””智能体的可共享状态。”””
@dataclass(slots=True) class Message: “””轻量级智能体消息。”””
============================================================
第二部分:低开销通信与状态共享机制
============================================================
class CommunicationBus: “””进程内异步发布/订阅总线,可替换为 MQTT、Redis 或网络传输层。”””
class CollaborativeAgent: “””能够发布增量状态、接收状态并维护共享视图的智能体。”””