跨Agent与跨Session通信:消息路由与状态同步设计

发布时间:2026/9/3 2:51:52
跨Agent与跨Session通信:消息路由与状态同步设计 跨Agent通信、跨Session通信这两个概念在Agent工程里变得越来越刚需。单Agent、单Session的模式很容易跑通用户发一条消息Agent带着上下文调用模型模型返回结果再调用工具闭环完成。可一旦进入真实业务场景会迅速变成多Agent协作、多Session并行用户可能在会话A里派了一个任务又在会话B里等结果planner Agent拆完任务后要交给coder Agentcoder Agent完成后要回传给另一个Session。这个时候缺的往往不是模型能力而是Agent之间、Session之间的通信设计。这篇文章我不打算绑定某个具体的Agent框架而是从通用工程视角把跨Agent通信和跨Session通信拆开讲清楚。文章会覆盖概念边界、消息结构设计、Topic路由、Session状态归属、本地可运行的通信演示、接口API与批量任务模板、资源占用观察以及常见问题排查。整个通信层并不依赖GPU所以哪怕你手上只有一台普通开发机也可以把链路完整跑一遍。适合的读者有两类一类是正在做Agent框架或对话系统发现“多个Agent一协作就乱”的开发者另一类是已经跑通单Agent但准备做任务编排、会话状态同步、批量任务投递的工程师。下面直接进入主题。1. 核心能力速览跨Agent通信与跨Session通信1.1 先讲清楚 Agent、Session 和通信“Agent”“Session”“通信”这三个词在不同技术语境下含义差别很大。在Web开发里讨论Session默认是服务端会话状态而Agent可能又是指某种自动化框架在嵌入式或网络编程里通信又可能是UDP、SPI、CAN这类协议。但这篇文章讨论的语境是AI Agent工程概念边界需要先统一Agent一个能自主规划、调用工具、执行任务的智能体单元常见的是planner规划智能体、coder编码智能体、reviewer审查智能体这类角色分工。Session一次业务执行过程的上下文容器里面保存消息历史、临时状态、当前目标、任务归属等信息。Session可能是一次用户对话也可能是一条长时间运行的后台任务。跨Agent通信解决Agent之间的消息传递、任务接力、结果回传。跨Session通信解决多个Session之间的状态同步、事件通知、任务会话切换。跨Agent通信和跨Session通信不是二选一而是经常同时发生。比如Session A中的planner Agent把一个编码任务交给Session B中的coder Agent这条消息既跨越了Agent边界也跨越了Session边界。对系统设计者来说如果一开始没有把Agent ID、Session ID、消息关系这些建模清楚后面很容易出现“消息不知道发给谁”“会话上下文互相污染”“任务重复执行”这些问题。1.2 能力速览能力项说明跨Agent消息路由消息按目标Agent或Topic路由实现任务定向分发跨Session状态同步多个Session之间共享必要上下文但不默认共享全量数据任务接力Session A中的Agent把任务交给Session B中的Agent执行完成后回传结果事件广播向一组订阅了某个Topic的Session或Agent下发通知批量任务将多个Session的任务写入队列设置并发后依次消费接口API通过HTTP或内部RPC接口触发Agent运行支持外部系统集成最小运行环境不依赖GPUPython或Node等语言环境即可完成本地验证生产推荐组件Redis Streams、RabbitMQ、Kafka、NATS、数据库消息表等典型风险Session隔离不严、消息幂等缺失、路由Topic不匹配、任务无限重试通信层本身不需要大显存或高端显卡核心是消息与状态的设计这也是为什么它比单纯调用Agent更容易被低估。硬件条件再好的团队如果消息落不到正确的SessionAgent之间的协作依然会失败。2. 适用场景与使用边界2.1 适合解决什么问题先说最典型的三类场景。第一类是多Agent任务编排。一个项目里不可能只有一个Agent常见的拆法是planner负责拆任务coder负责写代码reviewer负责审查tester负责跑测试。这些Agent如果每次都用直连函数互相调用系统会很快被耦合死但它们之间如果只靠一个共享变量传消息又很容易丢消息。更稳妥的方式是让Agent通过统一消息总线通信再配合Session ID做任务链路追踪。第二类是跨会话的任务接力。用户可能在会话A中创建了一个数据清洗任务随后又在会话B中发起了另一个模型训练任务。planner Agent需要把“数据清洗”的中间结果带给Session B中的Agent。这里的难点不是保存结果而是要让Session B知道这个结果来自哪个Session、由哪个Agent产生、应该在哪一步使用。第三类是事件驱动的协同状态更新。比如一个监控Agent发现线上服务异常它需要通知多个工作Session停止当前任务并进入等待状态。这本质上是发布订阅模型监控Agent发布“服务异常”事件订阅该事件的Session各自做出响应。跨Session通信在这里表现为一种事件扩散能力而不是简单的点对点调用。2.2 不该强行套用的场景与合规边界跨Agent通信和跨Session通信并不适合所有情况。如果两个Agent的调用逻辑在同一个事务中且不需要异步解耦直接用函数调用或RPC会更简单引入消息队列反而增加延迟和排错成本。低延迟场景也不能把消息队列当作唯一答案消息队列模型更擅长削峰和解耦不擅长微秒级通信。另一个必须强调的问题是数据边界。Session之间一旦具备通信能力就要防止越权读取。Session A不能因为能向Session B发消息就读取Session B里的私密内容也不能在用户未授权的情况下把Session A的用户信息广播给所有在线Session。涉及用户输入、人脸、音频、版权素材等数据时必须先确认授权口径再做脱敏和最小化传递通信日志要保留审计记录。跨Session通信是技术能力不能成为数据滥用的通道。3. 环境准备与前置条件3.1 最小运行环境本地验证这套通信能力并不需要一开始就引入队列、Redis或Kubernetes。先准备一台开发机建议安装Python 3.10以上版本再准备一个虚拟环境就可以把通信链路跑通。演示代码会使用asyncio实现进程内消息总线处理器之间通过异步函数传递消息清晰展示“跨Session”的语义。生产环境则建议引入一个真正的外部消息组件因为它能解决进程重启丢消息、多实例水平扩展、消息持久化这三个问题。组件选型按照团队熟悉度和数据规模来定小规模项目可以用Redis Streams或RabbitMQ复杂任务链路可以用Kafka或NATS。通信层建议封装接口业务模块不直接依赖具体消息组件方便后面替换。3.2 目录结构规划如果是自己搭一套演示代码建议按下面这个结构组织agent-demo/ ├── bus.py # 进程内消息总线 ├── runtime.py # Agent运行时与Session路由 ├── agents.py # 两个演示Agent ├── main.py # 启动入口 └── requirements.txt # Python依赖模型调用部分可以先省略用模拟耗时代替目的是把通信链路的逻辑验证清楚。后面接入真实大模型时只需要在Agent节点的处理函数中把模型调用补上不需要改动消息链路。4. 通信模型设计消息、Topic 与路由4.1 消息结构必须带上路由和追踪标记跨Agent通信和跨Session通信的第一个设计决策是消息格式。很多早期项目会把消息简化成{content: xxx}看起来简单但在多Agent、多Session场景下会立刻不够用。建议消息至少包含以下字段{ message_id: msg_20250401_001, request_id: req_20250401_001, source_session_id: session_a, source_agent: planner, target_session_id: session_b, target_agent: coder, topic: planner.to.coder.task, event_type: task.create, payload: { task: 实现登录接口, params: {} }, created_at: 2025-04-01T10:00:00Z }message_id用于幂等去重同一任务重复投递时消费端可以根据message_id判断是否已经处理过request_id是一条用户请求或业务请求的全局标识贯穿多个Agent和Sessionsource_session_id和source_agent负责溯源告诉接收方消息来自哪条链路target_session_id和target_agent用于精确路由。event_type描述消息的业务语义方便消费端做状态判断created_at用于排查延迟和超时。如果消息不用JSON也可以用Protobuf或Avro这类序列化协议但是字段语义要保持一致。消息是Agent之间唯一能达成共识的结构化对象字段设计得越清晰后面扩展路由时越省力。4.2 Topic 与订阅关系消息总线需要一个Topic来组织订阅关系。Topic的命名最好和Agent角色、业务动作绑定比如planner.to.coder.task表示planner发给coder的任务coder.to.planner.review表示coder交给planner审查的结果。这样命名有好处监控日志时一眼能看出消息链路做权限控制时也能基于Topic快速拦截非法消息。订阅关系可以简单理解为一个注册表# 伪代码向消息总线注册处理器 bus.subscribe(planner.to.coder.task, coder_on_task) bus.subscribe(coder.to.planner.review, planner_on_review)发布端和订阅端不需要知道彼此的地址发布端只负责把消息投递到总线订阅端注册了自己关心的Topic后就会收到消息。在进程内演示阶段这样已经够用在生产环境把bus替换成Redis或RabbitMQ的客户端这个模型依然成立。4.3 Agent 之间采用什么通信模型跨Agent通信至少有三种模型可以根据场景选择。点对点模型适合明确的单跳分工一个Agent发起请求另一个Agent处理并返回结果。优点是简单、可追踪缺点是如果链路长Agent之间仍然会产生直接依赖。发布订阅模型适合事件广播监控Agent发布告警事件多个下游Agent各自订阅并处理。优点是可扩展新增订阅者不影响发布者缺点是发布者并不知道事件是否被人真正处理可能需要在事件后增加确认机制。任务队列模型适合批量任务外部系统把大量任务写入队列一组Agent Worker并行消费。优点是天然支持削峰填谷和失败重试缺点是处理结果的顺序与请求顺序不一定一致需要靠request_id重新关联。很多复杂系统并不是只用一种模型而是以任务队列为主线在局部使用点对点回调再用事件广播做全局通知。通信层设计得好不好关键看你能不能把三种模型封装在同一个消息API下。5. 跨 Session 通信的核心难点与设计5.1 Session 状态到底放在哪里Session本身就是一份状态但通信模型会影响状态存放的位置。很多人误以为跨Session通信就是直接把两个Session的上下文塞进同一个内存池这是不安全的做法。Session状态应当分层存放运行时热数据放在进程内存或Redis里例如当前正在执行的工具调用信息、本轮未结束的上下文需要长期保存的消息历史放在数据库或对象存储中跨Session共享的中间产物放在一个显式的共享区域通过session_id和request_id关联。关键原则是Session可以持有私有状态只有被明确标记为共享状态的事件消息才进入通信总线。比如Session A的完整消息历史不需要同步给Session BSession B只需要知道“Session A已完成前端页面开发物料存放在某个路径下”。5.2 Session 访问边界跨Session通信最常见的故障是会话串号。现象是Session B收到了本该属于Session C的消息或者Session A的中间变量被Session B覆盖。原因通常是路由层没有校验Session关系或者接收方只按Agent ID过滤没有按Session ID过滤。设计上要把Session ID当作路由的核心字段。消息投递前路由层需要校验target_session_id和source_session_id是否具备合法的共享关系。可以把Session之间的共享关系做成一张调度表或白名单由业务方控制哪些Session允许共享事件哪些Session必须隔离。Session隔离和Agent实例可复用是两回事。同一个coder Agent实例可以处理多个Session的任务但处理完Session A的任务后必须清除运行时上下文不能把Session A的参数带到Session B里。代码实现上Agent处理函数不应当保留实例字段级别的可变上下文而是把上下文作为参数传入。5.3 从 Session A 派生子 Session跨Session通信的常见业务形态是子Session派发。用户进入Session Aplanner Agent判断任务链路较长于是创建一个子任务Session B再把编码任务派给Session B中的coder Agent。此时Session B需要记录parent_session_idsession_a形成一棵会话树。子Session的好处是隔离任务现场主Session可以继续响应用户新指令子Session专注后台任务互不阻塞。主Session如果中途关闭子Session需要根据策略决定是终止还是继续执行子Session完成后也需要把结果回写成主Session可见的事件。这些状态迁移需要明确的状态机来管理否则会出现任务完成后无人接收结果或者Agent被释放但Session仍被占用的问题。Session状态可以用以下状态简化表达状态含义典型流转createdSession已创建创建后进入running或waitingrunning有Agent正在执行任务Agent回调后进入completed或failedwaiting会话在等待外部事件或子Session收到事件后进入runningcompleted会话任务完成结果可被其他Session读取failed会话执行失败进入重试或终止timeout会话超时触发会话清理逻辑6. 本地可运行示例两个 Agent 的跨 Session 任务接力6.1 最小异步消息总线没有消息队列之前先用一个进程内消息总线把逻辑跑通。下面这段代码定义MessageBus支持按Topic注册异步处理器# bus.py # 进程内异步消息总线仅供本地演示 import asyncio from collections import defaultdict from typing import Awaitable, Callable Handler Callable[[dict], Awaitable[None]] class MessageBus: def __init__(self): self._handlers: dict[str, list[Handler]] defaultdict(list) def subscribe(self, topic: str, handler: Handler) - None: self._handlers[topic].append(handler) async def publish(self, topic: str, message: dict) - None: handlers list(self._handlers.get(topic, [])) for handler in handlers: await handler(message)这个总线的核心逻辑只有三部分按Topic保存处理器、注册订阅、发布消息时逐个调用处理器。它没有持久化也没有多进程能力但足以验证消息路由和跨Session语义。6.2 定义带 Session 上下文的 Agent 消息为了让消息携带跨Session信息需要定义一个消息类。source_session_id记录消息来自哪个Sessiontarget_session_id记录消息要送往哪个Sessiontarget_agent记录目标Agent角色。# models.py # 消息结构同时满足跨Agent通信和跨Session通信的最小模型 from dataclasses import dataclass, field from datetime import datetime, timezone from uuid import uuid4 def _now_iso() - str: return datetime.now(timezone.utc).isoformat() dataclass class AgentMessage: source_session_id: str target_session_id: str target_agent: str payload: dict source_agent: str message_id: str field(default_factorylambda: uuid4().hex) created_at: str field(default_factory_now_iso)消息投递时发布方只需要指定目标Session、目标Agent和payload消息总线从Topic中解析目标Agent目标Agent再根据target_session_id决定当前处理的会话上下文。字段越少演示越清楚真正接入生产时再补充request_id、event_type等字段。6.3 两个 Agent 完成跨 Session 接力下面模拟一个实际链路Session A中的planner Agent创建一个编码任务消息投递给Session B中的coder Agent。coder完成处理后将结果回传给Session A。# agents.py # 两个演示Agentplanner负责派发coder负责编码 import asyncio class PlannerAgent: def __init__(self, bus): self.bus bus # planner 监听结果回传主题 self.bus.subscribe(agent.io.planner.result, self.receive_result) self.agent_id planner async def dispatch_task(self, source_session_id: str, target_session_id: str): message { source_session_id: source_session_id, target_session_id: target_session_id, source_agent: self.agent_id, target_agent: coder, payload: { task: 实现一个登录接口包含参数校验与Token签发 } } print(f[planner] 向Session {target_session_id} 的coder投递编码任务) await self.bus.publish(agent.io.coder.task, message) async def receive_result(self, message: dict): print(f[planner] 收到来自Session {message[source_session_id]} 的编码结果) print(f[planner] 结果摘要: {message[payload].get(summary)}) class CoderAgent: def __init__(self, bus): self.bus bus # coder 只监听任务主题 self.bus.subscribe(agent.io.coder.task, self.execute_task) self.agent_id coder async def execute_task(self, message: dict): source_session_id message[source_session_id] source_agent message[source_agent] task_desc message[payload][task] print(f[coder] 收到来自Session {source_session_id} 的任务: {task_desc}) # 模拟耗时处理过程 await asyncio.sleep(0.5) # coder 任务完成后把结果回传给来源Session中的planner result_message { source_session_id: message[target_session_id], target_session_id: source_session_id, source_agent: