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,保持向后兼容。
数据结构:run_id 索引 + seen 去重集 self.metadata(run_id → Meta) run-abc → ((), {node:"agent",step:1}) run-def → ((), {node:"agent",step:3}) start 时登记,end 时 pop self.seen(已吐 message id) { "msg-001", "msg-002" } 碎片吐出时 add 完整消息 dedupe 时 check stream((ns,"messages",(msg,meta))) 最终推进 D57 那条统一流
图注:metadata 查"谁在说",seen 保"不重复",一起喂给统一流。
控制流:一次模型调用的回调时间线 on_chat_model_start登记 run_id token"北" token"京" token"晴" 各自推流(不去重) on_llm_end完整消息 dedupe→跳 on_chain_endpop run_id + dedupe
图注: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 里那些错误类型的总览。
← Day 57 流式 5 模式 Day 59 · 运行时与部署 →