目录

“”” 面向多智能体协作的低开销通信与状态共享系统

运行方式: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()
关于

本项目面向多智能体协作场景,研究一种低开销通信、状态传递与共享记忆机制,提高智能体之间的信息交互效率。通过优化通信策略和状态同步方法,实现多个智能体在复杂环境下的协同决策,为智能制造、智能机器人及人工智能应用提供支持。

31.0 KB
邀请码
    Gitlink(确实开源)
  • 加入我们
  • 官网邮箱:gitlink@ccf.org.cn
  • QQ群
  • QQ群
  • 公众号
  • 公众号

版权所有:中国计算机学会技术支持:开源发展技术委员会
京ICP备13000930号-9 京公网安备 11010802047560号