明明在跑,看起来却像卡住了
目标
把同一个图用五种方式流式输出,亲手数出什么在什么时候出现。只收集更新来重建最后的状态,并对照流式输出得到的最后状态与 invoke 的答案是否相同。
为什么重要
智能体需要流式输出的不只是令牌(token)。如果是最后只调用一次模型的图,即使加上令牌流式输出,前面的沉默依旧。用户需要的是“现在在哪一步”和“新确定了什么”,而这是图知道的,不是模型。
每种模式的单位不同。updates 是以节点为单位,所以并排运行的两个节点即使在同一个超级步里,事件也是分别出现;values 是以超级步为单位,所以两者的结果合并之后才出现一次。debug 还会告知超级步编号。不了解这个差别,就会把“为什么出现两次”误认为是 bug。
只收集更新来重建状态,必须了解 Reducer。图内部由 Reducer 做的事,在外部必须自己来做,如果把接续拼接的键覆盖,足迹就只剩一格。
而且不想留在状态里的进度提示,要用 custom 发出。放进状态,会被写进检查点,恢复时旧的提示会复活。
评分器不会相信你写下的说明。它会真正导入你的模块,每次用不同的会议纪要运行,并把事件的个数、顺序和形态与评分器另外数出的值对照。不测量时间。
步骤
- 在 /root/work/agstream/stream.py 中创建
State、四个节点(split、keypoints、actions、compose)、build_graph()、start(text)、run_values(text)。split之后,keypoints和actions并排运行,二者都汇聚到compose。 - 增加
run_updates(text),把stream_mode="updates"的事件原样作为列表返回。 - 增加
REDUCERS、rebuild(text)、rebuild_report(text),只收集更新来重建最后的状态,并对照它是否与values的最后一个相同。 - 增加
run_debug(text),让它返回{수퍼스텝번호: [노드이름...정렬됨]}(占位符依次为超级步编号、节点名称、已排序)。 - 增加
run_multi(text)和mode_sequence(text),处理stream_mode=["updates", "values"]的结果。 - 让
compose接收writer: StreamWriter,发出两个片段,并增加run_custom(text)。 - 增加
compare_invoke(text),一次性返回流式输出得到的最后状态是否与invoke的答案相同、最先到达的节点是什么、事件有几个。 - 把要测量的会议纪要保存到 /root/work/agstream/minutes.txt,用它测量,并记录到 /root/work/agstream/stream_report.json 和 /root/work/agstream/stream_report.md 中。
参考
- 执行契约:评分器会把
/root/work/agstream/stream.py当作 Python 模块导入,直接使用上面列出的名称。不会作为脚本运行。 - 状态键:
text、lines、points、todos、summary、trace。只有trace使用接续拼接的 Reducer,其余都是覆盖。节点名称和状态键不能重名——重名的话,编译时会出现ValueError: 'x' is already being used as a state key。 - 节点做的事:
split把text按行拆开,丢掉空行后放入lines。keypoints把包含"결정"(韩文,意为“决定”)的行放入points,actions把包含"하기로"(韩文,意为“决定要做”)的行放入todos。compose把"결정 N건 · 할 일 M건"(韩文,意为“决定 N 条 · 待办 M 条”)放入summary。四个节点都把自己的名称写入trace。 start(text)返回{"text": text, "trace": []}。run_values、run_updates、run_multi、run_custom要把stream(...)给出的内容不做加工,原样以列表返回。rebuild从start(text)开始,依次应用updates事件,返回得到的字典。REDUCERS是记录接续拼接的键的表,形如{"trace": "append"}。rebuild_report(text)的答案:{"same": 참거짓, "stream_order": [...], "merged_order": [...]}(占位符为布尔值)。stream_order是重建出的状态的trace,merged_order是values最后一个事件的trace。两者并不相同——为什么不同,请亲自查看后写下来。run_debug只看event["type"] == "task"的事件,收集event["step"]和event["payload"]["name"]。值要排序。mode_sequence(text)从run_multi的元组中按顺序只取出模式名称。compare_invoke(text)的答案:{"same": 참거짓, "first_node": 문자열, "values_events": 정수, "updates_events": 정수}(占位符依次为布尔值、字符串、整数、整数)。- 第 8 步的 /root/work/agstream/minutes.txt 是生成报告时用的会议纪要。评分器会读取该文件,用同样的输入重新测量,所以生成报告之后如果改了内容,数字就会对不上。
- 这个 Pod 没有互联网。langgraph 0.2.60 已经装好了。
- 官方文档:Streaming · Graph API overview · Types 参考
- 常见错误:以为
updates是以超级步为单位、重建时把接续拼接的键覆盖、把进度提示放进状态、把stream()的结果加工后返回(评分器看的是原本的形态)。
完整状态每一步都会出现
在 /root/work/agstream/stream.py 中创建 State、四个节点、build_graph()、start(text)、run_values(text)。split 之后 keypoints 和 actions 并排运行,二者都汇聚到 compose。
把 stream(입력, stream_mode="values")(占位符为输入)给出的内容原样作为列表返回即可。请亲眼确认最前面会出现一次输入状态。并排运行的两个节点的结果合并之后作为一个事件出现。
更新每个节点一个
增加 run_updates(text),把 stream_mode="updates" 的事件不做加工,原样作为列表返回。
一个事件是 {노드이름: 갱신}(占位符依次为节点名称与更新)。并排运行的两个节点即使在同一个超级步里,事件也会分别出现——所以事件数比 values 多。请亲自数一数出现了几个。
只靠更新重建状态
增加 REDUCERS = {"trace": "append"} 和 rebuild(text)。从 start(text) 开始依次应用 updates 事件,返回得到的字典。
图内部由 Reducer 做的事,在外部必须自己来做。全部覆盖的话,接续拼接的键只会剩下最后一格。但是,即使把 Reducer 模仿得很像,也不会完全相同——请把两个 trace 并排打印出来,亲自确认哪里不同、为什么不同。两者每次给出的答案都一样,但彼此不同。
第几个超级步里有哪些节点在一起运行
增加 run_debug(text),让它返回 {수퍼스텝번호: [노드이름...정렬됨]}(占位符依次为超级步编号、节点名称、已排序)。只看 event["type"] == "task" 的事件。
debug 模式会对每个节点给出 task 和 task_result 两个事件。task 一侧的 step 是超级步编号,payload["name"] 是节点名称。请确认并排运行的节点有相同的编号——这就是使用这个模式的理由。
同时监听多个模式
增加 run_multi(text) 和 mode_sequence(text)。给出 stream_mode=["updates", "values"],会出现 (모드이름, 값)(占位符依次为模式名称与值)元组。
以列表给出模式时,每个事件前面会附上它属于哪个模式。请亲眼看看两个模式的事件以什么顺序混合出现——不是一个模式全部给完再给另一个模式。mode_sequence 只取出名称,显示这个顺序。
不留在状态里,只展示
让 compose 接收 writer: StreamWriter 作为第二个参数,依次发出 {"stage": "compose", "points": 개수} 和 {"stage": "compose", "todos": 개수}(占位符均为个数),并增加 run_custom(text)。
把节点函数的第二个参数命名为 writer,并用 StreamWriter 标注,LangGraph 就会替你填上。用 writer(...) 放入的内容不会留在状态里,只会送到监听的一方——进度提示放进状态,会被写进检查点,恢复时还会复活。
流式输出的最后状态与一次性拿到的答案
增加 compare_invoke(text),让它返回 {"same": 참거짓, "first_node": 문자열, "values_events": 정수, "updates_events": 정수}(占位符依次为布尔值、字符串、整数、整数)。
如果 values 的最后一个事件与 invoke() 的答案相同,那么界面用流式输出绘制,记录另外用 invoke 留下,这种做法就没有必要运行两次。first_node 是 updates 第一个事件里的节点名称——它说明最先能展示给用户的是什么。
整理什么时候展示什么
把要测量的会议纪要保存到 /root/work/agstream/minutes.txt(决定事项和待办事项各至少一行,还要有一个空行)。用这份会议纪要测量,在 /root/work/agstream/stream_report.json 中写入 values_events、updates_events、first_node、supersteps、mode_sequence、custom_chunks、rebuild_same、stream_order、merged_order,在 /root/work/agstream/stream_report.md 中用 ## 다섯 가지 방식이 각각 무엇을 주나(韩文,意为“五种方式各给出什么”)、## 갱신만으로 상태를 다시 세우려면(韩文,意为“只靠更新重建状态时”)、## 화면에 먼저 보여 줄 수 있는 것(韩文,意为“最先能展示在界面上的内容”)、## 현장에서 무엇을 고르나(韩文,意为“在现场选哪一种”)四节来写。
全部是用下面的会议纪要运行得到的值。supersteps 放入把 run_debug 的答案改成字符串键之后的结果,rebuild_same、stream_order、merged_order 原样放入 rebuild_report 的答案。不要编造数字,要数出来再填。