Day 17 / 共 20 天 · 阶段5 进阶与收官

回调与流式:LangChain 的"神经系统"和打字机效果的真相

Day16 结尾埋了个钩子:with_listeners 凭什么能在链跑完时拿到完整记录?答案是贯穿一切的回调系统(callbacks/)——链开始、模型吐字、工具结束……每个环节都会向一批"听众"广播事件。今天拆三层:①听众长什么样(BaseCallbackHandler);②谁负责广播(CallbackManager + handle_event);③流式输出 stream/astream_events 是怎么在回调之上"长"出来的。学完它,Day19 的追踪系统就只剩临门一脚。

📍 你在 20 天里的位置(阶段5:进阶与收官 · D17-20)
S1 全景/LCEL S2 模型/消息 S3 数据/RAG S4 工具/Agent D17 回调/流式 D18 LCEL 高级 D19 追踪 D20 收官
💡 先用两个类比兜住今天 类比一:回调系统像体育比赛的现场直播体系。比赛(你的链)照常进行,场边架着一排机位(handlers):解说员(StdOutCallbackHandler,把过程念给你听)、记录员(LangChainTracer,Day19 见)、字幕机(流式 handler)。导播台(CallbackManager负责在每个节点喊"进球了!"(on_llm_start/on_llm_new_token…),所有机位各自反应;某台机器坏了也不影响比赛本身。类比二:astream_events 像给直播加了逐条弹幕的赛事数据流——不只给你最终比分(invoke 的返回值),而是把"谁在第几分钟做了什么"一条条推给你,前端想做多花哨的进度 UI 都有料。
L01

痛点:链是个黑箱,流式又从哪来

🤔 痛点一条 RAG 链跑了 8 秒才吐答案,慢在检索还是模型?中间 prompt 到底拼成了什么样?想在网页上做"逐字打出"的效果,可 invoke 要等全部生成完才返回。如果靠在业务代码里到处插 print 和计时器,链一改就全废。观测、日志、计费统计、流式 UI……这些需求都要求"看见链的内部过程",但又都不该侵入链本身的代码
💡 本质:观察者模式 + 事件总线LangChain 的答案是把执行过程事件化:每个组件在关键时刻(开始/出 token/结束/出错)调用 CallbackManager 广播事件,任何实现了 BaseCallbackHandler 的对象都可以订阅。链的代码只管"广播",不关心"谁在听、听了干什么"。日志、追踪、流式、甚至 Day16 的写历史,全是不同的听众——一套总线,喂饱所有横切需求
事件家族典型钩子谁触发
模型on_chat_model_start / on_llm_new_token / on_llm_end / on_llm_errorBaseChatModel(Day05)
on_chain_start / on_chain_end / on_chain_error各种 Runnable(Day03-04)
工具on_tool_start / on_tool_end / on_tool_errorBaseTool(Day13)
检索器on_retriever_start / on_retriever_endBaseRetriever(Day12)
大白话看这张表就明白:前 16 天学的每个组件,其实一直在偷偷"直播"自己——只是你没接机位而已。今天就是学怎么接。
L02

BaseCallbackHandler:一个"听众"长什么样

听众基类 BaseCallbackHandlerlibs/core/langchain_core/callbacks/base.py:496)——注意它是靠一堆 Mixin 拼出来的:

# libs/core/langchain_core/callbacks/base.py:496
class BaseCallbackHandler(
    LLMManagerMixin,        # 提供 on_llm_new_token / on_llm_end / on_llm_error
    ChainManagerMixin,      # 提供 on_chain_end / on_chain_error / on_agent_action...
    ToolManagerMixin,       # 提供 on_tool_end / on_tool_error
    RetrieverManagerMixin,  # 提供 on_retriever_end / on_retriever_error
    CallbackManagerMixin,   # 提供各种 on_*_start
    RunManagerMixin,        # 提供 on_text / on_retry / on_custom_event
):
    """Base callback handler."""

    raise_error: bool = False   # 听众自己抛异常时,要不要拖垮主流程?默认不要
    run_inline: bool = False    # 是否在主线程同步执行(追踪器需要顺序保证时置 True)

    @property
    def ignore_llm(self) -> bool:
        """Whether to ignore LLM callbacks."""
        return False            # 子类可声明"我不关心 LLM 事件",广播时会跳过它
