回调与流式:LangChain 的"神经系统"和打字机效果的真相
Day16 结尾埋了个钩子:with_listeners 凭什么能在链跑完时拿到完整记录?答案是贯穿一切的回调系统(callbacks/)——链开始、模型吐字、工具结束……每个环节都会向一批"听众"广播事件。今天拆三层:①听众长什么样(BaseCallbackHandler);②谁负责广播(CallbackManager + handle_event);③流式输出 stream/astream_events 是怎么在回调之上"长"出来的。学完它,Day19 的追踪系统就只剩临门一脚。
StdOutCallbackHandler,把过程念给你听)、记录员(LangChainTracer,Day19 见)、字幕机(流式 handler)。导播台(CallbackManager)负责在每个节点喊"进球了!"(on_llm_start/on_llm_new_token…),所有机位各自反应;某台机器坏了也不影响比赛本身。类比二:astream_events 像给直播加了逐条弹幕的赛事数据流——不只给你最终比分(invoke 的返回值),而是把"谁在第几分钟做了什么"一条条推给你,前端想做多花哨的进度 UI 都有料。痛点:链是个黑箱,流式又从哪来
invoke 要等全部生成完才返回。如果靠在业务代码里到处插 print 和计时器,链一改就全废。观测、日志、计费统计、流式 UI……这些需求都要求"看见链的内部过程",但又都不该侵入链本身的代码。CallbackManager 广播事件,任何实现了 BaseCallbackHandler 的对象都可以订阅。链的代码只管"广播",不关心"谁在听、听了干什么"。日志、追踪、流式、甚至 Day16 的写历史,全是不同的听众——一套总线,喂饱所有横切需求。| 事件家族 | 典型钩子 | 谁触发 |
|---|---|---|
| 模型 | on_chat_model_start / on_llm_new_token / on_llm_end / on_llm_error | BaseChatModel(Day05) |
| 链 | on_chain_start / on_chain_end / on_chain_error | 各种 Runnable(Day03-04) |
| 工具 | on_tool_start / on_tool_end / on_tool_error | BaseTool(Day13) |
| 检索器 | on_retriever_start / on_retriever_end | BaseRetriever(Day12) |
BaseCallbackHandler:一个"听众"长什么样
听众基类 BaseCallbackHandler(libs/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_start 记 self.t0 = time.time(),on_chain_end 打印 f"耗时 {time.time()-self.t0:.2f}s"。使用:chain.invoke({"q": "hi"}, config={"callbacks": [Timer()]}) → 控制台输出 耗时 3.41s。链的代码一个字都没改——这就是"无侵入观测"。CallbackManager:导播台怎么广播
广播方 CallbackManager(libs/core/langchain_core/callbacks/manager.py:1377)。看最常走的 on_chat_model_start(manager.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_event(manager.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 拍平成字符串重新广播。生态包袱在源码里的样子。CallbackManager 自己不处理事件,它做两件事:广播(handle_event 遍历听众)和续树(造出携带 run_id 血缘的子 manager 传给下一层)。Day04 看到的 run_manager.get_child()、今天的 CallbackManagerForLLMRun,都是同一棵树在不同层的化身。configure:听众是怎么被装配上车的
你可能好奇:我只在 config 里塞了个 handler,谁把它装进 manager 的?每个组件执行前都会调 CallbackManager.configure(libs/core/langchain_core/callbacks/manager.py:1683),它委托给模块级的 _configure(manager.py:2390)。核心逻辑三步:
① 合并两路听众inheritable_callbacks(可继承:会随 get_child() 传给子 run,你在 config 里传的就是这类)+ local_callbacks(本层专属:只听这一个组件)。区分继承性,才能做到"给根链装一个 handler,全树事件都听得到"。② 按环境变量加"官方机位"verbose=True → 自动挂 StdOutCallbackHandler(解说员);debug 模式 → ConsoleCallbackHandler;开了 LangSmith 追踪 → 自动挂 LangChainTracer(manager.py:2524-2532,Day19 主角,明天后天见)。③ 挂上下文钩子register_configure_hook 注册的 ContextVar handler(如 usage 统计)也在这里被拾取。所以"with 一个上下文管理器就能全局监听"不是魔法。LANGSMITH_TRACING=true 环境变量,什么代码都不改就有了全链路追踪。chain.invoke(x, config={"callbacks": [h]}) 传的是可继承听众,全树可见;而 chain.with_config(callbacks=[h]) 绑在特定 Runnable 上的行为不同层级会有差异。发现"我的 handler 只收到一半事件"时,先检查它是从哪条路上车的。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 的步骤"(比如吃整个输入才能算的自定义函数),流到它就断了——它上游的块会被攒齐再往下传。排查"为什么我的链不逐字出"时先找这种堵点。chain = prompt | model | StrOutputParser(),执行 for chunk in chain.stream({"topic": "冰淇淋"}) → 依次打印 "为"、"什"、"么"、"冰淇淋"…… 每个 chunk 是 str(因为最后一环 parser 的 transform 把 AIMessageChunk 转成了字符串片段)。同一时刻,如果挂了 handler,它的 on_llm_new_token 也在被逐个调用。stream 调用,组合链时不用关心每个成员支不支持流式——不支持的成员自动退化为"一大块"。代价是退化是静默的:你以为在流式,其实在等全量。LangChain 选择了可组合性优先,把识别堵点的责任留给使用者。astream_events:把回调变成一条事件流
stream 只能拿到最终输出的碎片;想要"检索完成了、模型开始了"这类中间事件,用 astream_events(libs/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": ...} 字典,塞进内存队列(_MemoryStream,event_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_start、on_chat_model_stream(data.chunk 是 AIMessageChunk)、on_retriever_end(data.output 是文档列表)……docstring 里有完整对照表(runnables/base.py:1399-1420)。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 实现"后台跑前台读"。
串起来 + 今日小结
with_listeners(on_end=_exit_history) 现在完全透明了:它就是往回调总线上挂了一个"只听 on_chain_end 的听众"(底层是 tracers 里的 RootListenersTracer),链跑完收到带完整输入输出的 Run → 写历史。记忆、日志、流式、追踪,四件看似无关的事,全是同一条总线上的不同听众。🧠 今天你应该能回答
- 回调系统解决什么问题?(无侵入地观察执行过程:观察者模式 + 事件总线)
- handler 抛异常会砸掉主流程吗?(默认不会,
raise_error=False,handle_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:自装听众+队列
runnables/ 目录里剩下的六件高级积木(passthrough/branch/fallbacks/retry/configurable/router)一次配齐。