Day 58 / 共 60 天 · 阶段 9 预制件与流式
流式底层:一个 token 从模型到你屏幕的旅程
昨天知道了 messages 模式靠"挂一个回调 handler"接入。今天把这个 handler——StreamMessagesHandler(_messages.py)拆开:模型每吐一个 token,它怎么捕获、怎么带上"这是哪个节点、哪个 LLM"的元数据、怎么防止"流式吐过的 token"和"节点返回的完整消息"被重复计数。这是 ChatGPT 式"逐字打字"体验在 LangGraph 里的最后一公里。
📍 阶段 9 · 预制件与流式(6 天)你在这里
D53 总览→
D54 逐行①→
D55 ToolNode→
D56 校验→
D57 流式5模式→
D58 流式底层
💡 用一个类比先兜住今天
想象模型是一台正在吐纸带的打字机,一个字一个字往外蹦。
StreamMessagesHandler 是装在打字机旁的传感器:模型开始打字时(on_chat_model_start)它记下"这是哪个部门在打"(元数据);每蹦一个字(on_llm_new_token)它立刻抄一份贴到公告栏(流);打完(on_llm_end)它核对一下别把整段又抄一遍(去重)。它不打断模型、只旁听——这就是"回调"的本质:搭个便车,模型该干嘛干嘛。L01
token 是怎么"冒"出来的
🤔 痛点模型内部是"一个 token 接一个 token"生成的。可
model.invoke() 只在全部生成完才返回一整条 AIMessage。LangGraph 怎么在不改模型代码的前提下,把中途每个 token 都截获下来送进流?💡 本质靠 LangChain 的回调(callback)机制。模型在生成过程中会主动"广播"事件:开始了(on_chat_model_start)、来了个新 token(on_llm_new_token)、结束了(on_llm_end)。任何注册了回调 handler 的对象都能旁听这些事件。
StreamMessagesHandler 就是这样一个旁听者——D57 里 run_manager.inheritable_handlers.append(...) 就是把它注册进去。它的类头 _messages.py:49:
# _messages.py:49
class StreamMessagesHandler(BaseCallbackHandler, _StreamingCallbackHandler):
"""A callback handler that implements stream_mode=messages.
Collects messages from:
(1) chat model stream events; and
(2) node outputs.
"""
run_inline = True
"""We want this callback to run in the main thread to avoid order/locking issues."""
继承 BaseCallbackHandler让它成为一个合法的 LangChain 回调,能接收 on_llm_* 等事件。收集两个来源注释点明:它既收模型 token 流(逐字),也收节点输出里的完整消息(比如节点直接返回一条消息)。两个来源 → 去重成了刚需(L05)。run_inline = True要求这个回调在主线程同步跑,不丢线程池。为什么?因为 token 是有顺序的,扔多线程会乱序、还要加锁。同步跑最省心。👶 一句话它是个"贴在模型旁的录音笔",模型说一个字它录一个字,全程不插嘴。
L02
Handler 的家当:stream + seen + metadata
构造函数 _messages.py:60,家当很少:
# _messages.py:60
def __init__(self, stream, subgraphs, *, parent_ns=None):
self.stream = stream # ← D57 传进来的 stream.put,往流里推数据
self.subgraphs = subgraphs # 是否也收子图的消息
self.metadata: dict[UUID, Meta] = {} # run_id → (命名空间, 元数据)
self.seen: set[int | str] = set() # 已经吐过的 message id(去重用)
self.parent_ns = parent_ns # 本 handler 创建时所在的命名空间
self.stream就是 D57 的 stream.put。handler 捕获到 token,最终就是调它 self.stream((ns,"messages",(msg,meta))) 推进流。self.metadata: run_id → Meta关键字典:每次模型调用有个唯一 run_id,这里记下它对应的命名空间 + 元数据(哪个节点、哪个 step)。token 来时靠 run_id 查回元数据。self.seen记录已经吐过的 message id。防止同一条消息被吐两次(L05)。parent_ns处理子图流式的边界:只有明确请求了子图消息时才放行(见 on_chat_model_start 的判断)。💡 本质:一个 handler 服务多个并发的模型调用Agent 里可能同时有多个模型调用(v2 并行工具后各自再想)。它们共用同一个 handler 实例,靠
run_id 这把钥匙在 self.metadata 里各认各的元数据、互不串味。用 run_id 做隔离键是整个 handler 能并发工作的地基。L03
on_chat_model_start:token 来之前先记好"户口"
模型一开始调用,先触发这个(_messages.py:130):
# _messages.py:130
def on_chat_model_start(self, serialized, messages, *, run_id, ..., tags=None, metadata=None, **kwargs):
if metadata and (not tags or (TAG_NOSTREAM not in tags)): # 没被标记"不要流"
ns = tuple(metadata["langgraph_checkpoint_ns"].split(NS_SEP))[:-1] # 算命名空间
if not self.subgraphs and len(ns) > 0 and ns != self.parent_ns:
return # 不收子图 → 跳过
if (filtered_tags := filter_to_user_tags(tags)) is not None:
metadata["tags"] = filtered_tags
self.metadata[run_id] = (ns, metadata) # ★ 登记 run_id → (ns, meta)
TAG_NOSTREAM not in tags边界:如果这次模型调用被打了 TAG_NOSTREAM(比如内部用的、不想流给用户的),直接不登记——后面 token 来了也无处可查、自然不吐。ns = ...split(NS_SEP)[:-1]从检查点命名空间算出这次调用属于哪个(子)图。[:-1] 去掉最后一段(本节点的 id)。not self.subgraphs and len(ns)>0 → return如果没要子图消息、而这次调用来自子图,跳过——不登记就不会吐。self.metadata[run_id] = (ns, metadata)核心动作:登记"户口"。之后 token/结束事件都靠 run_id 查这份 (ns, metadata)。没登记的 run_id,后续事件一律无视。💡 设计取舍①:为什么用"登记制"而不是"来了 token 现算元数据"?token 事件(on_llm_new_token)触发得极其频繁——一条回答几百个 token 就是几百次回调。如果每个 token 来时都去解析 tags、算命名空间、判断要不要流,那是几百次重复计算。改成在 start 时算一次、登记进字典,token 来时只做一次
self.metadata.get(run_id) 字典查询。把"是否该流、属于谁"这个不变判断前移到 start,让高频的 token 路径极致轻量。这和 D55 的"__init__ 预扫注入参数"是同一种智慧。L04
on_llm_new_token:每个 token 立刻推进流
今天的绝对核心——每来一个 token 触发(_messages.py:151):
# _messages.py:151
def on_llm_new_token(self, token, *, chunk=None, run_id, ..., **kwargs):
if not isinstance(chunk, ChatGenerationChunk):
return # 不是聊天模型的 chunk → 忽略
if meta := self.metadata.get(run_id): # ★ 查户口:登记过才处理
self._emit(meta, chunk.message) # 把这一小片消息推进流
推流动作 _emit(_messages.py:97):
# _messages.py:97
def _emit(self, meta, message, *, dedupe=False):
if dedupe and message.id in self.seen:
return # 去重(L05)
else:
if message.id is None:
message.id = str(uuid4()) # 没 id 就补一个
self.seen.add(message.id)
self.stream((meta[0], "messages", (message, meta[1]))) # ★ (ns, "messages", (消息, 元数据))
chunk.message每个 token 其实是一个 AIMessageChunk(消息碎片),只含这一小段增量内容。meta := self.metadata.get(run_id)查 L03 登记的户口。没登记(返回 None)→ 整个 if 不进,这个 token 被静默丢弃。这就是"NOSTREAM/子图不流"的落地点。self._emit(meta, chunk.message)注意这里 dedupe 默认 False——流式 token 每片都不同、都要吐,不去重。self.stream((meta[0], "messages", (message, meta[1])))最终推流:三元组正是 D57 学的 (ns, mode, data),其中 data 是 (消息碎片, 元数据)。你 for msg, meta in stream(...) 接住的就是它。📝 真实值模型回答"北京晴",你会依次收到:
(AIMessageChunk("北"), meta)、(AIMessageChunk("京"), meta)、(AIMessageChunk("晴"), meta)——三次回调、三次推流。前端把 content 拼起来就是逐字出现的效果。meta 里有 langgraph_node="agent" 等,能知道是哪个节点在说话。L05
seen 去重:为什么会重复、怎么防
🤔 痛点token 流吐了"北""京""晴"三片。模型结束时(on_llm_end)又拿到一条完整的 AIMessage("北京晴")。节点返回时这条完整消息还会再出现一次。如果都吐,用户就会看到内容出现两三遍。怎么办?
靠 seen 集合 + dedupe=True。看 on_llm_end(_messages.py:166)和节点结束 on_chain_end(_messages.py:223):
# _messages.py:166 模型结束:完整消息 dedupe=True
def on_llm_end(self, response, *, run_id, **kwargs):
if meta := self.metadata.get(run_id):
if response.generations and response.generations[0]:
gen = response.generations[0][0]
if isinstance(gen, ChatGeneration):
self._emit(meta, gen.message, dedupe=True) # ★ 去重吐完整消息
self.metadata.pop(run_id, None) # 用完清户口
# _messages.py:223 节点结束:节点输出里的消息也 dedupe
def on_chain_end(self, response, *, run_id, **kwargs):
if meta := self.metadata.pop(run_id, None):
... self._find_and_emit_messages(meta, response) # 内部 _emit(dedupe=True)
流式 token:dedupe=FalseL04 里每片碎片都吐、不查 seen——因为碎片本就该逐个出现。但每片吐出时会 seen.add(message.id)。关键:同一次生成的所有碎片共享同一个 message.id。on_llm_end:dedupe=True模型结束给的完整消息 id 和碎片相同 → 已在 seen 里 → _emit 开头 if dedupe and id in seen: return 直接跳过。不会把完整消息再吐一遍。on_chain_end:也 dedupe节点返回时那条消息同样 id 已 seen → 跳过。这就是"两个来源"(L01)却不重复的秘密。metadata.pop(run_id)模型/节点结束时清掉户口,防止字典无限增长(内存泄漏防护)。💡 本质:id 一致 + 先到先得去重的地基是"流式碎片和最终完整消息共享同一个 message.id"。流式碎片先到、逐个吐并把 id 记进 seen;完整消息后到、发现 id 已 seen 就闭嘴。谁先到谁负责显示,后到的同 id 一律沉默——用一个 id 集合优雅解决了"同一内容多来源"的重复难题。若模型不支持流式(没碎片),那 on_llm_end 的完整消息就是第一次见到、正常吐出——两种情况都对。
L06
_find_and_emit_messages:从节点输出里挖消息
节点可能返回各种形状:一条消息、消息列表、或 {"messages":[...]} 状态。_find_and_emit_messages(_messages.py:106)负责挖出所有消息:
# _messages.py:106
def _find_and_emit_messages(self, meta, response):
if isinstance(response, BaseMessage):
self._emit(meta, response, dedupe=True) # 直接是一条消息
elif isinstance(response, Sequence):
for value in response:
if isinstance(value, BaseMessage):
self._emit(meta, value, dedupe=True) # 消息列表
else:
for value in _state_values(response): # 状态 dict:遍历各字段值
if isinstance(value, BaseMessage):
self._emit(meta, value, dedupe=True)
elif isinstance(value, Sequence):
for item in value:
if isinstance(item, BaseMessage):
self._emit(meta, item, dedupe=True) # 字段值是消息列表
三种形状全覆盖消息 / 消息列表 / 状态 dict。_state_values(:38)把 dict 或 Pydantic 状态的字段值都取出来遍历。全部 dedupe=True节点输出里的消息都去重——因为它们大概率已经在流式阶段吐过了(同 id 已 seen)。on_chain_end 里处理 Command补充:on_chain_end(:223)还专门处理节点返回 Command(D15)的情况,从 command.update 里挖消息。覆盖工具用 Command 返回结果的场景。⚠️ 边界:节点直接返回消息、但模型没流式,怎么保证不漏?如果某节点不调模型、直接往 messages 里塞一条消息(比如硬编码回复),那就没有 on_llm_new_token 事件——只有 on_chain_end。此时这条消息第一次进
_find_and_emit_messages,id 不在 seen 里,dedupe=True 也会正常吐出(因为"没吐过")。反过来,on_chain_start(:191)还会预先把输入里已有消息的 id 记进 seen——防止把"上一轮就存在的历史消息"当新消息重吐。seen 既防重复吐,也防把老消息误当新消息,两头堵。L07
v2 事件流 on_stream_event + 阶段收官
新版模型支持"内容块(content-block)"事件(如推理块、工具调用块分开流)。StreamMessagesHandlerV2(_messages.py:259)用 on_stream_event(_messages.py:373)转发:
# _messages.py:373
def on_stream_event(self, event, *, run_id, ..., **kwargs):
if meta := self.metadata.get(run_id):
if event.get("event") == "message-start":
self._streamed_run_ids.add(run_id)
msg_id = event.get("message_id")
if msg_id:
self.seen.add(msg_id) # 提前记 id,让 on_chain_end 去重
v2_meta = {**meta[1], "run_id": str(run_id)}
self.stream((meta[0], "messages", (event, v2_meta))) # 把事件本身推进流
v2 vs v1v1(on_llm_new_token)吐 AIMessageChunk;v2 吐结构化事件(message-start / content-block-* / message-finish)。v2 能区分"这是推理内容还是最终答案"。message-start 时 seen.add(msg_id)同样的去重思路:v2 事件流一开始就把 message id 记进 seen,防止节点结束时把完整消息重吐。on_llm_new_token 被 v2 置空(:275)v2 handler 故意把 v1 的 token 回调设成 no-op——因为 v2 走事件流,不能让 v1 碎片也漏进来污染。何时用 v2只在内部 CONFIG_KEY_STREAM_MESSAGES_V2 开关打开时才挂 v2(D57 main.py:2810 的 use_stream_messages_v2);直接 stream_mode="messages" 仍用 v1,保持向后兼容。图注:metadata 查"谁在说",seen 保"不重复",一起喂给统一流。
图注:start 登记 → 每个 token 立刻吐 → end 时完整消息因去重被跳过。
🎓 阶段 9 收官六天我们从
create_react_agent 一行函数出发(D53),钻进它建的两个节点(D54 模型节点、D55 ToolNode),看了路由与校验的辅助件(D56),再从流式五模式(D57)一路追到单个 token 的旅程(D58)。核心感受:预制件不是魔法,是前 52 天底盘知识的精致组装;流式不是附加功能,是引擎的原生工作方式。你现在既能用现成 Agent,也看得懂它每一层在做什么。🧠 今天你应该能回答
- token 是怎么被截获的?(LangChain 回调机制,handler 旁听 on_llm_new_token)
- 为什么 handler 要 run_inline=True?(token 有序,同步跑避免乱序和加锁)
- on_chat_model_start 干嘛?(按 run_id 登记 (ns, metadata),高频 token 路径只需查字典)
- 流式 token 为什么 dedupe=False,完整消息 dedupe=True?(碎片都要吐;完整消息同 id 已 seen 要跳过)
- 去重的地基是什么?(碎片和完整消息共享同一 message.id,先到先得)
- 节点直接返回消息(没流式)会漏吐吗?(不会,第一次见 id 不在 seen 照样吐)
- v2 handler 和 v1 差别?(v2 吐结构化 content-block 事件,且把 v1 token 回调置空)
✋ 10 分钟动手
# 1. 核心回调逐行读
sed -n '97,165p' libs/langgraph/langgraph/pregel/_messages.py # _emit + start + token
sed -n '166,257p' libs/langgraph/langgraph/pregel/_messages.py # end + 去重 + 节点消息
# 2. 亲手看逐 token(需真实模型 key)
python - <<'PY'
from langgraph.prebuilt import create_react_agent
agent = create_react_agent("openai:gpt-4o-mini", [])
for msg, meta in agent.stream(
{"messages":[("user","用一句话介绍北京")]},
stream_mode="messages"):
print(repr(msg.content), "| node=", meta.get("langgraph_node")) # 一片一片打印
PY
明天预告 · Day 59:进入收官阶段——运行时与生产。看
runtime.py 的 Runtime/context(节点怎么拿到 store、stream_writer)、config.py、LangGraph Platform 与 CLI(libs/cli)、以及 errors.py 里那些错误类型的总览。