一堆 Mixin每个 Mixin 对应一个事件家族(L01 的表)。所有钩子默认都是"空实现可覆写"——你写自定义 handler 时只覆写关心的那几个,比如做计费统计只写 on_llm_end 读 token 用量。
raise_error=False★关键设计:听众坏了不能砸场子。广播函数会捕获 handler 的异常只打日志(L03 见),保证"直播机位故障不中断比赛"。
ignore_llm 等开关性能优化:handler 声明不关心某类事件,handle_event 广播时直接跳过,连方法都不调。
AsyncCallbackHandler同文件下面还有异步版(callbacks/base.py:548),所有钩子是 async def——L06 的事件流 handler 就继承它。
📝 真实值:一个最小计时听众 class Timer(BaseCallbackHandler): 里只覆写两个钩子:on_chain_startself.t0 = time.time()on_chain_end 打印 f"耗时 {time.time()-self.t0:.2f}s"。使用:chain.invoke({"q": "hi"}, config={"callbacks": [Timer()]}) → 控制台输出 耗时 3.41s。链的代码一个字都没改——这就是"无侵入观测"。
L03

CallbackManager:导播台怎么广播

广播方 CallbackManagerlibs/core/langchain_core/callbacks/manager.py:1377)。看最常走的 on_chat_model_startmanager.py:1431):

# libs/core/langchain_core/callbacks/manager.py:1431
def on_chat_model_start(self, serialized, messages, run_id=None, **kwargs
) -> list[CallbackManagerForLLMRun]:
    managers = []
    for message_list in messages:                 # 每批消息一个独立的 run
        run_id_ = run_id or uuid7()               # ① 发一个 run_id(追踪树的节点编号)
        handle_event(                             # ② ★向所有 handler 广播"模型开始了"
            self.handlers, "on_chat_model_start", "ignore_chat_model",
            serialized, [message_list],
            run_id=run_id_, parent_run_id=self.parent_run_id,
            tags=self.tags, metadata=self.metadata, **kwargs,
        )
        managers.append(CallbackManagerForLLMRun( # ③ 返回"本次运行专属"的子管理器
            run_id=run_id_, handlers=self.handlers,
            inheritable_handlers=self.inheritable_handlers,
            parent_run_id=self.parent_run_id, ...))
    return managers

广播的底层是个通用函数 handle_eventmanager.py:285):

# libs/core/langchain_core/callbacks/manager.py:285
def handle_event(handlers, event_name, ignore_condition_name, *args, **kwargs):
    for handler in handlers:
        try:
            if ignore_condition_name is None or not getattr(
                handler, ignore_condition_name          # ① 尊重 ignore_chat_model 等开关
            ):
                event = getattr(handler, event_name)(*args, **kwargs)  # ② 反射调钩子
                if asyncio.iscoroutine(event):
                    coros.append(event)                 #    异步钩子攒起来统一跑
        except NotImplementedError as e:
            if event_name == "on_chat_model_start":     # ③ ★老 handler 不认识 chat 事件?
                handle_event([handler], "on_llm_start", # 降级成 on_llm_start 再喂一次
                             "ignore_llm", args[0], message_strings, ...)
run_id / parent_run_id★每次"开始"都发一个 uuid,并带上父 run 的 id——所有事件因此能拼成一棵运行树。Day19 的追踪就是把这棵树画出来。
返回子管理器CallbackManagerForLLMRun 是"绑定了本次 run_id 的遥控器"——模型后续吐 token 时调它的 on_llm_new_token,事件自动归属到正确的 run 上。这就是"父发号、子汇报"的结构。
getattr(handler, event_name)广播就是遍历听众、按名字反射调钩子。外层 try 捕获异常(raise_error=False 时只 log 警告)——印证 L02"机位坏了不砸比赛"。
NotImplementedError 降级贴心的兼容:老 handler 只实现过 on_llm_start(收字符串),遇到 chat 事件就把消息 get_buffer_string 拍平成字符串重新广播。生态包袱在源码里的样子。
💡 本质:manager 是"树上的信使"CallbackManager 自己不处理事件,它做两件事:广播(handle_event 遍历听众)和续树(造出携带 run_id 血缘的子 manager 传给下一层)。Day04 看到的 run_manager.get_child()、今天的 CallbackManagerForLLMRun,都是同一棵树在不同层的化身。
L04

configure:听众是怎么被装配上车的

你可能好奇:我只在 config 里塞了个 handler,谁把它装进 manager 的?每个组件执行前都会调 CallbackManager.configurelibs/core/langchain_core/callbacks/manager.py:1683),它委托给模块级的 _configuremanager.py:2390)。核心逻辑三步:

