子图隔离与转换:schema 完全不同也能嵌套 + 流式下钻
昨天的共享靠"字段同名"。但真实世界里,你复用的子图往往有自己的一套 state schema,和父图一个字段都不重名。今天讲第二种接法——包一层转换函数做输入/输出映射,让父子图彻底解耦。再解决一个实战问题:流式运行时,怎么看到子图内部每一步(stream(subgraphs=True))?这背后是每个 Pregel 循环给自己的流式输出打的"命名空间戳"。
query、输出叫 docs;而你的父图里对应字段叫 user_input 和 retrieved。昨天的"同名共享"完全失效——名字都不一样。硬改子图的 schema?那子图就不通用了。怎么办?subgraph.invoke(...)、再把子图结果翻译回父图字段。子图完全不知道父图长啥样。类比:两个说不同语言的部门,中间放一个翻译(转换函数),各说各话、互不改口音。这是"适配器模式"在图上的体现。两种接法对照:共享 vs 隔离
| 维度 | 直接当节点(D49 共享) | 包一层转换(今天 隔离) |
|---|---|---|
| 写法 | add_node("sub", subgraph) | add_node("sub", call_sub)(call_sub 是函数) |
| state 关系 | 同名字段自动互通 | 手动翻译,schema 可完全无关 |
| 控制力 | 低(隐式) | 高(你决定传什么、收什么) |
| 适用 | 父子图本就共用一套状态 | 复用第三方/独立子图,字段命名不一致 |
包一层转换函数:手动 invoke 子图
核心就是"节点不是子图本身,而是一个调用子图的普通函数":
from langgraph.graph import StateGraph, START, END
from typing import TypedDict
# 子图:自己的一套 schema(query / docs)
class SubState(TypedDict):
query: str
docs: str
sub = StateGraph(SubState)
sub.add_node("search", lambda s: {"docs": f"docs-for({s['query']})"})
sub.add_edge(START, "search"); sub.add_edge("search", END)
subgraph = sub.compile()
# 父图:完全不同的 schema(user_input / retrieved)
class ParentState(TypedDict):
user_input: str
retrieved: str
def call_sub(state: ParentState) -> dict:
# ① 父 → 子:翻译输入字段
sub_out = subgraph.invoke({"query": state["user_input"]})
# ② 子 → 父:翻译输出字段
return {"retrieved": sub_out["docs"]}
parent = StateGraph(ParentState)
parent.add_node("sub", call_sub) # 节点是转换函数,不是子图
parent.add_edge(START, "sub"); parent.add_edge("sub", END)
app = parent.compile()
app.invoke({"user_input": "cats", "retrieved": ""})
# {'user_input': 'cats', 'retrieved': 'docs-for(cats)'}
subgraph.invoke({"query": ...})父图节点函数里显式调子图的 invoke,只传子图认识的字段 query。父图的 user_input 在这里被翻译成子图的 query。return {"retrieved": sub_out["docs"]}子图返回 docs,函数把它翻译回父图字段 retrieved。父子图字段名从头到尾不用一致。节点 = call_sub 函数对父图而言这就是一个普通函数节点,走 Day 03-05 的节点机制。子图被"藏"在函数体里。藏进函数里,引擎还认得它是子图吗?认得
你可能担心:子图藏在 call_sub 函数体里,引擎还能探测到吗?能——回看 D49 的 find_subgraph_pregel,它会挖函数闭包里引用的对象 pregel/_utils.py:63:
# pregel/_utils.py:63
elif isinstance(c, RunnableCallable):
if c.func is not None:
candidates.extend(
nl.__self__ if hasattr(nl, "__self__") else nl
for nl in get_function_nonlocals(c.func) # 挖函数引用到的外部对象
)
get_function_nonlocals(c.func)call_sub 函数体里用到了外层的 subgraph 变量(一个闭包自由变量/全局引用)。这个工具把这些被引用的对象挖出来。候选里发现 Pregel → 返回挖出的对象里有编译好的 subgraph(是 Pregel),于是识别成功。所以"包一层"的子图依然会出现在 get_subgraphs() 里、也能被流式下钻。find_subgraph_pregel 靠分析函数引用的外部变量。如果你在 call_sub 内部临时 compile 一张子图(sub = StateGraph(...).compile(); sub.invoke(...)),它是局部变量、不是被引用的 nonlocal,探测挖不到——于是 stream(subgraphs=True) 看不到它内部、get_state 也不展开它。想被下钻,就把子图 compile 在模块级/闭包外层,让节点函数去引用它。执行本身不受影响(invoke 照跑),受影响的只是"可观测性"。流式下钻的地基:每个循环给输出盖"命名空间戳"
子图内部跑的是它自己的一轮 Pregel 循环(PregelLoop)。这个循环启动时,把自己的 checkpoint 命名空间切成元组存起来 pregel/_loop.py:367:
# pregel/_loop.py:367
self.checkpoint_ns = (
tuple(cast(str, self.config[CONF][CONFIG_KEY_CHECKPOINT_NS]).split(NS_SEP))
if self.config[CONF].get(CONFIG_KEY_CHECKPOINT_NS)
else ()
)
然后每次它往外流式吐东西,都带上这个命名空间元组 pregel/_loop.py:1394:
# pregel/_loop.py:1394
for v in values(*args, **kwargs):
if mode in self.stream.modes:
self.stream((self.checkpoint_ns, mode, v)) # (命名空间, 模式, 数据)
checkpoint_ns 切成元组父图根循环的 ns 是 ()(空元组);子图循环的 ns 是 ("sub:<task_id>",);孙子图是 ("sub:x", "grand:y")。层级用元组长度表示。stream((ns, mode, v))每个流式块都是三元组:谁发的(ns)+ 什么模式 + 数据。父图和子图的输出都汇进同一个流,靠 ns 区分来源。那如果同一个父图节点里并发调了两次同一个子图,ns 岂不撞车?引擎用一个计数器给它们加数字后缀区分 pregel/_loop.py:329:
# pregel/_loop.py:329
if cnt := scratchpad.subgraph_counter(): # 本节点内第几次进子图
self.config = patch_configurable(self.config, {
CONFIG_KEY_CHECKPOINT_NS: NS_SEP.join(
(config[CONF][CONFIG_KEY_CHECKPOINT_NS], str(cnt)) # ns 末尾拼一个序号
)
})
subgraph_counter()同一节点内多次进入子图时递增。第一次不加后缀,之后拼上 |1、|2…让每次子图调用有唯一命名空间,存档和流式都不混。stream(subgraphs=True):放行子图输出
默认 stream() 只给你父图那层(ns=())的输出,子图内部悄悄跑完。传 subgraphs=True 就把带命名空间的子图输出也吐给你。看官方参数说明 pregel/main.py:2713:
# pregel/main.py:2713(stream 的 docstring)
# subgraphs: Whether to stream events from inside subgraphs, defaults to False.
# If True, the events will be emitted as tuples (namespace, data),
# or (namespace, mode, data) if stream_mode is a list,
# where namespace is a tuple with the path to the node where a subgraph is invoked,
# e.g. ("parent_node:", "child_node:").
# 默认:只见父图
for chunk in app.stream({"user_input": "cats", "retrieved": ""}):
print(chunk)
# {'sub': {'retrieved': 'docs-for(cats)'}} ← 只有父图节点的更新
# subgraphs=True:能下钻看到子图内部
for ns, chunk in app.stream({"user_input": "cats", "retrieved": ""}, subgraphs=True):
print(ns, chunk)
# () {'sub': {...}} ← 父图层
# ('sub:',) {'search': {'docs': ...}} ← 子图 search 节点那一步!
打开 subgraphs=True 后,你能看到子图里 search 节点执行的中间步骤,前缀是它所在的命名空间元组。这对调试嵌套系统至关重要。返回值多了 ns开了 subgraphs=True,每个 chunk 前面多一个命名空间元组,告诉你"这条是哪一层哪个节点发的"。ns 就是 L04 那个元组它正是子图循环打的戳。前端拿到后可以按 ns 做缩进展示、或过滤只看某一层。取舍、边界 + 今日小结
subgraphs=False 是刻意的:对大多数使用者,子图是一个封装好的黑盒,他们只关心"这个节点产出了什么",不关心子图内部十几步的碎碎念。默认扁平化(只看顶层)让输出简洁、心智负担低;需要调试时再显式 subgraphs=True 下钻。用"命名空间元组"给每条输出打标签、而非用多个独立流,好处是一条流承载所有层级,消费者一个循环就能处理、还能自由决定过滤到哪一层。这是"信息分层 + 按需展开"的设计。👶 用转换函数隔离后,子图还能有自己的 checkpointer 做断点续跑吗?
👨🏫 能。子图当节点(无论共享还是包一层)时,如果它编译时没显式关掉存档,会继承父图的 checkpointer,并在自己的命名空间下存档(D49 L05)。所以子图内部的 interrupt、时间旅行都正常工作,坐标是 父ns|子ns。只有显式 compile(checkpointer=False) 的子图才不参与持久化(也就不会被 find_subgraph_pregel 当作可下钻子图)。
🧠 今天你应该能回答
- 父子图 schema 完全不同怎么嵌套?(父图节点用普通函数,函数里手动翻译字段并 subgraph.invoke)
- 藏在函数里的子图还能被探测吗?(能,find_subgraph_pregel 会挖 get_function_nonlocals 引用的对象)
- 什么情况探测不到?(在节点函数内部临时 compile 的局部子图,非静态引用)
- 流式下钻的地基是什么?(每个 PregelLoop 把 checkpoint_ns 切成元组,输出都盖 (ns, mode, v) 戳)
- ns 元组的长度代表什么?(嵌套深度:() 父图、("sub:x",) 子图、两层则更长)
- stream(subgraphs=True) 做什么?(放行带命名空间的子图输出,chunk 前多一个 ns 元组)
✋ 10 分钟动手
# 1. 读命名空间戳与下钻源码
sed -n '360,372p' libs/langgraph/langgraph/pregel/_loop.py # ns 切成元组
sed -n '1388,1394p' libs/langgraph/langgraph/pregel/_loop.py # 输出打 ns 戳
# 2. 亲手下钻子图
python - <<'PY'
from langgraph.graph import StateGraph, START, END
from typing import TypedDict
class Sub(TypedDict):
query: str; docs: str
sub = StateGraph(Sub); sub.add_node("search", lambda s:{"docs":f"docs({s['query']})"})
sub.add_edge(START,"search"); sub.add_edge("search",END); subgraph = sub.compile()
class P(TypedDict):
user_input: str; retrieved: str
def call_sub(s): return {"retrieved": subgraph.invoke({"query": s["user_input"]})["docs"]}
p = StateGraph(P); p.add_node("sub", call_sub)
p.add_edge(START,"sub"); p.add_edge("sub",END); app = p.compile()
print("--- 默认 ---")
for c in app.stream({"user_input":"cats","retrieved":""}): print(c)
print("--- subgraphs=True ---")
for ns,c in app.stream({"user_input":"cats","retrieved":""}, subgraphs=True): print(ns, c)
PY
CachePolicy 怎么算缓存 key、TTL 怎么控制过期?明天读 types.py 的 CachePolicy 和 _internal/_cache.py 的默认 key 函数。