持久化、跨进程恢复与 Graph 版本
长任务不能依赖单个进程一直存活。恢复 Agent Graph 时,需要区分 Flow 快照、Agent Trace 和 ChatSession 消息历史。
1、三类状态
| 数据 | 所在位置 | FlowContext.toJson() 是否覆盖 |
|---|---|---|
| Graph 恢复位置 | FlowTrace | 是 |
| 业务输入和中间结果 | FlowContext.data() | 可序列化值是 |
| Team/Simple/ReAct Trace | Context 中的 Trace Key | 是 |
| 原始 Prompt | Agent Trace | 是 |
| pending 原因 | Context | 是 |
| ChatSession 历史消息 | Session 消息存储 | 否 |
| Session、模型、Driver、事件监听器 | 运行时对象 | 否 |
FlowTrace 按 Graph ID 保存最后节点,是恢复游标集合,不是完整节点历史。
2、保存 Context 快照
AgentSession session = InMemoryAgentSession.of("change-001");
changeGraph.prompt("升级订单服务")
.session(session)
.call();
if (session.isPending()) {
String snapshot = session.getContext().toJson();
snapshotRepository.save("change-001", snapshot);
}
快照可能包含用户输入、模型结果和审核数据,应加密、脱敏、授权并设置过期策略,不要直接打印到日志。
3、跨进程恢复
应用重启后先重新构建同版本 Agent 和 Graph,再绑定恢复 Context:
String snapshot = snapshotRepository.load("change-001");
FlowContext restored = FlowContext.fromJson(snapshot);
AgentSession session = InMemoryAgentSession.of(restored);
restored.put("approval.decision", "APPROVE");
restored.put("approval.operator", "alice");
TeamResponse response = changeGraph.prompt()
.session(session)
.call();
InMemoryAgentSession.of(restored) 只恢复并绑定 Context,不会恢复旧 InMemory Session 的 ChatSession 消息列表。若续跑依赖聊天历史,应由持久化 Session 同时恢复消息。
4、为什么使用空 prompt()?
changeGraph.prompt()
表示从 Agent Trace 中保存的原始 Prompt 继续任务。传入新 Prompt:
changeGraph.prompt("新的问题")
会被识别为新任务,并重新准备协作状态。一个从未执行过的新 Session 也不能靠空 prompt() 凭空获得任务。
5、运行时对象需要重新装配
快照不会替你恢复:
- ChatModel、Agent、TeamAgent 和 Graph 定义;
- FlowEngine、Driver、Executor;
- FlowOptions 和各类 Interceptor;
FlowContext.eventBus()的监听器;- 流式订阅、线程、锁和回调;
- 数据库、HTTP、Redis 等连接。
恢复顺序:
加载指定版本的 Graph 定义
-> 读取 Context 快照
-> 创建并绑定 AgentSession
-> 重新注册监听和运行时组件
-> 写入人工决定或外部事件
-> prompt() 续跑
6、Session 实现的持久化边界
InMemoryAgentSession
只适合开发、测试和单进程短任务,不会自动落盘。
FileAgentSession
可持久化消息和快照,但 Context 的普通 put(...) 不等于立即写盘。若写入人工决定后暂不恢复执行,应显式调用:
session.updateSnapshot();
RedisAgentSession 或自定义实现
多实例部署通常需要 Redis 或数据库 Session。以本系列核验时的源码基线 4.1.1-SNAPSHOT 为准,内置 RedisAgentSession 的快照写入键与构造时读取键并不一致,不能直接承诺 updateSnapshot() 后重新构造 Session 即可恢复。该问题后续版本可能修复,使用时应以对应版本源码和闭环测试为准;在修复前,可采用经过验证的自定义 AgentSession。
修正后必须增加“保存后重新构造 Session”的测试:
first.getContext().put("approval.decision", "APPROVE");
first.updateSnapshot();
AgentSession second = createPersistentSession("change-001");
Assertions.assertEquals("APPROVE",
second.getContext().get("approval.decision"));
生产实现还应加入分布式锁或乐观版本控制,避免同一实例被多个节点同时恢复。
7、Graph ID 和定义版本
TeamAgent 名称不等于内部 Graph ID。默认情况下:
team.name = change_graph
flow.graph.id = __change_graph
definition.version = 2
instance.id = change-001
发布时删除或改名恢复节点,旧快照可能无法继续。建议:
- 运行中实例继续使用旧定义;
- 保存
definition.version; - 为旧快照编写迁移器;
- 在稳定 Activity 边界设置检查点;
- 谨慎跨版本恢复 Parallel、Inclusive 和循环的中间状态。
Graph ID 应在共享执行树和同一 FlowEngine 注册空间内唯一。
8、恢复节点必须幂等
挂起节点恢复后会重新执行。外部副作用节点应使用业务幂等键:
spec.addActivity("execute_change")
.task((ctx, node) -> {
String key = ctx.getAs("execution.idempotencyKey");
String oldResult = operationRepository.find(key);
if (oldResult != null) {
ctx.put("change.result", oldResult);
return;
}
String result = executeChange();
operationRepository.save(key, result);
ctx.put("change.result", result);
});
9、恢复测试
不要只断言最终文本。为关键节点增加确定性计数:
int count = ctx.getOrDefault("count.proposal", 0);
ctx.put("count.proposal", count + 1);
恢复后至少验证:
- 已完成节点没有重复运行;
- 挂起节点按恢复语义重跑;
- pending 已清理;
- Context 中间数据保留;
- 父、子 TeamTrace 都能定位;
- Graph 最终到 End;
- 消息和快照持久化形成闭环。