① 合并两路听众inheritable_callbacks(可继承:会随 get_child() 传给子 run,你在 config 里传的就是这类)+ local_callbacks(本层专属:只听这一个组件)。区分继承性,才能做到"给根链装一个 handler,全树事件都听得到"。
② 按环境变量加"官方机位"verbose=True → 自动挂 StdOutCallbackHandler(解说员);debug 模式 → ConsoleCallbackHandler开了 LangSmith 追踪 → 自动挂 LangChainTracermanager.py:2524-2532,Day19 主角,明天后天见)。
③ 挂上下文钩子register_configure_hook 注册的 ContextVar handler(如 usage 统计)也在这里被拾取。所以"with 一个上下文管理器就能全局监听"不是魔法。
大白话一句话:每次组件开跑前,都会现场组一个"导播台"——你手动塞的、环境变量声明的、上下文里挂的听众全被拢到一起。这就是为什么设个 LANGSMITH_TRACING=true 环境变量,什么代码都不改就有了全链路追踪。
⚠️ 坑:handler 传错位置,子组件听不到chain.invoke(x, config={"callbacks": [h]}) 传的是可继承听众,全树可见;而 chain.with_config(callbacks=[h]) 绑在特定 Runnable 上的行为不同层级会有差异。发现"我的 handler 只收到一半事件"时,先检查它是从哪条路上车的。
L05

stream:打字机效果从哪来

先看一个反直觉的事实——Runnable.stream 的默认实现(libs/core/langchain_core/runnables/base.py:1182):

# libs/core/langchain_core/runnables/base.py:1182
def stream(self, input, config=None, **kwargs) -> Iterator[Output]:
    """Default implementation of stream, which calls invoke.
    Subclasses must override this method if they support streaming output."""
    yield self.invoke(input, config, **kwargs)   # ★默认:憋完一整个结果,一次 yield
默认 = 假流式基类的 stream 只是把 invoke 的完整结果包成"只有一块的流"。真流式必须由子类覆写:ChatModel 覆写成逐 token 吐(Day05 见过 _stream),RunnableSequence 覆写成 transform 管道——上一步吐一块、下一步立刻加工一块,块在管道里"流"过去。
token 与回调的关系模型每收到一个 token,做两件事:yield 给调用方(打字机)、同时调 run_manager.on_llm_new_token(...) 广播给听众。同一个 token,流式管道和回调总线各发一份——前者给你,后者给观测系统。
断流点链中间夹了一个"不会 transform 的步骤"(比如吃整个输入才能算的自定义函数),流到它就断了——它上游的块会被攒齐再往下传。排查"为什么我的链不逐字出"时先找这种堵点。
📝 真实值:stream 一条链 chain = prompt | model | StrOutputParser(),执行 for chunk in chain.stream({"topic": "冰淇淋"}) → 依次打印 "为""什""么""冰淇淋"…… 每个 chunk 是 str(因为最后一环 parser 的 transform 把 AIMessageChunk 转成了字符串片段)。同一时刻,如果挂了 handler,它的 on_llm_new_token 也在被逐个调用。
💡 取舍:为什么默认实现是"假流式"而不是报错?好处是接口统一:任何 Runnable 都可以被 stream 调用,组合链时不用关心每个成员支不支持流式——不支持的成员自动退化为"一大块"。代价是退化是静默的:你以为在流式,其实在等全量。LangChain 选择了可组合性优先,把识别堵点的责任留给使用者。
L06

astream_events:把回调变成一条事件流

stream 只能拿到最终输出的碎片;想要"检索完成了、模型开始了"这类中间事件,用 astream_eventslibs/core/langchain_core/runnables/base.py:1356)。它的 v2 实现在 tracers/event_stream.py:1008——原理妙极了,就是"自己给自己装个听众":

# libs/core/langchain_core/tracers/event_stream.py:1008
async def _astream_events_implementation_v2(runnable, value, config=None, ...):
    event_streamer = _AstreamEventsCallbackHandler(   # ① ★造一个特殊听众(tracers/event_stream.py:101)
        include_names=include_names, include_types=include_types, ...)

    config = ensure_config(config)
    callbacks = config.get("callbacks")
    if callbacks is None:
        config["callbacks"] = [event_streamer]        # ② 把这个听众塞进 config
    elif isinstance(callbacks, list):
        config["callbacks"] = [*callbacks, event_streamer]

    async def consume_astream() -> None:
        try:      # ③ 后台任务:正常跑 astream,输出也 tap 进事件流
            async with aclosing(runnable.astream(value, config, **kwargs)) as stream:
                async for _ in event_streamer.tap_output_aiter(run_id, stream):
                    pass                              # 内容都被听众截走了
        finally:
            await event_streamer.send_stream.aclose()

    task = asyncio.create_task(consume_astream())     # ④ 链在后台跑
    async for event in event_streamer:                # ⑤ ★前台:从队列里逐条吐事件
        ...
        yield event
