跨Agent通信,跨Session通信是一个非常实用的功能
在Agent开发从单机demo走向真正的复杂业务时,大多数人会遇到同一个瓶颈:多个Agent可以各自跑得很好,可一旦让它们协作,就不知道该把数据放在哪里。
如果你观察实际项目中的Agent系统,会发现一个典型矛盾——单个Agent的上下文做得越来越厚,Session记忆越来越长,但Agent与Agent之间、Session与Session之间却在各说各话。
比如你有一个“销售Agent”负责分析线索,有一个“调度Agent”负责分配任务,又有一个“质检Agent”负责检查结果。它们各自有独立的会话、独立的记忆,甚至部署在不同的服务中。表面上看这是一个“多Agent架构”,但实际运行起来,每个Agent仍然是孤岛。因为会话是私有的,Agent之间没有一条真正可用的“电话线”。
跨Agent通信、跨Session通信,这三个词的组合看起来像一个抽象理论,但它其实是实际项目里让人最头疼的工程问题。这篇博客我来拆解这个问题的本质,给出从本地服务到进程间通信的代码实现和踩坑记录。这篇文章默认你已经知道Agent是什么,Session大概是什么,但还没把通信链路想清楚——我把这部分补上。
1. 先看清楚:Agent和Session的隔离,是一个默认设计问题
先说一个容易混淆的点:Agent、Session、通信三件事,经常被人分开理解。
AI Agent本身是一个能感知环境、做决策、执行动作的自治实体。Session则是这个Agent与外界交互的过程性上下文,常见的是某个对话、某次任务执行区间。
但把Agent和Session放在一起看时,一个根本性矛盾就出现了:每个Agent大概率会维护自己的Session状态,希望保持上下文独立。可是业务协作又要求多个Agent共享一些信息,甚至某个Agent需要获知另一个Agent当前Session的实时状态。
举个例子。你在做一个智能客服系统,用户从网页端发起了一个“查快递”的Session,与此同时他可能又从微信端发起了“改地址”的Session。从用户视角看,这是同一个需求场景。但如果系统把两个Session隔离得干干净净,客服Agent就不能感知到快递查询Session中已经出现的异常状态。结果是它只能让用户复述一遍问题。
很多人把这个问题归咎于“Session管理没做好”,其实深层原因是“Session之间的通信链路根本没设计”。
如果只看表面,容易误以为跨Session通信只是把数据放到共享数据库里,让另一个Agent去读。但如果做得这么简单,会出现更多问题:命名冲突、并发覆盖、脏读、一个Agent改了数据库导致另一个Agent误判状态……通信不是数据共享的代名词。真正的跨Agent通信,是让一个Agent能向另一个Agent传递“有业务语义的事件”,并得到明确的响应或后续动作。
2. 通信要解决什么问题
站在AI应用的真实开发场景里,跨Agent通信与跨Session通信解决的是同一条链路上的两个不同困难:
- 跨Agent通信:强调实体之间协作。不同Agent拥有不同角色、不同职责、不同工具的权限边界。它们要互相传递任务、汇报结果、请求确认。
- 跨Session通信:强调会话上下文之间的事件同步。同一个Agent可能会在处理多个Session,也可能一个复杂的Session会分散到多个Agent中执行。
所以我更愿意把它们的组合称为**“通信层”**。这个通信层要解决四个核心问题:
- 路由问题:消息发给谁?是按Agent类型,还是按Session归属?
- 状态同步问题:接收方Agent处理事件时,需要知道是哪个Session触发的,以及这个Session当前处于什么阶段。
- 生命周期问题:如果接收方Agent还没有启动,或者对应Session已经过期,消息该怎么办?
- 安全边界问题:跨会话通信不能从一个Session越权到另一个Session,必须有身份和授权校验。
很多Agent开发学习路线从Prompt工程、工具调用开始,却很少提到“消息投递”这件事。这会导致你做一个多Agent的demo很容易,做生产化部署却发现通信逻辑乱成一团。
理解了这个背景,再来看代码实现就会觉得顺理成章。因为这不是在学一个新轮子,而是在补一个工程系统本来就该有的骨架。
3. 通信模式的选择:同步调用与异步消息
在设计跨Agent通信时,第一步不是写代码,而是决定通信模式。主流模式有两种:
| 模式 | 工作机制 | 适用场景 | 风险点 |
|---|---|---|---|
| 同步请求/响应 | Agent A直接调用Agent B接口,等待返回结果 | 任务分发、问答、单次请求 | 调用链过长会阻塞;B不可用时A容易雪崩 |
| 异步消息/事件 | Agent A把事件发到通道,Agent B自行监听处理 | 通知类场景、跨流程解耦 | 需要额外处理投递可靠性、顺序性 |
在实际项目中,我更加推荐以异步事件为底座,同时保留同步调用的封装。理由是跨Agent通信中的调用方往往不能确定被调用方什么时间能完成准备,尤其在两个Agent都有自己的LLM推理链时,同步等待会产生超长延迟。
举个例子,订单Agent处理完一个用户指令后,需要让消息Agent通知用户“处理完成”。如果采用同步调用,订单Agent必须等待消息Agent完成发送后才返回结果;一旦消息Agent在等待LLM生成通知文案,整个订单流程都会被卡住。
如果改成异步事件模式,订单Agent只需要发出一个AgentEvent,内容包含“订单ID、事件类型、目标Agent标识、所属SessionID”,就可以马上返回继续处理其他事情。消息Agent订阅该事件后再异步执行。
这种模式的核心概念叫作**“事件总线”**。事件总线负责接收消息、按照订阅关系路由消息、并在必要时保存消息记录。
4. 在本地多Agent架构里实现跨Session通信
很多早期项目并不需要分布式消息中间件。多个Agent跑在同一个进程内,只要设计一个轻量级事件总线,就能打通它们之间的Session。
4.1 最简模型:先设计AgentID与SessionID
在代码开始前,需要明确一点:Agent是实体,Session是上下文。消息里必须同时带上这两个标识。
# agent_event.py from dataclasses import dataclass, field from datetime import datetime from typing import Any @dataclass class AgentEvent: event_id: str source_agent: str # 发送方Agent标识 target_agent: str | None # 接收方Agent标识,None表示广播 session_id: str | None # 所属Session标识,None表示全局 event_type: str # 事件类型 payload: Any # 实际数据 timestamp: str = field(default_factory=lambda: datetime.now().isoformat())这里的关键点是session_id字段。它在跨Session通信中承担“定位上下文”的作用。如果不带它,接收方只能看到“有人发了一个事件”,但无法关联到当前的会话状态。
4.2 让事件总线支持跨Agent投递
接下来实现事件总线。它需要维护两个映射关系:
- 从AgentID到订阅回调函数。
- 从SessionID到需要感知该Session的Agent列表。
# event_bus.py import uuid from collections import defaultdict from agent_event import AgentEvent class EventBus: def __init__(self): # key: agent_name, value: list of callable self._subscribers = defaultdict(list) # key: session_id, value: set of agent_name self._session_agents = defaultdict(set) def register_agent(self, agent_name: str, session_id: str | None = None, handler=None): """注册Agent,可将其关联到某个Session""" if handler: self._subscribers[agent_name].append(handler) if session_id: self._session_agents[session_id].add(agent_name) def publish(self, event: AgentEvent): """发布事件到目标Agent或目标Session所关联的Agent""" targets = [] if event.target_agent: targets.append(event.target_agent) elif event.session_id and event.session_id in self._session_agents: targets.extend(self._session_agents[event.session_id]) else: # 广播到所有注册Agent targets.extend(self._subscribers.keys()) # 去重 for agent_name in set(targets): for handler in self._subscribers.get(agent_name, []): try: handler(event) except Exception as exc: print(f"[EventBus] handler for {agent_name} error: {exc}") bus = EventBus()这里隐藏了好几个容易出错的地方,需要单独说明。
第一,如果只想给某一个Session里的Agent发消息,就需要依赖_session_agents这个映射表。它维护的是“哪些Agent正在处理某个Session”。在实际项目里,这个映射表通常来自于Agent启动时的注册,而不是由EventBus自动管理。
第二,publish方法里的异常捕获非常重要。跨Agent通信中一个Agent的handler出错,不应该阻塞其他Agent接收同一条事件。否则一个下游Agent的异常会导致整个总线瘫痪。
第三,为什么要支持target_agent直接指定。因为很多场景下调用方明确知道消息要发给谁。例如某个工具Agent要向规划Agent汇报执行结果,就必须直接指定目标。
4.3 用例子验证跨Session语义
现在编写两个Agent跑一段时间。
# demo_cross_session.py from agent_event import AgentEvent from event_bus import bus def process_order_event(event: AgentEvent): print(f"[订单会话处理器] 收到事件: {event.session_id} - {event.event_type}") print(f"[订单会话处理器] 载荷: {event.payload}") def process_notify_event(event: AgentEvent): print(f"[通知Agent] 收到事件: {event.session_id} - {event.event_type}") print(f"[通知Agent] 内容: {event.payload}") # 假设当前进程中有两个Session: # session_001 是用户的订单查询会话, 由订单Agent处理 # session_002 是后台通知会话, 由通知Agent处理 # 注意: 它们可能属于不同的业务模块, 但需要通信 bus.register_agent("order_agent", session_id="session_001", handler=process_order_event) bus.register_agent("notify_agent", session_id="session_002", handler=process_notify_event) # 现在, 跨Agent同时跨Session发送一条事件 # 订单Agent处理完session_001中的用户请求后, # 希望让另一个Agent——通知Agent——在session_002相关上下文中发送通知 bus.publish(AgentEvent( event_id="evt_0001", source_agent="order_agent", target_agent="notify_agent", session_id="session_002", # 注意这里目标是另一个Session event_type="order.completed", payload={"order_id": "A1001", "message": "您的订单已完成"} ))运行上面的代码,输出是:
[通知Agent] 收到事件: session_002 - order.completed [通知Agent] 内容: {'order_id': 'A1001', 'message': '您的订单已完成'}可以看到,订单Agent的Handler并没有被调用,因为目标Agent是notify_agent,而不是order_agent。这就是跨Agent通信。而session_id从session_001切换到了session_002,这就是跨Session通信。
这里真正需要纠正的一个误区是:很多Agent框架看起来已经内置了“Session状态”,但那只代表当前Agent能读取到当前Session的上下文。它并不等于其他Agent能读取,更不等于其他Session能感知。Session是面向单个Agent私有上下文的边界;通信层则是穿过这个边界、同时避免私有数据直接暴露的桥梁。
5. 基于消息队列的生产级跨进程方案
当Agent被部署成多个独立服务,或者你需要保证消息不丢、可重放时,轻量级内存事件总线就不够用了。此时需要引入真正的消息中间件。
下面用一个通用示例演示生产级的跨Agent通信如何实现。我会以Redis Streams为例,因为它依赖少、理解成本低,适合说明思路。如果团队已经在使用RabbitMQ、Kafka,你只需把消费者和生产者换成对应组件即可,整体设计不变。
5.1 生产端:发布跨Agent事件到Stream
# producer_agent.py import json import redis import uuid from datetime import datetime r = redis.Redis(host="localhost", port=6379, db=0) def publish_cross_event(source_agent: str, target_agent: str, session_id: str, event_type: str, payload: dict): event = { "event_id": str(uuid.uuid4()), "source_agent": source_agent, "target_agent": target_agent, "session_id": session_id, "event_type": event_type, "payload": json.dumps(payload, ensure_ascii=False), "timestamp": datetime.now().isoformat() } # 使用 Redis Stream r.xadd("agent:event:stream", event)生产端的核心是使用Stream数据结构保存事件。使用Stream而不是简单的List,是因为Stream天然支持多个消费者组之间独立消费,方便多个Agent针对同一消息做不同处理。
很多开发者会在这一层犯一个错误:他们倾向于把payload做得特别复杂,塞入大量上下文。这个做法会让消息体变大,而且导致耦合。正确做法是payload只传递必要参数和引用ID,具体数据由接收方Agent自行查询。让事件的发送方与接收方只通过事件ID和SessionID耦合,而不是共享一份巨大的上下文数据。
5.2 消费端:接收事件并做路由
每个服务启动后台线程或使用异步任务框架消费Stream中的事件。
# consumer_agent.py import json import redis import threading import time r = redis.Redis(host="localhost", port=6379, db=0) STREAM = "agent:event:stream" GROUP = "notify_agent_group" CONSUMER = "notify_agent_instance_1" try: r.xgroup_create(STREAM, GROUP, id="0", mkstream=True) except Exception: pass # 组已存在 def handle_event(event): target = event.get("target_agent") # 该消费者只处理发给notify_agent的事件 if target != "notify_agent": return session_id = event.get("session_id") event_type = event.get("event_type") payload = json.loads(event.get("payload", "{}")) print(f"[notify_agent] 收到跨Session事件: session={session_id}, type={event_type}") print(f"[notify_agent] 携带数据: {payload}") # 这里可以调用具体的Agent执行逻辑 def listen_loop(): while True: # 从消费者组读取新消息 results = r.xreadgroup(GROUP, CONSUMER, {STREAM: ">"}, count=10, block=5000) if not results: continue for _, messages in results: for msg_id, msg in messages: handle_event(msg) # 确认消息已被处理 r.xack(STREAM, GROUP, msg_id) time.sleep(0.1) if __name__ == "__main__": t = threading.Thread(target=listen_loop, daemon=True) t.start() print("消费者监听已启动...") t.join()这里有几个生产级要点:
- 消费者组保证了一条消息不会在一个组内被多个实例重复处理,适用于“同一个Agent部署多个副本”的场景。
- xack确认机制记录消息是否被成功处理。一旦程序异常退出,Redis Streams不会把已读但未确认的消息标记为完成,重启后还能继续消费。
- 消费过滤逻辑不能少。多个Agent共享同一个Stream时,每个服务进程都要判断
target_agent,避免处理不属于自己的事件。
运行这段代码时,需要确保Redis 5.0以上版本且服务已启动。启动消费者进程,再运行生产者发送一条事件,可以看到输出:
[notify_agent] 收到跨Session事件: session=session_002, type=order.completed [notify_agent] 携带数据: {'order_id': 'A1001', 'message': '您的订单已完成'}相比第四节的本地EventBus,这个方案多了一个持久化层,服务重启后事件不会全部丢失,生产环境一般按这个思路落地。
5.3 更轻的分布式模式:AgentGateway
如果不想引入Redis Streams这样的中间件,还有一种变通方案,就是用一个HTTP服务作为AgentGateway,所有通信都通过Gateway转发。这本质上是一个中心化的路由代理。
# gateway.py from fastapi import FastAPI, Request app = FastAPI() # agent_name -> base_url AGENT_ROUTES = { "order_agent": "http://order-agent-service:8001", "notify_agent": "http://notify-agent-service:8002", } @app.post("/rpc/{target_agent}") async def rpc_to_agent(target_agent: str, request: Request): data = await request.json() base_url = AGENT_ROUTES.get(target_agent) if not base_url: return {"error": "unknown agent"} # 这里可以使用 httpx 异步转发请求 return {"target": target_agent, "data": data}Gateway模式的好处是路由逻辑统一,便于统一做鉴权、限流和日志追踪。坏处是它引入了中心节点,网关一旦故障整个通信就断了,需要额外考虑网关高可用。在Agent数量不多且语言环境复杂(不同Agent用不同语言编写)时,这个方案反而比重型消息队列更容易接入。
选择哪种方式,可以按这个标准判断:
- Agent在同一个服务内,不需要持久化:直接采用本地EventBus。
- Agent跨服务、需要消息可靠投递:消息队列加消费者组。
- Agent数量少,但使用语言复杂、需要快速打通:HTTP网关模式。
6. 事件与命令的区分
在编写跨Agent逻辑时,很多新手会写出一种混乱的消息:让一个Agent告诉另一个Agent“你应该做什么”,但这个消息内部包含大量判断规则,导致接收方只能死板地执行。
更好的做法是区分两类消息:
事件:表示“已经发生了什么”,是事实描述,例如OrderCreated、PaymentConfirmed。发布事件的一方不关心谁消费,也不期待响应。
命令:表示“请你去做某件事”,需要明确指定接收方,通常期待执行结果或确认。例如SendNotification(order_id="A1001")。
设计跨Agent通信时,尽量让业务逻辑表现为事件驱动。因为Agent本身是有决策能力的,给它发送事件后,它可以自己判断是否需要响应,以及怎么响应。而命令的发送隐含了“发送方已经替接收方做了决策”,这会削弱Agent自治性。
例如:订单Agent完成了订单后,更合理的事件是:
{ "event_type": "order.completed", "payload": { "order_id": "A1001", "completed_at": "2025-01-01T10:00:00Z" } }而不是命令:
{ "event_type": "send_notification", "payload": { "order_id": "A1001", "phone": "13800000000", "message": "您的订单已完成" } }前者允许通知Agent自行决定通过什么渠道通知、使用什么文案模板;后者则让通知Agent变成了一个死板的执行器。事件设计比命令设计更适合AI Agent的自主决策机制。
7. 跨Session通信的常见问题与排查思路
真正运行起来后,以下问题会让开发者排查很久。
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 消息发出去但目标Agent没有处理 | 目标Agent名称不匹配,或target_agent字段写错 | 打印收到的原始消息,检查target_agent值 | 建立Agent注册表,避免硬编码字符串名称 |
| Session上下文错乱,Agent处理了其他Session的事件 | 发送事件时session_id透传错误 | 在消息入口和出口分别打印session_id | 为事件增加source_session_id和target_session_id两个字段 |
| 重复处理同一条消息导致副作用 | 消费端未做消息幂等,或消息重放 | 检查消息ID是否唯一处理过一次 | 在数据库中记录event_id消费表,做幂等处理 |
| 下游Agent异常导致事件一直堆积 | 消费队列阻塞,没有快速失败 | 查看消费组堆积数量和日志 | 对异常情况做重试次数限制,超过阈值进入死信队列 |
| 跨Session访问了不该访问的私有状态 | 接收Agent没有对Session所有者校验 | 检查登录态和Session归属字段 | 在总线层强制校验会话所有权,禁止默认放行 |
其中“幂等”问题最容易遗漏。在进程内EventBus中,不存在消息重复投递的问题。但引入Redis Streams或Kafka后,断线重连、消费超时都会导致同一条消息被投递两次。如果接收方Agent不能识别重复,就可能执行两次转账通知、两次状态更新,或者生成两条重复记录。
幂等实现的典型方式:
# idempotent.py import redis r = redis.Redis(host="localhost", port=6379, db=1) def is_processed(event_id: str) -> bool: # 如果event_id已存在,说明处理过 return r.exists(f"processed_event:{event_id}") == 1 def mark_processed(event_id: str): r.set(f"processed_event:{event_id}", "1", ex=86400)在消费端,每次处理事件前先做幂等判断:已处理则跳过,未处理才执行业务逻辑,处理完成后立刻标记。过期时间可以按业务需要设置,一般保存一天到一周即可。
8. 工程落地时的安全边界与最佳实践
跨Agent通信这个短语听起来很酷,但它也是一种危险能力。一个Agent能向另一个Agent发消息,本身就意味着信任关系的建立。在实际项目中,安全问题不可跳过,这里做个强调。
8.1 Session所有权校验
当A Session的建设Agent准备与B Session的建设Agent打交道时,不要把B Session的内部数据一起透传。接收端必须校验:发送方Agent是否有权读取本Session的信息。每个Agent都应当维护一个最小权限原则,只接受与自身职责相关的事件,并在代码中显式声明。
# permission_check.py ALLOWED_SESSION_SCOPES = { "order_agent": {"read": ["order.*"], "write": ["order.*"]}, "notify_agent": {"read": ["order.completed", "payment.confirmed"], "write": []}, } def check_permission(target_agent: str, event_type: str, action: str = "read") -> bool: scopes = ALLOWED_SESSION_SCOPES.get(target_agent, {}) allowed = scopes.get(action, []) for pattern in allowed: if event_type.startswith(pattern.rstrip("*")): return True return False切勿把跨Session通信当成可以穿透Session安全边界的“后门”。就算业务上需要两个Session共享信息,也要通过一个明确的授权上下文来传递,例如用户同意共享、组织级授权或服务端签名令牌。
8.2 通信日志必须独立记录
跨Agent通信比普通内部函数调用更难追踪。普通函数调用可以依赖栈信息,而跨Agent消息是跳跃式的。生产环境里强烈建议给每条消息生成唯一traceId,并在消费端把该traceId写入日志。参考前面的代码里已经有event_id字段,但还不够,还需要一个从业务请求进入系统开始就一直存在的traceId。
# logger_setup.py import logging import sys logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(levelname)s] trace=%(trace_id)s agent=%(agent_name)s %(message)s", handlers=[logging.StreamHandler(sys.stdout)] ) def make_logger(name: str): logger = logging.getLogger(name) logger = logging.LoggerAdapter( logger, {"trace_id": "-", "agent_name": name} ) return logger实际项目中建议在HTTP入口或MQ消费入口处生成traceId,通过消息头或消息字段随链路传递。排查问题时,按traceId过滤日志,可以将一次跨Agent协作的完整过程还原出来。
8.3 通信过程中避免传递大对象
很多工程事故是通信层堆积大量数据导致的。AI Agent在生产过程中经常会产生长文本、工具调用结果、文档快照。把这些内容直接塞进跨Agent事件里,会造成两个问题:
- 消息体膨胀,序列化压缩和网络传输花费大量时间。
- 接收方Agent不需要原始内容,只需要引用,但会因为消息里已经带了完整数据而跳过了对数据源头的校验。
正确方式是使用“存储+引用”模式:消息里只放URI、主键ID或文件路径,数据存储于共享存储中心,接收方按需拉取。这样通信层职责清晰,只负责路由和通知,不负责数据搬运。
8.4 设置超时和重试策略
同步RPC式跨Agent调用必须设置超时。Agent内部往往包含多次LLM调用,耗时可长可短,因此超时时间不能设得太短,但也不能无限等待。
经验区间是200ms到60s之间。如果目标Agent快速处理,超时应小于1s;如果目标Agent可能要调用LLM或工具链,超时应放宽到10s以上。推荐采用两层超时设计:
- 连接超时:例如1s,防止目标不可达时长时间占用线程。
- 处理超时:例如30s或60s,取决于目标Agent业务量级。
重试策略则要区分消息类型:
- 查询类操作可以安全重试。
- 写入类和通知类操作必须配合幂等机制。
建议使用指数退避加抖动:1s、2s、4s、8s,最多重试5次。超过重试次数后放入死信队列,由人工或补偿流程处理。
9. 通信层会如何决定Agent架构的上限
从本文的示例可以看出,跨Agent通信和跨Session通信不是一个单点功能,而是整套Agent系统走向复杂化时必须建立的基础设施。
当你在写单体Agent应用时,可以在内部函数调用之间自由传参,通信似乎没有必要单独设计。可是Agent系统的趋势是把大量任务拆给不同角色承担,每个角色有自己的数据偏好和策略偏好,又必须服务于同一个用户目标。这时候通信层就决定了:
- 协作效率:靠轮询共享数据库,还是靠事件驱动,效率完全不是一个量级。
- 上下文同步及时性:Session A里用户修改了一个约束,Session B如果几秒内感知不到,就会给出旧答案。
- 安全可控性:通信链路清晰则权限边界清晰;各Agent通过私有数据库悄悄交互,权限就成了一笔糊涂账。
- 可靠性:消息有持久化、幂等、重试机制,业务流程才能承受Agent重启和网络抖动。
对正在做Agent开发的团队,我建议把这个通信层的设计提前到Agent角色定义完成后就开始,而不要等各个Agent独立实现完了再试图拼接。事后拼接的成本往往会高出三倍以上,因为Agent内部早已把Session私有状态写死了。
刚开始可以先用第四节的内存EventBus让本地逻辑跑通,再增强为第五节的Stream方案,最后再扩展网关与HTTP RPC。核心设计上,始终把AgentID、SessionID、event_id、event_type四个字段保留在消息体内,架构演进会顺畅很多。如果你正在开发多Agent协作系统,值得把跨Agent通信作为第一号技术基建来规划。