5.1 多智能体通信机制 — Agent智能体开发实战 本节导读:深入理解多智能体系统的通信机制,掌握消息传递、共享内存、事件驱动等核心通信模式,构建高效的多智能体协作通信架构。 学习目标 理解多智能体通信的核心概念和模式 掌握消息传递协议和实现技术 学习共享内存和事件驱动通信机制 实现智能体的异步通信和同步机制 核心概念 多智能体通信机制是多智能体系统的核心,决定了智能体之间如何交换信息、协调行为、实现协作。 通信模式分类 分步实战 步骤 1:消息传递系统设计与实现 步骤 2:共享内存通信机制 常见问题 FAQ Q1:如何处理多智能体通信的延迟和超时?
本节导读:深入理解多智能体系统的通信机制,掌握消息传递、共享内存、事件驱动等核心通信模式,构建高效的多智能体协作通信架构。
多智能体通信机制是多智能体系统的核心,决定了智能体之间如何交换信息、协调行为、实现协作。
B --> F[阻塞式消息传递] B --> G[共享内存] C --> H[非阻塞消息队列] C --> I[发布订阅] D --> J[事件总线] D --> K[消息代理]
</div> ## 分步实战 ### 步骤 1:消息传递系统设计与实现 ```python import time import threading import queue import json import uuid import logging from dataclasses import dataclass, asdict from typing import Dict, List, Any, Optional, Callable from enum import Enum from abc import ABC, abstractmethod class MessageType(Enum): """消息类型枚举""" REQUEST = "request" RESPONSE = "response" NOTIFICATION = "notification" COMMAND = "command" HEARTBEAT = "heartbeat" ERROR = "error" STATUS = "status" class MessagePriority(Enum): """消息优先级""" LOW = 1 NORMAL = 2 HIGH = 3 CRITICAL = 4 @dataclass class Message: """消息数据结构""" id: str type: MessageType priority: MessagePriority sender: str receiver: str content: Dict[str, Any] timestamp: float ttl: float = 30.0 metadata: Dict[str, Any] = None def is_expired(self) -> bool: """检查消息是否过期""" return time.time() - self.timestamp > self.ttl def to_dict(self) -> Dict[str, Any]: """转换为字典格式""" data = asdict(self) data['type'] = self.type.value data['priority'] = self.priority.value return data class MessageBus: """消息总线""" def __init__(self, transport_layer): self.transport_layer = transport_layer self.message_queue = queue.PriorityQueue() self.router = MessageRouter() self.response_handlers = {} self.lock = threading.RLock() self.running = False self.logger = logging.getLogger('MessageBus') def start(self): """启动消息总线""" self.transport_layer.connect() self.running = True self.message_processor_thread = threading.Thread(target=self._process_messages, daemon=True) self.message_processor_thread.start() self.logger.info("Message bus started") def stop(self): """停止消息总线""" self.running = False self.transport_layer.disconnect() self.message_processor_thread.join(timeout=5) self.logger.info("Message bus stopped") def send_message(self, message: Message) -> bool: """发送消息""" return self.transport_layer.send(message) def request_response(self, message: Message, timeout: float = 10.0) -> Optional[Message]: """请求响应""" message_id = message.id response_queue = queue.Queue() with self.lock: self.response_handlers[message_id] = response_queue if self.send_message(message): try: response = response_queue.get(timeout=timeout) return response except queue.Empty: self.logger.error(f"Timeout waiting for response to message {message_id}") return None return None def _process_messages(self): """处理消息的内部方法""" while self.running: try: message = self.transport_layer.receive(timeout=1.0) if message: if message.is_expired(): self.logger.warning(f"Message {message.id} expired") continue self.router.route_message(message) except Exception as e: self.logger.error(f"Error processing messages: {e}") class AgentCommunicator: """智能体通信器""" def __init__(self, agent_id: str, message_bus: MessageBus): self.agent_id = agent_id self.message_bus = message_bus self.logger = logging.getLogger(f'AgentCommunicator-{agent_id}') def send_message(self, receiver_id: str, message_type: MessageType, content: Dict[str, Any], priority: MessagePriority = MessagePriority.NORMAL) -> bool: """发送消息""" message = Message( id=str(uuid.uuid4()), type=message_type, priority=priority, sender=self.agent_id, receiver=receiver_id, content=content, timestamp=time.time(), ttl=30.0 ) return self.message_bus.send_message(message) def request_response(self, receiver_id: str, message_type: MessageType, content: Dict[str, Any], timeout: float = 10.0) -> Optional[Message]: """请求响应""" message = Message( id=str(uuid.uuid4()), type=message_type, priority=MessagePriority.HIGH, sender=self.agent_id, receiver=receiver_id, content=content, timestamp=time.time(), ttl=30.0 ) return self.message_bus.request_response(message, timeout) def broadcast_notification(self, message_type: MessageType, content: Dict[str, Any], ttl: float = 10.0): """广播通知""" message = Message( id=str(uuid.uuid4()), type=message_type, priority=MessagePriority.LOW, sender=self.agent_id, receiver="*", content=content, timestamp=time.time(), ttl=ttl ) return self.message_bus.send_message(message) def start(self): """启动通信器""" self.message_bus.start() def stop(self): """停止通信器""" self.message_bus.stop()
import mmap import os import struct from typing import Dict, Any, Optional @dataclass class SharedMemoryMessage: """共享内存消息结构""" message_id: int sender_id: str receiver_id: str message_type: int priority: int content_size: int timestamp: float ttl: float def pack(self) -> bytes: """打包为字节""" return struct.pack( 'Q 20s 20s I I I d d', self.message_id, self.sender_id.encode('utf-8').ljust(20), self.receiver_id.encode('utf-8').ljust(20), self.message_type, self.priority, self.content_size, self.timestamp, self.ttl ) @classmethod def unpack(cls, data: bytes) -> 'SharedMemoryMessage': """从字节解包""" (message_id, sender_id, receiver_id, message_type, priority, content_size, timestamp, ttl) = struct.unpack( 'Q 20s 20s I I I d d', data ) return cls( message_id=message_id, sender_id=sender_id.decode('utf-8').strip(), receiver_id=receiver_id.decode('utf-8').strip(), message_type=message_type, priority=priority, content_size=content_size, timestamp=timestamp, ttl=ttl ) class SharedMemoryCommunicator: """共享内存通信器""" def __init__(self, memory_name: str = "agent_comm", memory_size: int = 1024 * 1024): self.memory_name = memory_name self.memory_size = memory_size self.message_size = struct.calcsize('Q 20s 20s I I I d d') self.max_message_size = memory_size - self.message_size self.lock = threading.RLock() self.message_counter = 0 # 创建共享内存文件 self._create_shared_memory() def _create_shared_memory(self): """创建共享内存""" self.memory_file = mmap.mmap(-1, self.memory_size, self.memory_name) self.memory_file.seek(0) self.memory_file.write(b'\x00' * self.memory_size) def send_message(self, sender_id: str, receiver_id: str, message_type: MessageType, content: Dict[str, Any], priority: MessagePriority = MessagePriority.NORMAL) -> bool: """发送消息""" with self.lock: content_str = json.dumps(content) if len(content_str) > self.max_message_size: return False message = SharedMemoryMessage( message_id=self.message_counter, sender_id=sender_id, receiver_id=receiver_id, message_type=message_type.value, priority=priority.value, content_size=len(content_str), timestamp=time.time(), ttl=30.0 ) message_data = message.pack() content_data = content_str.encode('utf-8') # 简化实现:写入固定位置 offset = 4 + self.message_counter * (self.message_size + len(content_data)) if offset + len(message_data) + len(content_data) > self.memory_size: return False self.memory_file.seek(offset) self.memory_file.write(message_data) self.memory_file.write(content_data) self.message_counter += 1 return True
A:处理通信延迟和超时的策略:
A:确保通信可靠性的策略:
A:处理消息冲突的策略:
本节深入探讨了多智能体通信机制,从消息传递到共享内存,全面介绍了多智能体系统的通信技术和实现方法。通过学习本节内容,读者应该能够:
下一节我们将探讨多智能体协作协议与策略,继续深入多智能体系统的核心技术。
关键词:Agent智能体开发实战, 多智能体通信, 消息传递, 共享内存, 协作协议
难度:进阶
预计阅读:30 分钟