_AstreamEventsCallbackHandler一个"翻译官"听众:收到 on_chat_model_start 回调 → 翻译成 {"event": "on_chat_model_start", "name": ..., "run_id": ..., "data": ...} 字典,塞进内存队列(_MemoryStreamevent_stream.py:147-149)。
后台跑 + 前台读★生产者-消费者:链作为 asyncio task 在后台执行、事件不断入队;astream_events 的调用方在前台 async for 逐条消费。这就是"边跑边报"的实现。
include/exclude 过滤事件可能非常多,支持按 name/type/tags 过滤(_RootEventFilter)——比如只订阅 on_chat_model_stream 做打字机、忽略其余。
事件命名规律on_[组件类型]_(start|stream|end)on_chain_starton_chat_model_streamdata.chunkAIMessageChunk)、on_retriever_enddata.output 是文档列表)……docstring 里有完整对照表(runnables/base.py:1399-1420)。
回调总线:一次广播,多个听众;事件流 = 把回调装进队列 你的链在执行 retriever → prompt → model 每个节点触发 on_*_start/end/token CallbackManager handle_event 广播 StdOutCallbackHandler 解说员:打印过程 LangChainTracer 记录员:Day19 见 _AstreamEventsHandler 翻译官:事件入队 内存队列 → async for event in ... astream_events 的调用方逐条消费 听众互不干扰:加一个机位 不用改比赛(链)的任何代码
图注:广播是"推",事件流是在推的末端接了个队列变成"拉"。同一套回调总线,喂出了日志、追踪、事件流三种能力。
📝 真实值:astream_events 吐出的前几条事件 async for ev in chain.astream_events({"topic": "冰淇淋"}, version="v2"),依次收到:
{"event": "on_chain_start", "name": "RunnableSequence", "run_id": "0198...", "data": {"input": {"topic": "冰淇淋"}}}
{"event": "on_prompt_start", "name": "ChatPromptTemplate", ...} → ③ on_prompt_end(output 是拼好的 ChatPromptValue
{"event": "on_chat_model_start", ...} → ⑤⑥⑦… {"event": "on_chat_model_stream", "data": {"chunk": AIMessageChunk(content="为")}} 一大串
on_chat_model_end → ⑨ on_chain_end(data.output 是完整答案)。前端拿 ⑤ 做打字机、拿 ② ④ 做"正在思考/正在生成"状态灯。

👶 小白:stream 和 astream_events 我该用哪个?

👨‍🏫 老师:只要最终答案逐字出——用 stream/astream,简单直接。要中间过程(Agent 正在调哪个工具、检索到了什么、每一步进度)——用 astream_events,按事件名过滤你关心的。经验法则:聊天气泡用前者,Agent 执行面板用后者。另外注意 astream_events 只有异步版——因为它靠 asyncio task 实现"后台跑前台读"。

L07

串起来 + 今日小结

💡 回望 Day16,闭环了昨天的 with_listeners(on_end=_exit_history) 现在完全透明了:它就是往回调总线上挂了一个"只听 on_chain_end 的听众"(底层是 tracers 里的 RootListenersTracer),链跑完收到带完整输入输出的 Run → 写历史。记忆、日志、流式、追踪,四件看似无关的事,全是同一条总线上的不同听众。

🧠 今天你应该能回答

  • 回调系统解决什么问题?(无侵入地观察执行过程:观察者模式 + 事件总线)
  • handler 抛异常会砸掉主流程吗?(默认不会,raise_error=Falsehandle_event 捕获后只记日志)
  • run_id/parent_run_id 是干嘛的?(给每次运行发编号并记血缘,事件才能拼成运行树——Day19 的地基)
  • 为什么设个环境变量就有 LangSmith 追踪?(_configure 组装 manager 时按环境自动挂 LangChainTracer
  • stream 的默认实现是什么?(调 invoke 一次性 yield——假流式;真流式靠子类覆写 transform/_stream)
  • astream_events 的原理?(自己给自己挂一个把回调翻译成事件、塞进队列的听众;链后台跑、前台逐条消费)

✋ 10 分钟动手

cd /Users/bitmart/work/codes/github/AI_WORK/langchain/libs/core/langchain_core

# 1. 听众:Mixin 拼出来的钩子全家福
sed -n '496,546p' callbacks/base.py          # BaseCallbackHandler + ignore_* 开关
grep -n "def on_" callbacks/base.py | head -25

# 2. 导播台:广播与续树
sed -n '285,335p'  callbacks/manager.py      # handle_event(含降级兼容)
sed -n '1431,1483p' callbacks/manager.py     # on_chat_model_start
sed -n '2510,2545p' callbacks/manager.py     # _configure 自动挂 tracer

# 3. 流式:假流式默认值 + 事件流实现
sed -n '1182,1202p' runnables/base.py        # stream 默认=invoke
sed -n '1008,1075p' tracers/event_stream.py  # astream_events v2:自装听众+队列
明日预告 · Day 18:回调让链"可观测",但链本身还能更聪明——输入不齐怎么补默认值?按条件走不同分支?主模型挂了怎么自动换备胎?失败了怎么优雅重试?明天把 runnables/ 目录里剩下的六件高级积木(passthrough/branch/fallbacks/retry/configurable/router)一次配齐。
← Day 16 记忆与对话历史 Day 18 · LCEL 高级组合 →