多 Agent 通信协议选型:gRPC、消息队列还是共享存储(续篇)

📅 2026/7/30 3:43:21
多 Agent 通信协议选型:gRPC、消息队列还是共享存储(续篇)
多 Agent 通信协议选型gRPC、消息队列还是共享存储续篇三个 Agent 用 HTTP 互相调用链路延迟叠加到 2 秒——选错协议不是慢一点是慢一倍。一、场景痛点你的多 Agent 系统有三个 Agent编排器、知识检索器、回答生成器。编排器用 HTTP REST 调用知识检索器200ms知识检索器再用 HTTP 调用回答生成器300ms总链路延迟 500ms HTTP 开销 800ms。你改用 gRPCHTTP 的 JSON 序列化/反序列化耗时约 50ms/次gRPC 用 Protobuf 二进制序列化耗时约 5ms/次。链路延迟降到 510ms。但你发现 gRPC 的长连接维护成本高——Agent 实例崩溃后 gRPC 连接不会自动断开编排器不知道下游 Agent 已经挂了请求超时才报错。你又考虑用消息队列Agent 之间不直接调用通过消息队列异步通信。好处是解耦Agent 崩溃不影响消息投递坏处是延迟更高消息队列的中转延迟约 10-50ms。核心矛盾gRPC 快但耦合紧消息队列慢但解耦好共享存储最慢但最简单——没有全能方案按场景选择。二、底层机制与原理剖析2.1 三种通信协议的架构对比2.2 选型决策矩阵维度gRPC消息队列共享存储延迟5-50ms50-200ms10-50ms可靠性连接断开需重试消息持久化不丢Redis 持久化可选解耦度低A 知道 B 地址高A 不知道 B最高A 不知道 B复杂度中proto stub高队列运维低读写 Redis并发控制流控在 gRPC 层队列自然缓冲需要手动并发控制可观测性gRPC 内置 tracing需额外接入Redis 没有原生 tracing三、生产级代码实现3.1 gRPC Agent 通信实现// agent_comm.proto —— gRPC Agent 通信协议定义 syntax proto3; package agent_comm; service AgentCommunication { // 同步调用编排器 → 知识检索器 rpc RetrieveKnowledge(RetrieveRequest) returns (RetrieveResponse); // 流式调用编排器 → 回答生成器生成过程中逐步返回 rpc GenerateAnswer(GenerateRequest) returns (stream GenerateChunk); } message RetrieveRequest { string query 1; string agent_id 2; int32 top_k 3; mapstring, string context 4; // 上下文变量 } message RetrieveResponse { repeated KnowledgeHit hits 1; string agent_id 2; int32 latency_ms 3; bool success 4; string error 5; } message KnowledgeHit { string content 1; float score 2; string source 3; } message GenerateRequest { string query 1; repeated KnowledgeHit context 2; string agent_id 3; } message GenerateChunk { string chunk_text 1; bool is_final 2; int32 tokens_so_far 3; }3.2 gRPC 客户端带健康检查// grpc-agent-client.ts —— gRPC Agent 客户端带健康检查和重试 import * as grpc from grpc/grpc-js; export class GrpcAgentClient { private client: any; private target: string; private maxRetries: number; private retryDelayMs: number; private healthCheckIntervalMs: number; private isHealthy: boolean true; constructor( target: string, protoPath: string, maxRetries: number 3, retryDelayMs: number 500, healthCheckIntervalMs: number 10000, ) { this.target target; this.maxRetries maxRetries; this.retryDelayMs retryDelayMs; this.healthCheckIntervalMs healthCheckIntervalMs; // 创建 gRPC 客户端长连接 健康检查 const credentials grpc.credentials.createInsecure(); this.client this.createClient(protoPath, target, credentials); // 启动健康检查循环定期检查连接是否正常 this.startHealthCheck(); } /** 带重试的 gRPC 调用 */ async retrieveKnowledge(request: any): Promiseany { for (let attempt 0; attempt this.maxRetries; attempt) { if (!this.isHealthy) { // Agent 不可用等待健康检查恢复或直接失败 if (attempt this.maxRetries) { throw new Error(Agent ${this.target} is unhealthy after ${this.maxRetries} retries); } await this.sleep(this.retryDelayMs); continue; } try { const response await new Promise((resolve, reject) { // gRPC 调用设置超时 5 秒 this.client.RetrieveKnowledge( request, { deadline: this.getDeadline(5000) }, (err: any, resp: any) err ? reject(err) : resolve(resp), ); }); return response; } catch (err: any) { if (attempt this.maxRetries) { await this.sleep(this.retryDelayMs * (attempt 1)); } else { throw new Error(gRPC call failed after ${this.maxRetries} retries: ${err.message}); } } } } /** 健康检查定期验证 Agent 是否可用 */ private startHealthCheck(): void { const interval setInterval(async () { try { // gRPC 健康检查协议HealthCheck service await new Promise((resolve, reject) { this.client.Check( { service: }, { deadline: this.getDeadline(2000) }, (err: any, resp: any) err ? reject(err) : resolve(resp), ); }); this.isHealthy true; } catch { this.isHealthy false; } }, this.healthCheckIntervalMs); } private getDeadline(timeoutMs: number): Date { return new Date(Date.now() timeoutMs); } private sleep(ms: number): Promisevoid { return new Promise((r) setTimeout(r, ms)); } private createClient(protoPath: string, target: string, credentials: any): any { // 实际实现需要加载 proto 定义 return {} as any; // 简化 } }3.3 消息队列 Agent 通信# mq_agent_comm.py —— 基于 Redis Streams 的 Agent 消息队列通信 import json import time import logging import redis logger logging.getLogger(mq-agent-comm) class MessageQueueAgentComm: Agent 间通过 Redis Streams 消息队列通信 def __init__(self, redis_url: str, agent_id: str): self.redis redis.from_url(redis_url) self.agent_id agent_id def send_message( self, target_agent: str, message_type: str, payload: dict, timeout_ms: int 30000, ) - dict: 发送消息到目标 Agent通过 Redis Streams # 每个 Agent 有自己的 streamagent:{agent_id}:input stream_key fagent:{target_agent}:input message { source_agent: self.agent_id, message_type: message_type, payload: json.dumps(payload), timestamp: str(time.time()), reply_to: fagent:{self.agent_id}:output, # 回复流 } # 写入 streamRedis Streams 自动分配消息 ID msg_id self.redis.xadd(stream_key, message) # 等待回复从自己的 output stream 读取 reply_stream fagent:{self.agent_id}:output start_time time.time() while time.time() - start_time timeout_ms / 1000: # 从 output stream 读取消息XREAD 非阻塞 replies self.redis.xread( {reply_stream: $}, # 从最新消息开始读 count1, block1000, # 阻塞 1 秒等待回复 ) if replies: for stream, messages in replies: for msg_id, data in messages: # 检查是否是回复我们的请求 if data.get(correlation_id) msg_id: return json.loads(data.get(payload, {})) # 超时目标 Agent 没有回复 raise TimeoutError(fNo reply from {target_agent} in {timeout_ms}ms) def receive_messages(self, handler: callable): 接收消息从自己的 input stream 消费 stream_key fagent:{self.agent_id}:input last_id 0 # 从头开始消费 while True: # XREADGROUP消费者组模式保证每条消息只被处理一次 messages self.redis.xreadgroup( agent-group, self.agent_id, {stream_key: last_id}, count10, block2000, ) if messages: for stream, msgs in messages: for msg_id, data in msgs: try: # 处理消息 payload json.loads(data.get(payload, {})) result handler(data.get(message_type), payload) # 发送回复到 reply_to stream reply_stream data.get(reply_to) if reply_stream: self.redis.xadd(reply_stream, { correlation_id: msg_id, payload: json.dumps(result), source_agent: self.agent_id, timestamp: str(time.time()), }) # 确认消息已处理XACK self.redis.xack(stream_key, agent-group, msg_id) except Exception as e: logger.error(fFailed to process message {msg_id}: {e})四、边界分析与架构权衡4.1 gRPC 的 Proto 文件维护成本每次修改 Agent 间通信的数据结构都需要更新 proto 文件、重新生成 stub、重新部署所有 Agent。这在快速迭代的 Agent 系统中是负担。消息队列方案不需要 proto 文件——消息格式是 JSON任意结构都可以传递不需要预定义 schema。4.2 消息队列的延迟叠加Redis Streams 的消息传递延迟约 10-30ms对于大多数场景可接受。但 Kafka 的延迟约 50-200ms磁盘持久化开销不适合实时 Agent 交互。4.3 适用边界与禁用场景协议适用场景禁用场景gRPC链路 ≤3 步、实时性要求高、Agent 数量少Agent 数量多连接管理复杂、Agent 可能频繁崩溃消息队列链路 ≥5 步、需要解耦、Agent 可能崩溃实时交互延迟 50ms 不可接受、简单场景队列运维成本 收益共享存储≤2 个 Agent、简单数据传递需要流式传输、高并发写入冲突五、结语多 Agent 通信协议选型不是哪个最好而是哪个最适合当前场景。gRPC 适合短链路实时交互延迟 5-50ms消息队列适合长链路解耦延迟 50-200ms共享存储适合最简场景1-2 个 Agent。gRPC 的代价是 proto 文件维护和连接管理消息队列的代价是延迟和运维复杂度共享存储的代价是并发控制和缺乏 tracing。链路 ≤3 步选 gRPC链路 ≥5 步或 Agent 可能崩溃选消息队列≤2 个 Agent 简单场景选共享存储。资料说明本文中的协议、版本、性能、成本和行业趋势应以可核验的一手资料为准。未标注统计口径的比例、时间表和预测仅作工程讨论不应视为行业事实。可参考 0730 资料来源索引并在发布前将具体来源贴到对应断言之后。