ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

跨Agent与跨Session通信:从消息总线到SQLite持久化

跨Agent与跨Session通信:从消息总线到SQLite持久化 Agent 项目一旦从单轮问答走向真实业务最先暴露的问题往往不是提示词写得不够好而是多个 Agent 之间不知道怎么传结果、多轮会话之间怎么续上上一次的进度。很多团队在演示环境里用单个 Agent 加一个聊天窗口跑得很顺进入开发环境后却发现编排层把任务拆给了执行 Agent结果执行 Agent 跑完没人拿结果用户刷新页面、换设备或者隔十分钟回来Session 里的上下文全部丢失只能重头再来。跨 Agent 通信和跨 Session 通信这两个能力就是用来解决这两类问题的。本文从概念讲起给出一个不依赖第三方框架的最小可运行示例先实现进程内消息总线再把会话历史落到 SQLite最后讨论生产环境下的选型取舍。学完后你可以把这套思路用于自研 Agent 编排模块也可以作为理解常见编排框架消息机制的基础。1. 先摸清两种通信分别卡在哪里1.1 跨 Agent 通信是“多个角色之间传任务”如果整个业务流程只由一个 Agent 在单次请求内完成其实不需要通信。需要通信的场景是任务被拆分规划 Agent 负责拆解目标执行 Agent 负责调用工具或生成内容审核 Agent 负责检查结果入口 Agent 负责汇总。这些角色之间必然要交换“任务描述”“执行结果”“错误信息”“状态变更”等数据这就是跨 Agent 通信。跨 Agent 通信和普通 RPC 调用有一个本质区别Agent 内部依赖了模型调用而模型调用存在推理慢、超时、返回格式不稳定、中间步骤失败等不确定性。因此消息不能只传一个“文本”还需要携带完整的任务标识、发送方、接收方、消息类型、优先级和可追溯信息。跨 Agent 通信的设计重点不是“能不能把变量传过去”而是“任务失败后由谁重试、结果返回后由谁消费、多个 Agent 并发执行时会不会冲突”。1.2 跨 Session 通信是“换一个会话仍然延续业务”Session 这个词在不同场景下的含义不同。Web 开发里的 Session 通常指服务端保存的用户会话状态用于维持登录态在 Agent 应用里Session 更多指向一次对话过程例如用户在聊天窗口里连续发送多条消息或者一条任务在浏览器刷新后仍需继续执行。跨 Session 通信要解决的核心问题是同一业务场景下新的技术 Session 能不能拿到历史 Session 产生的上下文。最直观的例子是用户第一次访问时提交了一个“整理一份面试复习计划”的任务Agent 生成了前半部分用户第二天打开页面浏览器新建了一个 Session如果系统不去查询历史记录用户只能重新发起请求。真正可用的 Agent 产品必须把“业务任务”和个人技术 Session 解耦用业务 ID 关联历史用存储介质保存进度让新的会话可以恢复旧的任务。1.3 用一张表分清两种通信对比维度跨 Agent 通信跨 Session 通信本质问题多个协作单元怎么交换任务和结果同一业务在不同会话期怎么保持上下文通信对象Agent 到 Agent用户在旧 Session 中产生的上下文到新的会话消息生命周期任务完成后消息通常可以清理需要长期保存用于恢复、审计和继续执行存储选型内存队列、消息中间件、Redis StreamSQLite、MySQL、Redis 快照、对象存储典型故障消息没人消费、重复消费、Agent 崩溃后任务丢失刷新后上下文丢失、跨设备无法找回、上下文太长超过模型限制这两类通信并不孤立。一个完整的 Agent 系统中跨 Agent 通信产生的消息记录同时也是跨 Session 恢复的数据来源Session 持久化保存的内容里核心部分往往就是 Agent 之间传递过的事件和结果。2. 设计通信前先约定会话、消息和上下文2.1 用统一 ID 把零散请求串成链路跨 Session 恢复之所以难并不是因为没有存储方案而是因为没有把“业务链路”标识出来。一次完整的用户目标可能经历多个技术 Session、多个 Agent、多次工具调用如果每段过程只有各自的随机 ID恢复时就没有办法按业务维度聚合。常见做法是同时存在两类 IDtask_id表示一次业务目标的完整生命周期例如“导出本月销售报表”这一个目标。session_id表示一次前端会话例如用户某一天打开的聊天窗口。当用户换了一个session_id继续同一条任务时系统应该通过task_id或者用户 ID 找到历史会话记录再把历史记录追加为后续处理的上下文。否则服务端只会看到两个互不关联的 Session任务自然无法延续。2.2 消息不要用裸字符串用信封结构很多自研 Agent 项目在初期会把通信写成“把这段文本发给另一个 Agent”。这种方式在只有两个角色时勉强能跑角色一多就会出现问题执行 Agent 不知道这条消息是该立即执行还是仅通知也不清楚消息来自规划器还是用户。消息需要一个稳定的信封结构。{ msg_id: msg_001, task_id: task_8f21a0c1, sender: planner, receiver: executor, msg_type: task.execute, payload: { task: 导出本月销售报表 }, create_time: 2025-01-01T12:00:00Z }msg_type比随机文本更可靠因为后续代码可以根据类型做分发、过滤、监控和重试。sender和receiver让通信可追踪payload只放业务数据。消息结构一旦确定跨 Agent 通信才有可能沉淀出可测试的协议而不是靠字符串拼接。2.3 上下文要区分“原始记录”和“可恢复快照”跨 Session 恢复时一个常见误区是把全部历史消息原样带入下一次模型调用。这么做在早期有效但任务执行几十步后历史消息可能超过模型窗口还会导致响应变慢、费用升高。更合理的做法是分层管理原始消息记录每条 Agent 消息都写入日志表用于故障排查和审计。压缩快照每一轮或每几轮执行后生成一份“当前任务做到哪一步、已确认信息、下一步计划”的摘要。近期明细只保留最近几条原始消息让模型能看到用户刚说过什么。恢复新 Session 时把“压缩快照 近期明细”拼接到提示词里而不是把整张历史表都塞进去。这条原则在后面实现示例时会用到。3. 最小可运行示例用进程内消息总线打通两个 Agent3.1 环境准备与项目结构这个示例使用 Python 3.10 以上版本只依赖标准库。这样做的目的是先把通信机制讲清楚避免读者在理解业务前还要先配置 Redis 或消息中间件。实际项目不要照搬内存方案但可以把这里的消息模型和角色划分保留下来。先创建目录结构agent_comms_demo/ ├── messages.py ├── bus.py ├── agents.py └── main.py3.2 定义消息模型消息模型对应上一节提到的信封结构。使用dataclass定义字段并给msg_id和create_time设置默认值这样调用方可以少传两个字段。# messages.py from dataclasses import dataclass, field import time import uuid dataclass class Message: sender: str receiver: str msg_type: str payload: dict msg_id: str field(default_factorylambda: uuid.uuid4().hex) create_time: float field(default_factorytime.time)这里的关键点是msg_id必须由发送方生成而不是由接收方生成。因为接收方如果需要做幂等处理就要用发送方的消息 ID 去重。如果 ID 在接收时重新生成重复消息就无法识别。3.3 实现一个线程安全的内存消息总线内存消息总线的思路是每个 Agent 对应一个队列send方法按receiver把消息投递到对应队列receive方法从该 Agent 自己的队列取消息。加锁是为了避免多线程同时读写同一个队列时出现数据错乱。# bus.py from collections import defaultdict, deque import threading from messages import Message class MessageBus: def __init__(self): self._queues defaultdict(deque) self._lock threading.Lock() def send(self, message: Message): with self._lock: self._queues[message.receiver].append(message) def receive(self, agent_name: str) - Message | None: with self._lock: queue self._queues[agent_name] if queue: return queue.popleft() return None def size(self, agent_name: str) - int: with self._lock: return len(self._queues[agent_name])receive在队列为空时返回None而不是阻塞等待。这种轮询方式在单进程示例里足够直观也容易理解和调试。生产环境建议换成阻塞队列或消息中间件的消费者组否则空轮询会浪费 CPU。3.4 编写两个轻量 Agent示例里规划 Agent 负责把用户请求拆成三步执行 Agent 负责消费并回传最终结果。实际产品中handle方法内部会调用模型或工具这里用字符串模拟让读者只看通信骨架。# agents.py from messages import Message class PlannerAgent: name planner def handle(self, msg: Message): task msg.payload.get(task, ) plan [ f校验任务{task}, 选择执行工具并生成命令, 验证产出文件并生成交付说明, ] print(f[planner] 拆解任务 - executor) return Message(self.name, executor, task.execute, {plan: plan}) class ExecutorAgent: name executor def handle(self, msg: Message): plan msg.payload.get(plan, []) summary - .join(plan) print(f[executor] 执行完成 - entry) return Message(self.name, entry, task.done, {summary: summary})注意发送给其他 Agent 的消息类型和普通业务消息不一样。task.execute表示执行任务task.done表示任务完成。接口设计上要避免让接收方通过猜测文本语义来决定行为类型字段才是分发依据。3.5 入口逻辑与运行结果入口函数先让“用户”把初始请求投递到 planner 的队列然后依次消费 planner 和 executor 的队列最后从 entry 队列取结果。# main.py from bus import MessageBus from messages import Message from agents import PlannerAgent, ExecutorAgent def drain(bus: MessageBus, agent_name: str, agent, max_iter: int 10): for _ in range(max_iter): msg bus.receive(agent_name) if msg is None: break bus.send(agent.handle(msg)) def main(): bus MessageBus() planner PlannerAgent() executor ExecutorAgent() # 用户请求以消息形式进入 bus.send(Message(user, planner, user.request, {task: 导出本月销售报表})) drain(bus, planner, planner) drain(bus, executor, executor) for _ in range(10): final_msg bus.receive(entry) if final_msg is None: break print(最终输出:, final_msg.payload[summary]) if __name__ __main__: main()在项目目录下执行python main.py预期输出类似[planner] 拆解任务 - executor [executor] 执行完成 - entry 最终输出: 校验任务导出本月销售报表 - 选择执行工具并生成命令 - 验证产出文件并生成交付说明这个示例说明了一个关键设计planner 和 executor 没有互相持有实例它们只依赖消息总线。后续替换通信实现时两个 Agent 的外部接口可以保持不变。4. 跨 Session把消息和上下文落到 SQLite4.1 为什么不能只留在内存进程内消息总线解决了“同一个进程内多 Agent 通信”的问题但没有解决“进程重启后数据丢失”和“多个会话之间状态不共享”的问题。用户刷新页面后如果状态还停留在 Python 进程的内存里进程一重启或者前端换了一个实例历史上下文就找不到了。跨 Session 通信的基础是持久化。只要 Agent 交互的历史被写入稳定存储新的 Session 就可以通过task_id或user_id找回旧数据。本节用 SQLite 做演示因为它是零配置的文件数据库适合学习阶段生产环境可以使用 MySQL、PostgreSQL 或 Redis 等组件设计思路不变。4.2 设计保存模型需要两张表一张存业务会话概要一张存 Agent 消息明细。实际场景还会增加工具调用记录、产物地址、错误日志等示例里只保留最小字段。CREATE TABLE IF NOT EXISTS session ( session_id TEXT PRIMARY KEY, user_id TEXT, task_id TEXT, context TEXT, updated_at TEXT ); CREATE TABLE IF NOT EXISTS message_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, session_id TEXT NOT NULL, sender TEXT NOT NULL, receiver TEXT NOT NULL, msg_type TEXT NOT NULL, payload TEXT NOT NULL, create_time TEXT NOT NULL );session表里的context保存的是上一节所说的压缩快照message_log表保存的是原始消息。恢复会话时两种数据配合使用。4.3 用 SessionStore 封装读写SessionStore的核心方法有三个保存一条消息、读取某个会话的全部消息、保存业务上下文快照。注意 SQLite 连接在示例里设置了check_same_threadFalse这只说明该连接允许多线程调用实际写操作仍然需要靠事务保证一致性。# session_store.py import json import sqlite3 from datetime import datetime, timezone class SessionStore: def __init__(self, db_pathsession.db): self.db_path db_path self.conn sqlite3.connect(self.db_path, check_same_threadFalse) self.conn.row_factory sqlite3.Row self._init_tables() def _init_tables(self): self.conn.executescript( CREATE TABLE IF NOT EXISTS session ( session_id TEXT PRIMARY KEY, user_id TEXT, task_id TEXT, context TEXT, updated_at TEXT ); CREATE TABLE IF NOT EXISTS message_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, session_id TEXT NOT NULL, sender TEXT NOT NULL, receiver TEXT NOT NULL, msg_type TEXT NOT NULL, payload TEXT NOT NULL, create_time TEXT NOT NULL ); ) self.conn.commit() def save_message(self, session_id, sender, receiver, msg_type, payload): now datetime.now(timezone.utc).isoformat() self.conn.execute( INSERT INTO message_log(session_id, sender, receiver, msg_type, payload, create_time) VALUES (?, ?, ?, ?, ?, ?), (session_id, sender, receiver, msg_type, json.dumps(payload, ensure_asciiFalse), now), ) self.conn.commit() def load_messages(self, session_id): rows self.conn.execute( SELECT sender, receiver, msg_type, payload, create_time FROM message_log WHERE session_id ? ORDER BY id, (session_id,), ).fetchall() return rows def save_context(self, session_id, user_id, task_id, context): now datetime.now(timezone.utc).isoformat() self.conn.execute( INSERT INTO session(session_id, user_id, task_id, context, updated_at) VALUES (?, ?, ?, ?, ?) ON CONFLICT(session_id) DO UPDATE SET user_id excluded.user_id, task_id excluded.task_id, context excluded.context, updated_at excluded.updated_at , (session_id, user_id, task_id, json.dumps(context, ensure_asciiFalse), now), ) self.conn.commit()每次修改都执行commit在示例里没有太大问题但在高频调用场景会带来磁盘 IO 压力。生产项目建议按批提交或者直接使用为高并发设计的内存型存储。4.4 演示新的 Session 恢复旧的历史下面用一段代码模拟跨 Session 场景。第一个session_id是用户昨天的会话第二个session_id是用户今天重新打开页面产生的新会话。恢复的核心不是把两个 Session 强行合并而是通过业务用户 ID 找到旧会话再把旧历史拼接进新会话的处理上下文。# session_demo.py from session_store import SessionStore def run_turn(store, session_id, user_text): store.save_message(session_id, user, planner, user.request, {content: user_text}) # 真实场景这里会调用模型和 Agent本示例只记录结果 store.save_message(session_id, planner, entry, task.done, {summary: 已经生成 user_text}) def main(): store SessionStore(session.db) # 第一次会话用户提出需求Agent 完成两轮处理 run_turn(store, session-001, 整理一份技术面试复习清单) run_turn(store, session-001, 重点补充计算机网络部分) # 第二次会话模拟用户第二天重新打开页面 old_history store.load_messages(session-001) print(恢复出的旧会话消息条数:, len(old_history)) for row in old_history: print(row[sender], -, row[receiver], row[msg_type]) # 新会话写入新消息同时继续引用旧会话 ID run_turn(store, session-002, 在上次复习清单基础上补充算法题) print(新会话也已保存) if __name__ __main__: main()执行python session_demo.py预期输出恢复出的旧会话消息条数: 4 user - planner user.request planner - entry task.done user - planner user.request planner - entry task.done 新会话也已保存运行后用 SQLite 命令行可以直接看到表里的数据sqlite3 session.db SELECT session_id, sender, msg_type FROM message_log ORDER BY id;这里要特别说明示例里的session-002并不是跨 Session真正起作用的是“通过业务用户 ID 找到session-001”这一步。如果代码里只保存了session-002自己的消息没有关联旧的task_id或user_id那无论存多少数据都无法实现续聊。5. 生产环境的演进路径从内存到 Redis 再到消息中间件5.1 不同存储方案的适用边界内存消息总线和 SQLite 只适合教学和原型验证。真实项目需要根据部署形态选择通信存储这里整理几种常见方案。方案适合场景主要限制进程内队列 内存单进程、单实例、学习 Demo重启丢消息、无法跨实例Redis List / Stream多实例共享队列、需要简单重试需要运维 Redis消息追溯能力弱Kafka / RabbitMQAgent 数量多、任务量大、需要可靠投递组件重学习成本高MySQL / PostgreSQL需要保存会话记录并支持业务查询不适合做高吞吐临时消息通道选型时不要只看某个组件流行要先确认你要解决的是“多个 Agent 同步调用”还是“异步任务可靠分发”。示例里的同步轮询方式在 Agent 数量少时够用但一旦出现一个 Agent 同时处理多个任务就要给消费过程增加确认机制否则可能会把一个任务分给多个执行者。5.2 协议字段和幂等策略要先定无论底层用哪个组件消息字段和消息消费语义都应该先定好。消息里最好包含schema_version方便以后协议升级时做兼容处理。比如第一版字段是payload第二版改成data没有版本号就只能靠猜。消费语义也要明确是“最多一次”还是“至少一次”。如果允许 Agent 执行失败后重试消息接收方必须幂等也就是同一个msg_id执行两次不能产生两份结果。最简单的幂等方案是维护一张processed_message表处理前先查msg_id是否已经存在不存在才继续执行。5.3 学习环境与生产环境的差异关注点学习环境生产环境消息存储内存队列Redis Stream 或消息中间件会话存储SQLite 文件数据库或 Redis考虑高可用消息确认不需要必须处理成功确认、失败重试日志直接打印按链路 ID 记录完整日志监控不关注采集消息积压、消费延迟、失败率清理策略无配置消息 TTL 和过期会话清理跨 Session 通信在生产环境还要额外处理一个问题会话数据不能无限增长。用户长时间不使用后旧会话可以转存为压缩摘要删除明细真正需要审计的原始记录可以移到冷存储。6. 常见问题排查为什么消息收不到、会话找不回6.1 消息收不到或没人消费如果消息总线投递正常但接收方取不到消息优先按以下顺序检查。先确认receiver名称是否一致。很多问题出在发送方写的是executor接收方注册的名字是executor_agent。再确认send是否真的执行了。检查总线某个队列的大小如果大小为 0说明问题在发送方。接着确认取消息的循环是否在发送动作之前就已经退出。示例里drain使用空队列退出如果时序写错接收方可能先轮询一遍发现为空就退出。最后看角色名是否和业务角色对应。结果投给了entry消费者却在监听done同样会拿不到。6.2 Agent 挂了之后消息丢失现象是执行 Agent 处理到一半进程退出重新启动后任务没有继续。可能原因是消息已经出队但尚未处理完进程退出后队列里不再有这条消息。解决方式是把“消息已接收”和“消息已处理完成”分开接收后先记录日志处理完成后再发送确认或删除消息。使用 Redis Stream 时可以利用消费者组和 Pending Entries 重投机制使用数据库时可以在消息表里增加status字段。6.3 会话记录存在但恢复出来为空这类问题最常见的根因不是存储坏了而是查询时用错了session_id。用户在新页面里产生的是session-002代码却只取了session-002的消息忘记了通过用户 ID 关联session-001。检查方式是在应用里打印业务用户关联的旧 Session ID再单独查询数据库确认数据是否存在。另一种可能是在多个环境切换时把数据库文件指向了不同路径。SQLite 是单文件数据库测试环境和本地环境用同一个session.db路径会互相覆盖思维混乱。解决方式是让数据库路径和部署环境绑定并在日志里打印实际路径。6.4 上下文越来越长导致模型调用失败现象是刚开发时运行正常执行几轮后请求超时或者提示超出最大 token 限制。原因是每次新 Session 都把所有历史明细拼进提示词。正确做法是先把旧历史做摘要再只保留最近几条原始消息。摘要可以放在session表的context字段里每次新任务开始前先更新摘要恢复时读取摘要而不是全量历史。6.5 多线程读写同一 SQLite 连接报 database is locked示例里为了简单复用了同一个连接但多线程高并发写同一个 SQLite 文件时容易出现锁冲突。生产环境不要让多个 Agent 线程共享同一个写连接推荐做法是每个写操作使用新的短连接或者把 SQLite 换成真正的数据库服务。问题现象常见原因检查方式处理建议消息收不到receiver 名字不一致打印发送方和接收方名称统一角色命名规范Agent 崩溃后任务丢失消息出队后未确认查看进程崩溃日志增加已处理状态或消费者组会话恢复为空查询了新会话没有关联旧会话检查业务用户 ID 映射建立 task_id 与 user_id 关联上下文超长全量历史塞进模型统计每次 prompt 字符数摘要 最近 N 条消息database is locked多线程共用一个 SQLite 连接查看异常堆栈
返回列表