学习笔记 · Obsidian
Streaming、Human-in-the-loop、Subgraph 与前端
两层 streaming API
LangGraph 当前有两层流式接口:
| 层 | 推荐场景 | 主要接口 |
|---|---|---|
| Event streaming | 新应用的业务层消费 | stream_events / astream_events,version=v3 |
| Stream-mode API | 需要 Pregel 低层事件和兼容既有代码 | stream / astream,version=v2 |
Event streaming 建立在低层 stream modes 之上:Pregel 先产生 raw events,event router 再把事件送入 transformers,最终暴露 messages、values、subgraphs、output、interrupts 等类型化 projection。
不要把两者当作互斥协议。业务 UI 通常消费 event streaming;调试、定制运行时事件或旧客户端可能直接消费 stream modes。
Event streaming:推荐的应用层模型
一个 run,多种独立 projection
run stream 暴露:
| projection | 用途 |
|---|---|
| stream 本身 | 遍历全部 protocol events |
| stream.messages | 模型消息、文本 token、reasoning 与 tool call delta |
| stream.values | 每步完整 State;可等待最终值 |
| stream.output | 等待最终输出 |
| stream.subgraphs | 发现和观察嵌套图 |
| stream.interrupts | 当前暂停请求的 payload |
| stream.interrupted | 本次 run 是否暂停 |
| stream.extensions | 自定义 transformer 的 projection |
多个消费者可以并发读取同一 run 的不同 projection。消费 messages 不会把 values 或 output 的事件“抢走”。异步代码可用并发任务分别消费;同步代码可用 interleave 保持严格到达顺序。
Messages
messages channel 使用 content block 生命周期:
- message-start;
- content-block-start;
- 零个或多个 content-block-delta;
- content-block-finish;
- message-finish。
这使文本、reasoning、tool call 与多模态内容不依赖具体供应商格式。message.text 可逐 token 遍历,也可在完成后转成完整字符串;message.reasoning 与 message.tool_calls 是独立 projection。
如果必须严格保留文本、reasoning 和 tool call chunk 的原始交错顺序,应读取 raw message events,而不是分别消费三个 projection。
Subgraphs
stream.subgraphs 直接提供嵌套执行的 path、graph_name、messages 与 values,不需要手动解析 namespace 字符串。生命周期事件的 cause 可以把子图或子 agent 关联回发起它的工具调用或边。
Interrupt 恢复
一次 run 完成后:
- 检查 stream.interrupted;
- 从 stream.interrupts 读取暂停 payload;
- 收集外部输入;
- 用相同 thread_id 再次 stream_events,输入 Command(resume=...);
- 重复,直到 interrupted 为 false;
- 从 stream.output 获取最终 State。
使用 interrupt 必须配置 checkpointer 和 thread_id。
Raw protocol channels
| channel | 事件 |
|---|---|
| values | 完整 State |
| updates | 节点增量 |
| messages | content block 消息流 |
| tools | tool-started、delta、finished、error |
| lifecycle | started、running、completed、failed、interrupted |
| checkpoints | branching 与 time travel 所需轻量信息 |
| input | HITL 请求与响应 |
| tasks | Pregel task 创建和结果 |
| custom | 节点自定义数据 |
| custom:name | 命名 transformer 输出 |
raw event 的 namespace 是从根图到当前 scope 的路径。每段由稳定名字与本次运行 ID 组成;如只关注某个子树,可按 namespace 过滤。
Stream transformer 与自定义 projection
Transformer 是只观察事件的投影层,不回调 graph runtime。每个 transformer 可:
- init:创建 projection;
- process:查看、修改或抑制事件;
- finalize:成功结束时完成 projection;
- fail:错误结束时传播失败。
required_stream_modes 声明它需要 Pregel 发出哪些底层模式。未声明的 mode 不会由图生成;声明 mode 只决定上游发出,不会自动过滤 process 收到的事件。
StreamChannel
- 命名 StreamChannel:projection 出现在 stream.extensions,同时每次 push 进入主 raw stream,事件名为 custom:name;payload 必须可序列化。
- 未命名 StreamChannel:只在进程内 side channel 暴露,可以承载 promise、async iterable 或类实例,不进入 wire protocol。
框架负责 channel 的 close 与 fail,transformer 只负责 push。
适合自定义 channel 的信号
- PII 脱敏计数与合规命中;
- 长任务进度、阶段与百分比;
- token、延迟和费用累计;
- 检索来源与引用;
- 不应写进聊天文本的领域事件。
Transformer 可以在事件到浏览器前做统一脱敏,但安全设计仍应默认源头少收集、服务端最小化输出,不能把前端 transformer 当作唯一保护。
低层 Stream-mode API v2
统一 StreamPart
LangGraph 1.1+ 使用 version=v2 时,每个 chunk 都有一致结构:
- type:mode 类型;
- ns:父子图 namespace;
- data:该 mode 的 payload。
相比 v1:
- 单 mode 不再返回裸 dict;
- 多 mode 不再返回二元组;
- subgraph 不再改变 tuple 形状;
- 根据 type 可做静态类型收窄。
invoke(version=v2) 返回 GraphOutput,主要读取 value 与 interrupts。旧式字典访问为兼容保留但已 deprecated。
stream modes
| mode | 输出 |
|---|---|
| values | 每步完整 State |
| updates | 每个节点的 State 增量 |
| messages | LLM token + metadata |
| custom | get_stream_writer 发出的业务数据 |
| checkpoints | checkpoint event,需要 checkpointer |
| tasks | task 开始、完成、错误,需要 checkpointer |
| debug | checkpoints、tasks 与更多元数据 |
updates 更省带宽;values 方便快照式 UI,但可能包含 private channel。debug 数据量最大,只用于诊断。
Token 过滤
messages mode 的 metadata 可按:
- LLM tag;
- langgraph_node;
- subgraph namespace;
过滤。给某次模型调用加 nostream tag 可以继续运行但不发 token,适合内部结构化输出或避免重复显示。
Custom data
节点或工具通过 get_stream_writer 写自定义数据,调用方包含 custom mode 才能收到。非 LangChain 模型也可把其原生流映射为 custom。
Python 3.11 以下的 async context 传播有限:
- 需要显式把 RunnableConfig 传给异步模型调用;
- 不能可靠使用 get_stream_writer,应显式注入 writer。
当前项目若已使用 Python 3.11+,仍应在库文档中保留这条兼容边界。
Subgraph stream
低层 API 需要 subgraphs=True 才会把子图事件发到父 stream。即使内层 create_agent 自己能流式输出,把它作为父图节点后如果未开启 subgraphs,父 messages mode 也看不到内层模型 token。
Interrupt 的精确语义
生命周期
节点调用 interrupt(payload) 时:
- runtime 抛出内部控制异常;
- 当前 graph state 由 checkpointer 保存;
- JSON 可序列化 payload 暴露给调用者;
- run 可无限期等待;
- 调用者以相同 thread_id + Command(resume=value) 恢复;
- value 成为节点内 interrupt 的返回值。
thread_id 是持久游标。换一个 thread_id 会启动新状态,无法恢复原暂停点。
恢复会从节点开头重跑
它不会从 interrupt 那一行继续。节点从头执行,interrupt 之前的所有代码再次运行。由此得到四条硬规则:
- interrupt 之前的副作用必须幂等;
- 更安全的做法是把副作用放在 interrupt 之后或独立节点;
- 不要在 interrupt 前创建不可查重的新记录;
- 不要依赖局部变量保留执行现场,所需数据写入 State。
不要用 try/except 包住 interrupt
interrupt 通过特殊异常向 runtime 冒泡。裸 try/except 会截获它,使图无法正确暂停。应把可能失败的业务代码与 interrupt 分开,或只捕获明确异常类型。
多 interrupt 的顺序
同一 task 内多个 interrupt 的 resume value 按索引匹配。上线后不能在恢复点之前:
- 重排 interrupt;
- 条件跳过某个 interrupt;
- 引入非确定性循环改变调用次数。
并行分支同时 interrupt 时,恢复输入应按 interrupt ID 映射各自值,避免把答案配给错误分支。
输入校验的正确模式
不要在一个节点内使用 while True + interrupt。每次恢复都从节点开头重放,循环会不断重复历史迭代。
正确模式:
- State 保存 pending_question;
- 节点每次只调用一次 interrupt;
- 无效答案更新 pending_question;
- conditional edge 路由回同一节点;
- 有效答案进入下一节点。
典型 HITL
- 批准或拒绝外部动作;
- 编辑模型输出或工具参数;
- 工具函数内部审批;
- 多字段表单和逐步澄清;
- 高风险 SQL、支付、邮件发送前人工确认。
静态 interrupt_before / interrupt_after 更像调试 breakpoint,不推荐作为业务 HITL。
Subgraph
Subgraph 是作为父图 node 使用的已编译 graph,适合:
- 多 agent;
- 重用一组节点;
- 多团队以稳定 input/output schema 并行开发;
- 把复杂流程封装成明确模块。
父子图通信
| 模式 | 适用条件 | 实现 |
|---|---|---|
| 在父 node 内 invoke 子图 | State schema 不同或需要转换 | wrapper 映射父 State → 子输入 → 父更新 |
| 直接把 compiled subgraph 加为 node | 共享 State key | 直接 add_node,无 wrapper |
不同 schema 的 wrapper 适合为每个 subagent 保留私有消息;共享 messages 等 channel 时,直接 subgraph node 更简单。
persistence 三种模式
| 模式 | compile 参数 | 跨调用记忆 | interrupt | 并行同一子图 |
|---|---|---|---|---|
| per-invocation | checkpointer=None,默认 | 无 | 有 | 支持 |
| per-thread | checkpointer=True | 有 | 有 | 不支持并行写同一 namespace |
| stateless | checkpointer=False | 无 | 无 | 支持但无 durable execution |
Per-invocation
每次调用从新 State 开始,但本次调用内继承父 checkpointer,因此可以 interrupt、恢复和容错。多数一次性 subagent tool 应使用默认模式。
Per-thread
同一 thread 多次调用会累积子图 State,适合持续研究或编码 assistant。代价是:
- 同一个 per-thread subgraph 不能并行调用,否则 checkpoint namespace 冲突;
- 要在模型层禁用并行工具调用,或加调用限制;
- 多个不同子图必须有稳定、唯一 namespace;
- 在父 node 内按调用顺序 invoke 多个 per-thread 子图,重排代码可能错配历史。
把不同子 agent 包进具有唯一 node name 的 StateGraph,可以获得稳定 namespace。直接作为父图节点的 subgraph 已自动获得按名字隔离的 namespace。
Stateless
像普通函数一样运行,减少 checkpoint 开销,但不能 pause/resume,进程崩溃后只能从头执行。
State inspection
get_state(config, subgraphs=True) 可读内部 State,但要求 runtime 能静态发现 subgraph:
- 作为 node 添加;
- 或在可识别父 node 内调用。
在 tool 函数深处动态调用的 subgraph 通常不能被静态检查,但 interrupt 仍能向顶层冒泡。
Subgraph time travel
默认 per-invocation 子图在父图只表现为一个 super-step;要从子图内部 checkpoint travel,需要 checkpointer=True 并使用内部 config。
前端:把图结构变成产品 UX
LangGraph 前端不是只显示一条 assistant 消息,而是可以直接映射运行时概念:
| runtime | UI |
|---|---|
| named nodes | 卡片、步骤、状态 badge |
| State keys | 分类、来源、分析、最终结论区域 |
| streaming metadata | 把 token 路由到产生它的节点 |
| checkpoints | 历史查看、恢复与审计 |
| interrupts | 审批、修改与补充输入 |
| subgraphs | 按需展开嵌套执行 |
useStream 与节点发现
前端 SDK 的 useStream 暴露:
- stream.subgraphs:当前 thread 已观察到的节点;
- useMessages(stream, node):该节点范围内的消息;
- stream.values:完整 graph state;
- node.status:pending、running、complete、error。
UI 应从 stream.subgraphs 动态发现节点,而不是写死固定管线。条件分支跳过的节点不会出现;可以只渲染实际节点,或把预期但未出现的节点显示为 dim 状态。
节点卡片
推荐:
- 一张卡对应一个 node;
- scoped messages 展示流式与最终内容;
- 只有确需业务字段时才读 stream.values;
- 完成节点自动折叠,当前节点展开;
- 单节点错误显示在对应卡片,不直接抹掉已完成分支;
- markdown renderer 能处理未闭合的流式语法;
- 显示总体步骤和合理的历史耗时预估。
不要假定 node name 与 State key 同名。节点消息用 namespace scoped selector,最终汇总字段再显式读取真实 State key。
自定义 channel 前端选择器
- useExtension(stream, name):返回该 custom channel 最新的、已解包 payload;适合进度、计数、状态 badge;
- useChannel(stream, full-channel-id):返回有界 raw event buffer;适合事件日志、审计或无高层 selector 的 channel。
useChannel 要配置 bufferSize 和 replay,避免无限内存增长。常见做法是同一 channel:
- useExtension 驱动当前摘要;
- useChannel 驱动滚动历史。
React、Vue、Svelte、Angular 的返回值遵循各自响应式模型,初始化前可能是 undefined。
安全与可靠性检查
- values stream 是否泄漏 private State?
- custom payload 是否可序列化并经过租户过滤?
- PII 是否在到达浏览器前脱敏?
- stream buffer 是否有上限与断线重连策略?
- interrupt payload 是否只含必要、JSON 安全数据?
- 恢复是否校验调用者有权访问该 thread?
- interrupt 前副作用是否幂等?
- 同节点多个 interrupt 的顺序是否会被版本升级改变?
- per-thread subgraph 是否禁止并行同实例调用?
- UI 是否只把状态展示给授权用户,而不是因“可观察”就默认公开?