存过了,所以能回到那个位置
目标
在加了 checkpointer 的图中,亲手数出什么留在了哪里。确认同一个 thread_id 会延续,把检查点读成列表,回溯到过去的坐标,修改值分叉出新分支,把状态序列化后数字节。最后重现重新建立账簿后会有什么消失。
为什么重要
checkpointer 在每一步把那一时刻的状态整个留下。这一句话同时解释了好事和坏事。
好的一面是可以接着运行。即使在第八步失败,也不必从头再来,对于“如果第三步用了别的资料会怎样”这个问题,也可以在同样的条件下回答。用过去快照的 config 调用 invoke(None, ...),就会在那个位置重新运行,用 update_state 修改值放进去,就会从那里产生分支。此时原来的分支不会被删除,但线程的“当前”会转移到新分支,所以想重新看原来的结果,就得把那一时刻的 config 拿在手里。
坏的一面是保存的是状态整体。在状态中放入一个大的值,那个值会按检查点的数量被重复保存。本实验中用字节而不是时间来测量——时间会因机器和负载而不同,而同样的状态大小总是一样的。
评分器不会相信你写下的说明。它会真正导入你的模块运行图,自己数出检查点,并与你返回的值对照。主题和放入的值的大小每次运行都会变。
步骤
- 在 /root/work/agckpt/ckpt.py 中创建
ANGLES、State、三个节点(plan、write、review)、build_graph(checkpointer=None)、new_saver()、cfg(thread_id)。相同的thread_id要能延续,不同的要各自保留。 - 增加
history(app, config)、history_len(app, config)、pending(app, config),从最旧的开始读取并统计检查点。 - 增加
snapshot_before(app, config, node),找出还没有运行该节点的最近一个检查点。 - 增加
replay(app, config, node),在该坐标上用invoke(None, snapshot.config)重新运行。 - 增加
fork(app, config, node, values),用update_state分叉出新分支,同时返回原分支的末尾。 - 增加
state_bytes(values)、history_bytes(app, config)、weigh(topic, payload)和状态键bulk,以字节衡量检查点的大小。 - 增加
lost_demo(topic),重现重新建立账簿之后,即使用同一个thread_id也什么都不剩。 - 在 /root/work/agckpt/ckpt_report.json 和 /root/work/agckpt/ckpt_report.md 中记录测得的值。
参考
- 执行契约:评分器会把
/root/work/agckpt/ckpt.py当作 Python 模块导入,直接使用ANGLES、State、build_graph、new_saver、cfg、history、history_len、pending、snapshot_before、replay、fork、state_bytes、history_bytes、weigh、lost_demo。它不会作为脚本运行,所以可以没有if __name__ == "__main__"。 ANGLES = {"기본": 1, "요약": 2, "비교": 3}(韩文,意为“基本”“摘要”“对比”)。它记录每个 angle 名称下草稿要重复几次。- 节点做的事正是这样。
plan返回{"stage": "plan", "angle": 지금 angle 또는 "기본", "notes": ["plan:<topic>"]}(韩文,意为“当前的 angle,或‘基本’”),write把"<topic>/<angle> "重复ANGLES[angle]次并去掉两端空格后放入draft,并返回{"stage": "write", "notes": ["write:<angle>"]},review返回{"stage": "review", "score": len(draft), "notes": ["review:<score>"]}。 notes加上operator.addReducer。其余的键是覆盖。build_graph(checkpointer=None)以graph.compile(checkpointer=checkpointer)收尾。不给 checkpointer 就什么都不会保存。cfg(thread_id)返回{"configurable": {"thread_id": thread_id}}。get_state_history(config)从最新的开始返回。history()要把它反转,从最旧的开始返回。pending()是一个列表,装着每个快照的next的第一个元素,为空则为空字符串。replay()返回{"from": node, "result": 결과, "added": 늘어난 체크포인트 수}(占位符依次为结果与增加的检查点数量)。fork()返回{"forked_config": ..., "result": ..., "original": 분기 전 끝의 값 딕셔너리, "total": 지금 체크포인트 수}(占位符依次为分叉之前末尾的值字典与当前检查点数量)。state_bytes(values)是把json.dumps(values, ensure_ascii=False, sort_keys=True, default=str)得到的字符串以 UTF-8 编码后的长度。评分器会用同样的方式计算并对照,所以请原样使用这三个参数。weigh(topic, payload)返回{"checkpoints": ..., "thin_total": ..., "fat_total": ..., "payload_bytes": ..., "carrying": ...}。carrying是把相同位置的检查点逐一比较,大小之差达到payload_bytes以上的检查点个数。lost_demo(topic)返回{"kept": ..., "after_restart": ..., "values_after_restart": ...}。线程名称请使用run-1。- 这个 Pod 没有互联网。
pip install无法使用。langgraph 0.2.60 已经装好了(python3 -c "import langgraph")。 - 官方文档:Persistence · Use time-travel · Graph API overview
- 常见错误:没给
compile()传 checkpointer 就调用get_state(ValueError: No checkpointer set)、回溯时给输入传状态而不是None(会重新开始)、丢掉update_state返回的 config 而用原来的 config 接着运行、不反转get_state_history的顺序。
同一个线程会延续,不同的线程各自保留
在 /root/work/agckpt/ckpt.py 中创建 ANGLES、State、三个节点(plan、write、review)、build_graph(checkpointer=None)、new_saver()、cfg(thread_id)。加了 checkpointer 的图用同一个 thread_id 再次运行时,要累积在之前的 notes 之上,不同的 thread_id 要从空白开始。
一行 graph.compile(checkpointer=checkpointer) 就够了。把 cfg(thread_id) 做成返回 {"configurable": {"thread_id": thread_id}} 的小函数,后面的步骤会方便。必须给 notes 加上 operator.add,延续才看得见——不加的话,虽然保存了,但值会被覆盖,看不出延续。
检查点每一步都会留下
增加 history(app, config)、history_len(app, config)、pending(app, config)。history() 把检查点从最旧的开始排列,pending() 返回一个列表,装着每个快照接下来要运行的节点名称(已全部结束的位置为空字符串)。
app.get_state_history(config) 是生成器,从最新的开始给出。用 list() 接住后 reversed()。快照的 next 是元组——为空就意味着没有要运行的了。数一数会比节点数多两个,因为接收输入的位置和全部结束的位置各多出一个。
找出回溯的坐标
增加 snapshot_before(app, config, node)。返回还没有运行该节点的最近一个检查点,没有这样的位置就返回 None。
查找条件是快照的 next 为 (node,)。把从最旧的开始排列的列表从后往前扫,会先遇到“最近的”。返回的是快照本身——后面的步骤用的是其中的 config,里面除了 thread_id,还多了 checkpoint_id。
在那个位置重新运行
增加 replay(app, config, node)。用 snapshot_before 找到坐标,用 app.invoke(None, snapshot.config) 接着运行,并返回 {"from": node, "result": 결과, "added": 늘어난 체크포인트 수}(占位符依次为结果与增加的检查点数量)。没有坐标时抛出 ValueError。
关键是把输入设为 None。给状态就是重新开始,None 的意思是“从保存的那一时刻接着运行”。added 是在回溯前后测量 history_len 再相减得到的值——重新运行了几个节点,就增加几个。
修改值,走向另一个分支
增加 fork(app, config, node, values)。先记下分叉之前线程的末尾,用 app.update_state(snapshot.config, values) 创建分支,再用那个 config 接着运行。请返回 {"forked_config": ..., "result": ..., "original": 분기 전 끝의 값 딕셔너리, "total": 지금 체크포인트 수}(占位符依次为分叉之前末尾的值字典与当前检查点数量)。
update_state 返回带有新 checkpoint_id 的 config。丢掉它而用原来的 config 接着运行,就不会分叉。而且分叉之后,只用 thread_id 去问得到的“当前”指向新分支的末尾,所以原来的结果要用分叉之前抓住的快照的 config 另外读取。
用大小来衡量
增加状态键 bulk 以及 state_bytes(values)、history_bytes(app, config)、weigh(topic, payload)。state_bytes 是 json.dumps(values, ensure_ascii=False, sort_keys=True, default=str) 得到的字符串的 UTF-8 长度。weigh 把同一个图运行两次:一次不带 bulk,一次在 bulk 中放入 payload,并返回 {"checkpoints": ..., "thin_total": ..., "fat_total": ..., "payload_bytes": ..., "carrying": ...}。
carrying 是把两个列表按相同位置比较,大小之差达到 payload_bytes 以上的检查点个数——也就是真正带着那个大值的位置的数量。运行两次时要分别使用账簿(调用两次 new_saver())。不要测量时间——它会因机器和负载而不同。字节数在同样的状态下总是一样。
重新建立账簿就会消失
增加 lost_demo(topic)。用新账簿创建图,把 run-1 线程运行一次并测量检查点数量,然后再用另一个新账簿创建一个图,用同一个 thread_id 去问。请返回 {"kept": ..., "after_restart": ..., "values_after_restart": ...}。
MemorySaver 放在进程内存里。新建一个 MemorySaver,是模拟进程重新启动最诚实的办法。向新账簿询问不会报错,而是返回空值——所以这个问题是悄无声息的。这就是它适合用于实验和测试,却不适合生产环境的原因。
用测得的值来记录
在 /root/work/agckpt/ckpt_report.json 中写入 topic、checkpoints_per_run、pending、base_score、fork_angle、fork_score、checkpoints_after_fork、payload_bytes、carrying、survives_restart,在 /root/work/agckpt/ckpt_report.md 中用 ## 무엇이 저장되는가(韩文,意为“保存了什么”)、## 되감기와 분기는 어떻게 다른가(韩文,意为“回溯与分叉有何不同”)、## 크기로 잰 것(韩文,意为“用大小测出的内容”)、## 사라지는 것(韩文,意为“会消失的内容”)四节来写。
数字不要手写,要用实际运行你的模块得到的值来填。topic 原样使用你用过的主题字符串,fork_angle 是不同于 기본(韩文,意为“基本”)的 angle 名称——评分器会用这两个值重新计算分数来对照。要在 write 之前分叉。payload_bytes 必须在 300 以上,差别才看得出来。survives_restart 是真或假。