回调系统 (Callbacks)
⚙️ LangChain 的回调系统解决了什么问题?为什么需要它?¶
在 LangChain 的早期版本中,想要监控一个链的运行状态,通常只能修改链的内部源码,比如在 _call 方法里硬插入 print 语句。这种做法极其脆弱:一旦框架升级,修改就会丢失;而且无法根据环境(开发/生产)灵活开启或关闭。回调系统正是为了用一种无侵入、可组合的方式解决这些横切关注点而生的。
🔍 回调系统本质上是一个事件驱动的观察者模式实现。 它在链、LLM、工具、Agent 等关键执行路径上预埋了数十个“钩子”。开发者只需实现一个或多个回调处理器,并将它们注册到链上,就能在这些事件发生时自动收到通知并执行自定义逻辑,而完全不需要触碰框架核心代码。
🎯 它具体解决的核心问题包括:
-
全链路可观测性:在一次对话中,LLM 被调用了几次?每次 Prompt 是什么?返回了什么?消耗了多少 Token?Agent 的每一步推理(Thought/Action/Observation)是怎样的?这些信息通过回调可以被完整地记录下来,形成类似分布式追踪的调用链。每个事件都带有
run_id和parent_run_id,让你能串联起从用户请求到最终回答的全部过程。 -
无侵入的日志与调试:开发复杂 Agent 时,最头疼的就是“模型为什么选了工具 A 而不是 B?”、“为什么在这一步卡死了?”。通过
on_agent_action和on_tool_end回调,你可以实时打印出 Agent 的思维过程,快速定位问题,而不需要在整个 Agent 代码里到处打补丁。 -
成本控制与资源管理:大模型调用费用高昂。你可以在
on_llm_end回调中实时累加 Token 用量,并在达到预算阈值时抛出异常或发送告警,实现自动熔断。这在自动化任务(如 AutoGPT)中至关重要。 -
流式输出的实现:
on_llm_new_token回调是前端实现“打字机效果”的基础。它每产生一个新 Token 就会触发一次,服务端可以将这些 Token 通过 WebSocket 或 SSE 推送给前端。 -
安全与合规审计:在某些行业,所有 AI 交互必须存档。回调可以在不修改业务链的前提下,将所有对话原文、工具调用结果加密后存入审计系统。
-
A/B 测试与实验追踪:你可以通过回调将不同链、不同模型的运行数据发送到 MLflow 或 LangSmith,对比效果。
💡 为什么不能只用日志?
普通日志只能记录零散的字符串,无法携带结构化的上下文(如运行 ID、父子关系、标签)。回调系统提供了标准化的上下文对象,使得监控工具能够自动构建出完整的调用拓扑图,这是日志做不到的。
📋 列出至少 8 种 LangChain 支持的回调事件类型。¶
LangChain 的回调事件覆盖了所有关键组件。以下是我在实际项目中使用最频繁的 12 种,按组件分类:
🔗 链 (Chain) 级别
-
on_chain_start:任何链开始执行时触发。可获取链的名称、输入。 -
on_chain_end:链执行完毕时触发。可获取输出。 -
on_chain_error:链执行过程中发生异常时触发。可获取异常对象。
🧠 LLM 级别
-
on_llm_start:LLM 开始生成前触发。可获取 Prompt、模型名称。 -
on_llm_end:LLM 生成完成后触发。可获取完整的生成结果及 Token 消耗。 -
on_llm_error:LLM 调用出错时触发。 -
on_llm_new_token:流式生成时,每产生一个新 Token 就触发一次。仅在流式模式有效。
🛠️ 工具 (Tool) 级别
-
on_tool_start:工具被调用前触发。可获取工具名称和输入参数。 -
on_tool_end:工具执行完毕后触发。可获取返回值。 -
on_tool_error:工具执行出错时触发。
🤖 Agent 级别
-
on_agent_action:Agent 决定调用某个工具时触发。可获取 Action 对象(工具名、输入、日志)。 -
on_agent_finish:Agent 输出最终答案时触发。可获取AgentFinish对象。
📝 其他
-
on_text:某些链产生文本输出时触发,用于兼容旧版。 -
on_retriever_start/end:检索器调用前后触发。
这些事件已经足够构建一个完整的监控系统。
📝 写出一个自定义回调处理器,用于将 LLM 调用日志写入文件。¶
下面是一个生产可用的 LLMFileLogger 实现。它考虑了日志轮转、JSON 格式化、长文本截断、以及多线程安全。
import json
import time
import logging
from datetime import datetime
from logging.handlers import RotatingFileHandler
from langchain.callbacks import BaseCallbackHandler
class LLMFileLogger(BaseCallbackHandler):
"""将 LLM 调用详细信息写入文件,支持日志轮转和多线程安全"""
def __init__(self, log_file: str = "llm_calls.log", max_bytes: int = 10*1024*1024, backup_count: int = 5):
self._start_times = {}
# 为每个实例创建独立的 logger,避免与 root logger 冲突
logger_name = f"LLMCallback_{id(self)}"
self.logger = logging.getLogger(logger_name)
self.logger.setLevel(logging.INFO)
# 使用 RotatingFileHandler 自动轮转
handler = RotatingFileHandler(
log_file, maxBytes=max_bytes, backupCount=backup_count, encoding='utf-8'
)
handler.setFormatter(logging.Formatter('%(message)s'))
self.logger.addHandler(handler)
def on_llm_start(self, serialized, prompts, **kwargs):
run_id = str(kwargs.get("run_id"))
self._start_times[run_id] = time.time()
# 构建日志条目,长 prompt 截断
log_entry = {
"event": "llm_start",
"model": serialized.get("name", serialized.get("id", ["unknown"])[0]),
"prompts": [p[:500] for p in prompts], # 截断,避免单条日志过大
"timestamp": datetime.now().isoformat(),
"run_id": run_id
}
self.logger.info(json.dumps(log_entry, ensure_ascii=False))
def on_llm_end(self, response, **kwargs):
run_id = str(kwargs.get("run_id"))
start = self._start_times.pop(run_id, None)
latency = time.time() - start if start else None
# 兼容不同 LLM 的 token 返回格式
token_usage = {}
if hasattr(response, 'llm_output') and response.llm_output:
token_usage = response.llm_output.get('token_usage', {})
# 获取生成文本
text = ""
if response.generations:
text = response.generations[0][0].text
log_entry = {
"event": "llm_end",
"run_id": run_id,
"latency_sec": round(latency, 3) if latency else None,
"prompt_tokens": token_usage.get("prompt_tokens", 0),
"completion_tokens": token_usage.get("completion_tokens", 0),
"total_tokens": token_usage.get("total_tokens", 0),
"response_preview": text[:200],
"timestamp": datetime.now().isoformat()
}
self.logger.info(json.dumps(log_entry, ensure_ascii=False))
def on_llm_error(self, error, **kwargs):
run_id = str(kwargs.get("run_id"))
log_entry = {
"event": "llm_error",
"run_id": run_id,
"error": str(error),
"timestamp": datetime.now().isoformat()
}
self.logger.error(json.dumps(log_entry, ensure_ascii=False))
使用:
llm_logger = LLMFileLogger("my_llm.log")
chain = ConversationChain(llm=OpenAI(), callbacks=[llm_logger])
chain.predict(input="Hello")
这会生成结构化的 JSON 日志文件,可以直接被 Logstash 或 Filebeat 采集。
📊 如何通过回调统计每次 LLM 调用的 token 消耗和延迟?¶
上面 LLMFileLogger 已经实现了基础功能。如果你需要一个专门的内存统计器用于实时监控,可以这样写:
from collections import defaultdict
import time
from langchain.callbacks import BaseCallbackHandler
class TokenMetricsCollector(BaseCallbackHandler):
"""统计 LLM 调用的 Token 消耗和延迟,并提供查询接口"""
def __init__(self):
self.metrics = defaultdict(list)
self._start_times = {}
def on_llm_start(self, serialized, prompts, **kwargs):
run_id = str(kwargs.get("run_id"))
self._start_times[run_id] = time.time()
def on_llm_end(self, response, **kwargs):
run_id = str(kwargs.get("run_id"))
start = self._start_times.pop(run_id, None)
if start:
self.metrics["latency"].append(time.time() - start)
token_usage = response.llm_output.get("token_usage", {}) if response.llm_output else {}
self.metrics["prompt_tokens"].append(token_usage.get("prompt_tokens", 0))
self.metrics["completion_tokens"].append(token_usage.get("completion_tokens", 0))
def get_summary(self):
"""返回汇总统计"""
latencies = self.metrics["latency"]
return {
"total_calls": len(latencies),
"avg_latency": sum(latencies)/len(latencies) if latencies else 0,
"p99_latency": sorted(latencies)[int(len(latencies)*0.99)] if len(latencies) > 100 else None,
"total_prompt_tokens": sum(self.metrics["prompt_tokens"]),
"total_completion_tokens": sum(self.metrics["completion_tokens"]),
"total_cost_estimate": self._estimate_cost()
}
def _estimate_cost(self):
# 按 OpenAI 定价估算,实际请根据模型调整
prompt_cost = sum(self.metrics["prompt_tokens"]) * 0.03 / 1000
completion_cost = sum(self.metrics["completion_tokens"]) * 0.06 / 1000
return prompt_cost + completion_cost
💡 注意事项:
-
对于异步 LLM 调用,
on_llm_start和on_llm_end可能不在同一个线程,run_id是关键。必须使用kwargs.get("run_id")来关联。 -
部分 LLM 实现(如 Anthropic)的 Token 统计可能不在
llm_output中,而在generations的generation_info里,需要做兼容处理。
🔗 如何将回调处理器绑定到一条链上?全局设置和局部设置的区别。¶
局部绑定:在链初始化或运行时通过 callbacks 参数传入。这是最推荐的方式,因为它精确、无副作用。
# 方式1:构造函数传入
chain = ConversationChain(llm=llm, callbacks=[my_handler])
# 方式2:调用时传入(优先级最高)
chain.run("Hello", callbacks=[my_handler])
局部绑定的回调只会影响当前链及其内部调用(如 LLM、工具),其他链不受干扰。
全局绑定:通过全局回调管理器设置,对所有后续创建的链生效。
from langchain.callbacks import CallbackManager
from langchain.globals import set_global_callback_manager
manager = CallbackManager([my_global_handler])
set_global_callback_manager(manager)
⚠️ 全局 vs 局部的坑:
-
全局污染:在单元测试中,如果启用了全局回调,可能导致日志输出混乱。务必在测试后重置。
-
合并策略:当你同时设置了全局和局部回调时,它们会被合并。LangChain 内部会创建一个
CallbackManager,将全局回调放在前面,局部回调放在后面,全部执行。这意味着一个事件可能被两个处理器都监听到。 -
不可变性:局部回调列表在链内部会被转化为元组,保证执行期间不会被意外修改。
📌 最佳实践:为生产环境设置全局回调(如错误追踪、成本统计),在开发调试时使用局部回调覆盖详细日志。
🕵️ 在 Agent 执行过程中,你如何通过回调获取每一步的 Action 和 Observation?¶
Agent 的回调主要依赖 on_agent_action 和 on_tool_end。
class AgentStepRecorder(BaseCallbackHandler):
def __init__(self):
self.steps = [] # 每一步的完整记录
def on_agent_action(self, action, **kwargs):
# action 是 AgentAction 对象,包含 tool, tool_input, log
step = {
"thought": action.log, # 包含 Thought: ... 部分
"tool": action.tool,
"tool_input": action.tool_input,
"observation": None
}
self.steps.append(step)
def on_tool_end(self, output, **kwargs):
# 工具返回结果,补充到最新的一步
if self.steps:
self.steps[-1]["observation"] = output
def on_agent_finish(self, finish, **kwargs):
# finish.return_values 包含最终的输出
pass
在 ReAct Agent 中,action.log 是类似这样的字符串:
通过记录每一步,你可以完整回溯 Agent 的推理路径。这对于调试“Agent 为什么做出某个决定”非常有价值。
🎚️ 如果只在开发环境启用详细日志回调,而在生产环境关闭,你如何控制?¶
使用环境变量或配置文件来动态决定启用哪些回调。典型的实现是工厂模式:
import os
from langchain.callbacks import BaseCallbackHandler
def get_callback_handlers():
env = os.getenv("APP_ENV", "development")
handlers = []
# 生产环境只保留成本统计
handlers.append(TokenMetricsCollector())
# 开发环境添加详细日志
if env == "development":
handlers.append(LLMFileLogger("dev_llm.log"))
handlers.append(AgentStepRecorder())
# 测试环境可能不需要任何回调
elif env == "test":
return []
return handlers
# 在创建链时使用
handlers = get_callback_handlers()
chain = SomeChain(callbacks=handlers)
这样,你只需要设置环境变量 APP_ENV=production 就能轻松切换。高级用法还可以使用 Feature Flag 服务,实现动态下线和灰度。
🔗 LangChain 的回调是如何实现“责任链模式”的?多个回调处理器的执行顺序是怎样的?¶
LangChain 的 CallbackManager 内部维护了一个回调处理器列表。当事件发生时,它会按顺序遍历这个列表,依次调用每个处理器的相应方法。
责任链模式体现在:每个处理器都有机会处理事件,但不能中断后续处理器的执行(除非抛出异常)。这与经典的责任链(可终止)略有不同,更类似于“广播”模式。
执行顺序:
-
全局回调(
global_callbacks)排在前面。 -
局部回调(
local_callbacks)排在后面。 -
在同一个列表内部,按照注册的顺序执行。
如果某个处理器的方法抛出异常,默认行为是停止执行后续处理器,并将异常向上传播。你可以通过 CallbackManager.configure(raise_error=False) 来改变这一行为,让异常被静默吞掉,但这通常不推荐,因为它会掩盖问题。
💡 注意:异步回调的执行顺序同样遵循上述规则,但由于异步并发,同一个处理器的不同事件方法可能在时间线上是交错的。
⚡ 异步回调(async callbacks)和同步回调在实现上有什么不同?¶
同步回调继承 BaseCallbackHandler,实现 on_llm_start 等方法。这些方法会在主线程(或当前线程)中同步执行,如果其中做了耗时操作(如写入磁盘),会阻塞链的继续执行。
异步回调继承 AsyncCallbackHandler,必须实现 on_llm_start 等方法的异步版本(即 async def on_llm_start)。当你在异步链(如 chain.arun())中使用异步回调时,LangChain 会 await 这些方法,从而不会阻塞事件循环。
关键区别:
-
如果异步链使用了同步回调,LangChain 会通过
asyncio.to_thread将同步方法调度到线程池中执行,以避免阻塞事件循环。但这会引入线程切换的开销,并且可能无法保证严格的执行顺序。 -
同步链使用异步回调时,会直接报错,因为同步环境无法执行协程。
-
因此,在异步应用中必须使用异步回调,以获得最佳性能和正确性。
示例异步回调:
from langchain.callbacks import AsyncCallbackHandler
class AsyncLLMLogger(AsyncCallbackHandler):
async def on_llm_start(self, serialized, prompts, **kwargs):
await self._async_write_log("start", prompts)
async def _async_write_log(self, event, data):
async with aiofiles.open("log.txt", "a") as f:
await f.write(f"{event}: {data}\n")
📌 在自定义回调中,你可以获取到哪些上下文信息?(run_id, parent_run_id, tags 等)¶
在回调的 kwargs 参数中,LangChain 会注入丰富的上下文信息,这些信息是构建全链路追踪的关键。
典型用法:
-
结合
run_id和parent_run_id生成调用链路图,类似于 OpenTelemetry 的 Span。 -
通过
tags实现按实验组分流统计。 -
利用
metadata传递用户标识,实现用户级的成本核算。
def on_llm_start(self, serialized, prompts, **kwargs):
user_id = kwargs.get("metadata", {}).get("user_id", "unknown")
self.logger.info(f"User {user_id} LLM call started, run_id={kwargs['run_id']}")
这些上下文信息让回调系统超越了简单的日志打印,成为构建可观测性平台的基础设施。
🖨️ 如何利用回调实现一个简单的流式输出到 Stdout?StreamingStdOutCallbackHandler 的原理是什么?¶
流式输出是现代对话应用的基本体验。LangChain 内置的 StreamingStdOutCallbackHandler 正是通过回调机制,将 LLM 生成的新 token 实时打印到标准输出。
原理
-
当 LLM 开启流式模式(例如
ChatOpenAI(model="gpt-4", streaming=True))时,API 返回的是一个生成器,逐块返回生成的文本。 -
LangChain 在调用 LLM 的过程中,每收到一个 token,就会触发
on_llm_new_token(self, token: str, **kwargs)回调事件。 -
StreamingStdOutCallbackHandler重写了这个方法,将接收到的 token 直接print(token, end="", flush=True)输出到标准输出,并且可选的当生成开始时打印一个前缀,结束后打印换行符。
底层源码简析(概念):
class StreamingStdOutCallbackHandler(BaseCallbackHandler):
def on_llm_new_token(self, token: str, **kwargs) -> None:
sys.stdout.write(token)
sys.stdout.flush()
自定义流式输出回调
你可以仿照它实现更复杂的功能,比如只显示前 N 个字符,或者把 token 发送到 WebSocket。
from langchain.callbacks import BaseCallbackHandler
class MyStreamHandler(BaseCallbackHandler):
def on_llm_new_token(self, token: str, **kwargs) -> None:
# 自定义处理,例如发送给前端
print(f"[{token}]", end="")
使用
将 streaming=True 的 LLM 和该回调一起传入链,即可实现流式输出。
llm = ChatOpenAI(model="gpt-4", streaming=True)
chain = ConversationChain(llm=llm, callbacks=[MyStreamHandler()])
chain.run("你好") # token 逐个出现在 stdout
注意事项
-
必须设置
streaming=True,否则on_llm_new_token不会被触发。 -
流式回调只在 LLM 生成阶段有效,Agent 的中间推理步骤不会经过此回调,除非你单独处理。
🧪 你如何用回调将链的执行轨迹发送到 LangSmith?¶
LangSmith 是 LangChain 官方的可观测性平台。将执行轨迹发送到 LangSmith 非常简单,通常不需要自己写回调,只需要设置环境变量。LangChain 内部已经集成了一个专门用于 LangSmith 的回调处理器,它会自动捕获链、LLM、工具的所有事件并上报。
启用方法:
-
注册 LangSmith 账号,获取 API Key。
-
设置环境变量:
export LANGCHAIN_TRACING_V2=true
export LANGCHAIN_API_KEY=your_api_key
export LANGCHAIN_PROJECT=your_project_name # 可选
- 运行任何 LangChain 链或 Agent,执行轨迹会自动出现在 LangSmith 控制台。
工作原理
当 LANGCHAIN_TRACING_V2 设置为 true 时,LangChain 在初始化全局回调管理器时,会自动向其中注入 LangChainTracer 回调处理器。这个处理器内部实现了所有回调事件的方法,将每个事件及其上下文(run_id, parent_run_id, 输入输出等)序列化并异步发送到 LangSmith 后端。
自定义上报
如果你需要自定义一些字段(比如为每次运行添加用户 ID、实验标签),可以在链调用时通过 metadata 和 tags 参数传递,它们会被 LangSmith 记录。
注意
-
该功能需要在能访问外网的环境中使用。
-
它会产生一定的网络延迟,但基本不影响主链性能,因为上报是异步非阻塞的。
-
可以在代码中动态关闭,例如测试环境不发送。
💥 在回调中,如果某个 handler 抛出异常,会影响主链的执行吗?如何隔离?¶
默认行为:会影响。BaseCallbackManager 在处理事件时,会按顺序调用每个 handler 的方法。如果任何一个 handler 抛出异常,异常会立即向上传播,导致当前正在执行的链或 Agent 失败。
这意味着一个“脆弱”的回调(比如写日志时磁盘满了)可能导致整个业务流程中断。在生产环境中这是不可接受的。
隔离方法:
- 在 handler 内部捕获异常 这是最推荐的做法。每个 handler 应该自己保证健壮性,将异常控制在内部,不影响主链。
def on_llm_end(self, response, **kwargs):
try:
self._write_to_db(response)
except Exception as e:
# 记录错误,但不抛出
print(f"回调异常: {e}")
- 使用全局配置
raise_error=False在创建CallbackManager时,可以设置raise_error=False。这样,当某个 handler 抛出异常时,管理器会静默捕获并继续执行后续 handler。但这也会隐藏潜在的错误,不推荐作为唯一手段。
-
将高风险回调包装为“安全回调” 编写一个通用的装饰器或包装类,在调用 handler 方法时自动加 try/except。
-
使用异步回调和非阻塞机制 对于耗时或高风险操作,应该放在异步回调中执行,或者将任务提交给线程池/消息队列,彻底解耦。
最佳实践:永远不要在回调处理器中抛出未捕获的异常。每个 on_* 方法应该是“防御性”的。你可以把异常记录到日志中,并通过监控告警,而不是让主流程崩溃。
🎯 如何只对链中的某个特定 LLM 调用附加回调,而其他部分不受影响?¶
LangChain 的链支持层级化回调。你可以利用局部回调的作用域特性,为某个特定的 LLM 实例或子链单独设置回调,而不影响外层链。
方法 1:在 LLM 实例化时绑定回调
special_llm = ChatOpenAI(model="gpt-4", callbacks=[MySpecialHandler()])
normal_llm = ChatOpenAI(model="gpt-3.5-turbo")
chain = MyChain(llm=special_llm, other_llm=normal_llm)
这里 special_llm 的每次调用都会触发 MySpecialHandler,而 normal_llm 不会。
方法 2:使用 with_config 动态绑定(LCEL)
chain_with_callback = my_chain.with_config(callbacks=[MyHandler()])
result = chain_with_callback.invoke({"input": "..."})
这种方法不会修改原链,只会在这次调用中注入回调。
方法 3:通过 Runnable 的 config 传递
在 LCEL 中,可以在运行配置里传 callbacks:
只有这次调用的所有子步骤会继承该配置。
注意:回调有“继承”特性。当你在父链上设置了回调,子链和 LLM 会自动继承这些回调(除非子组件显式屏蔽)。但反过来,如果你在子组件上设置了回调,父链不会受影响。所以,要实现“只对特定 LLM 生效”,直接在目标 LLM 上设置回调即可。
🔗 在 LCEL 链中,如何通过 .with_config({"callbacks": [...]}) 动态传入回调?¶
LCEL 链支持通过 config 动态传入回调,这是实现“一次调用使用特定监控”的最优雅方式。
使用 with_config
from langchain_core.runnables import RunnableLambda
my_chain = prompt | llm | output_parser
# 创建一个携带回调的新链副本
monitored_chain = my_chain.with_config(callbacks=[MyHandler()])
result = monitored_chain.invoke({"question": "..."})
with_config 返回一个新的 Runnable,它会在执行时将指定的 callbacks 合并到运行上下文中。
直接在 invoke 时传递 config
这两种方式效果相同,都会在这次调用中激活回调,且不会影响其他并发调用。
原理:LCEL 在执行时,会从 config 中提取 callbacks,并与当前已有的全局回调合并,生成一个 CallbackManager。这个管理器会被传递给当前调用树中的所有子 Runnable。
最佳实践:
-
对于需要不同监控策略的请求(例如 VIP 用户的请求需要详细追踪),使用
invoke时动态传入回调。 -
避免在
with_config中硬编码复杂回调逻辑,保持配置简洁。
🛡️ 你如何使用回调来实时收集生成文本,并做内容安全检测?¶
你可以利用 on_llm_new_token 回调实时收集 LLM 生成的文本片段,并在收集到一定长度或遇到特定分隔符时,触发内容安全检测。
示例实现:
from langchain.callbacks import BaseCallbackHandler
import re
class SafetyCheckHandler(BaseCallbackHandler):
def __init__(self, block_words=None):
self.buffer = ""
self.block_words = block_words or ["暴力", "色情"]
self.blocked = False
def on_llm_new_token(self, token: str, **kwargs) -> None:
if self.blocked:
return # 已经检测到违规,不再处理
self.buffer += token
# 每累积一定长度检查一次
if len(self.buffer) >= 20 or '\n' in self.buffer:
for word in self.block_words:
if word in self.buffer:
print(f"警告:检测到敏感词 '{word}',生成内容被阻止。")
self.blocked = True
# 这里可以抛出异常来中断生成,但需要小心
raise StopIteration("安全检测中断")
self.buffer = "" # 清空缓存,继续检测下一段
def on_llm_end(self, response, **kwargs):
if self.buffer and not self.blocked:
# 处理剩余未检测的内容
for word in self.block_words:
if word in self.buffer:
print(f"警告:生成内容包含敏感词 '{word}'。")
使用:将该处理器传入 LLM 的 callbacks,并启用流式模式。
注意:
-
敏感词检测只是最简单的方式。生产环境应使用专门的内容安全 API(如 OpenAI Moderation API)或更复杂的 NLP 模型。
-
如果你希望在检测到违规时立即中断生成,可以在回调中抛出异常。但需要外层链路捕获该异常并返回错误提示。
📢 编写一个回调,当 LLM 返回包含特定关键词时发出告警。¶
下面是一个实用的告警回调,它会在 LLM 生成结束时检查完整回复,如果包含特定关键词,则通过日志、邮件或消息推送告警。
import logging
from langchain.callbacks import BaseCallbackHandler
class KeywordAlertHandler(BaseCallbackHandler):
def __init__(self, keywords, alert_func=None):
self.keywords = set(keywords)
self.alert_func = alert_func or (lambda msg: logging.warning(msg))
def on_llm_end(self, response, **kwargs):
text = response.generations[0][0].text if response.generations else ""
found = [kw for kw in self.keywords if kw in text]
if found:
msg = f"告警:LLM 回复中包含关键词: {', '.join(found)}"
self.alert_func(msg)
# 使用示例
handler = KeywordAlertHandler(
keywords=["机密", "内部文件", "密码"],
alert_func=lambda msg: requests.post(webhook_url, json={"text": msg})
)
为什么不在 on_llm_new_token 中做?
流式检测容易出现截断导致的误报(比如关键词被切成两半),在生成结束后做全量检查更准确。
注意:如果生成内容很长,一次性检查可能会消耗内存。可以结合流式缓冲,但告警仍建议在终点触发。
☁️ 在分布式部署时,回调产生的日志如何汇总?¶
分布式环境下,每个服务实例(Pod)都会产生自己的回调日志。要汇总分析,需要一个集中的日志收集系统。
方案一:标准 ELK/EFK 栈
-
每个实例将回调日志以 JSON 格式输出到标准输出或文件。
-
使用 Filebeat 采集日志文件,或通过 Docker 日志驱动收集标准输出。
-
发送到 Elasticsearch,用 Kibana 进行可视化。
方案二:直接发送到外部监控平台
-
在自定义回调处理器内部,将事件直接通过 HTTP 发送到 LangSmith、Datadog、Prometheus Pushgateway 等。
-
优点是不依赖文件采集,实时性好。
-
注意网络故障时的容错处理。
方案三:消息队列中转
-
回调处理器将事件写入本地消息队列(如 Redis Streams、Kafka)。
-
独立的消费者服务从队列拉取,进行聚合、存储和分析。
-
这种方式解耦性好,不会阻塞主链。
关键设计:
-
所有日志必须包含实例标识(如 pod name、host IP)和全局 Trace ID(run_id),以便串联跨实例的调用链。
-
使用异步发送,避免 I/O 拖慢主链。
-
设置合理的采样率,高流量时只记录部分请求,防止日志爆炸。
⏱️ 如何防止回调中执行耗时操作(如网络请求)拖慢主链?¶
回调在主线程中同步执行(除非使用异步回调且在异步环境中),如果有耗时操作,会直接增加用户感知的延迟。
防护策略:
-
使用异步回调(AsyncCallbackHandler):如果你的应用是基于
asyncio的,实现异步回调,耗时操作使用await,不会阻塞事件循环。 -
将耗时任务提交到线程池或任务队列:在同步回调中,使用
concurrent.futures.ThreadPoolExecutor或asyncio.run_coroutine_threadsafe将任务转交给后台线程执行。
from concurrent.futures import ThreadPoolExecutor
executor = ThreadPoolExecutor(max_workers=2)
def on_llm_end(self, response, **kwargs):
executor.submit(self._slow_report, response)
-
写入内存队列,异步消费:回调将事件放入本地的
queue.Queue,另一个后台线程不断从队列取数据并批量发送。 -
设置回调超时:如果使用线程池,可以用
future.result(timeout=...)来限制最大等待,但回调方法本身不应等待结果。 -
降低回调频率:对于高频事件(如
on_llm_new_token),只在满足特定条件(如每 10 个 token 或遇到换行)时才触发耗时逻辑。
总结:永远不要在回调方法体内做同步的网络 I/O 或磁盘写入。如果必须做,一定要“甩”出去。
🧪 你如何测试自定义回调处理器的正确性?¶
回调处理器的测试属于单元测试范畴,核心是验证“在特定事件发生时,处理器是否执行了预期的动作”。
方法:
- 直接实例化处理器,手动调用事件方法:
handler = TokenMetricsCollector()
# 模拟 on_llm_start
handler.on_llm_start({"name": "test"}, ["prompt"], run_id="123")
# 构造一个模拟的 response 对象
mock_response = MagicMock()
mock_response.llm_output = {"token_usage": {"total_tokens": 100}}
handler.on_llm_end(mock_response, run_id="123")
# 断言
assert handler.get_summary()["total_calls"] == 1
- 使用 mock 的 LLM 和链,在集成环境中触发回调:利用
FakeLLM或unittest.mock来模拟链的执行,验证回调的调用次数和参数。
from unittest.mock import patch, call
with patch.object(handler, 'on_llm_end') as mock_end:
chain.run("test")
mock_end.assert_called_once()
-
测试异常处理:模拟回调内部抛出异常,验证你的 try/except 是否生效,以及主链是否被影响。
-
验证副作用:比如检查日志文件是否生成、内存统计是否正确、外部 API 是否被调用。
建议:回调处理器应该尽量保持“纯函数”特性,依赖注入外部资源(如文件路径、API 客户端),这样测试时容易替换。
📦 什么是“BaseCallbackManager”?它和列表式的回调处理有什么区别?¶
BaseCallbackManager 是 LangChain 回调系统的核心管理器。它不是一个简单的回调列表,而是一个具有事件分发、异常处理和上下文传播能力的容器。
主要区别:
-
事件分发:它实现了所有
on_*方法,当这些方法被调用时,会遍历内部的 handler 列表,依次调用每个 handler 的对应方法。你可以把它看作一个“代理”。 -
继承与合并:它支持通过
copy()方法创建子管理器,并继承父管理器的回调。LCEL 中的链式调用会使用这一机制,保证回调在调用树中正确传播。 -
元数据管理:除了 handler 列表,它还携带了
tags、metadata、inheritable_handlers等属性,这些属性会随事件传递给每个 handler,使得 handler 可以获取运行上下文。 -
异常处理策略:通过
raise_error属性控制 handler 异常是否向上传播。 -
全局与局部整合:当一条链启动时,LangChain 会合并全局回调管理器和局部回调,生成一个最终的
CallbackManager实例。
相比之下,普通的列表式回调只是简单的 handler 集合,没有上述高级功能。你直接传 callbacks=[handler1, handler2] 时,LangChain 内部会立即用这个列表构造一个 CallbackManager。
总结:BaseCallbackManager 是 LangChain 回调功能得以运作的“幕后大脑”,它使得回调可以优雅地传播、合并和配置。
🔧 如何通过环境变量或配置文件来控制是否启用回调?¶
最灵活的方式是使用工厂模式结合配置中心。
示例:
import os
from langchain.callbacks import BaseCallbackHandler
def get_callbacks():
env = os.getenv("APP_ENV", "development")
handlers = []
# 基础回调:所有环境都启用成本统计
handlers.append(TokenMetricsCollector())
# 开发环境启用详细日志和 LangSmith
if env == "development":
handlers.append(LLMFileLogger())
# 可选:开启 LangSmith
os.environ["LANGCHAIN_TRACING_V2"] = "true"
# 生产环境只启用关键告警和采样追踪
elif env == "production":
handlers.append(ErrorAlertHandler())
# 采样:基于配置的比率决定是否启用详细追踪
if os.getenv("ENABLE_DETAILED_TRACE", "false").lower() == "true":
handlers.append(DetailedTraceHandler())
return handlers
# 创建链时
chain = MyChain(callbacks=get_callbacks())
更高级的控制:
-
使用
ConfigParser读取.ini或yaml文件。 -
通过环境变量
LANGCHAIN_CALLBACKS指定完全限定类名(需框架支持)。 -
在微服务中,通过配置中心(如 Apollo、Nacos)动态刷新,而不需要重启应用。此时可以用一个
DynamicCallbackManager定期拉取最新配置。
核心原则:将“是否启用回调”与“如何实现回调”分离。生产环境的回调应极轻量且高可用,避免因回调本身的故障影响主业务。
📊 在 Agent 执行中,如何通过回调实现一个“执行进度条”?¶
Agent 的执行步数不确定,用户可能等待较长时间。通过回调向用户展示进度,能有效缓解等待焦虑。
实现思路:
-
在
on_agent_action中递增步数,并更新进度条。 -
在
on_agent_finish中标记完成。
基于 tqdm 的简单进度条:
from tqdm import tqdm
from langchain.callbacks import BaseCallbackHandler
class AgentProgressBar(BaseCallbackHandler):
def __init__(self, max_steps=10):
self.pbar = tqdm(total=max_steps, desc="Agent 推理", unit="step")
self.current_step = 0
def on_agent_action(self, action, **kwargs):
self.current_step += 1
self.pbar.update(1)
# 可以在描述中显示当前工具
self.pbar.set_postfix_str(f"工具: {action.tool}")
def on_agent_finish(self, finish, **kwargs):
self.pbar.close()
def on_chain_end(self, outputs, **kwargs):
if self.pbar.n < self.pbar.total:
self.pbar.close()
对于流式前端,你可以通过 WebSocket 发送进度事件。在回调中调用 websocket.send_json({"step": self.current_step, "tool": action.tool})。
注意:进度条的总步数是不确定的,你可以设置一个上限(如 max_iterations),或者使用无限进度条模式。
🔗 回调系统的设计会不会引入隐式耦合?你怎么看?¶
确实会引入隐式耦合,但这是一种有意识的设计权衡,并且可以通过规范来管理。
隐式耦合的表现:
-
回调处理器依赖于特定的事件名称和参数格式。如果 LangChain 版本升级改变了事件签名,回调可能静默失效。
-
多个回调之间存在执行顺序依赖(比如 A 必须在 B 之前执行),这种顺序是隐式的,容易出问题。
-
回调通过修改外部状态(如全局变量)来传递信息,这会导致组件之间不透明的交互。
我的看法:
回调本质上是一种横切关注点的实现,它天然会带来一定耦合。但它的优势——无侵入、可组合——远大于缺点。关键在于如何控制耦合程度:
-
保持回调独立:每个回调只做一件事,不依赖其他回调的存在。
-
不通过回调传递业务数据:回调主要用于观测和辅助,不要把核心业务逻辑放在回调里。
-
版本兼容:尽量依赖官方文档记录的稳定事件,避免使用内部 API。
-
测试覆盖:将回调纳入自动化测试,确保升级后行为一致。
总结:回调是“好的耦合”,是框架提供给开发者的观察窗口。只要不滥用,它不会成为代码腐烂的根源。
🎓 总结一下回调在 LangChain 中的最佳实践,以及你踩过的坑。¶
最佳实践:
-
分离关注点:日志、监控、安全检测使用独立的回调处理器,不要混在一个大类里。
-
防御性编程:每个回调方法都要用 try/except 包裹,确保不会中断主链。
-
异步优先:在异步应用中使用异步回调,避免阻塞事件循环。对同步回调中的耗时操作,务必提交到线程池。
-
环境区分:通过配置动态加载回调,开发环境详细记录,生产环境只保留关键指标。
-
携带上下文:充分利用
metadata和tags传递用户 ID、会话 ID 等信息,方便追踪。 -
使用 LangSmith 作为标准监控,自定义回调作为补充。
-
流式回调要处理截断:敏感词检测等要注意 token 切割问题。
-
测试回调:对每个回调编写单元测试,验证正常和异常场景。
我踩过的坑:
-
坑1:全局回调污染测试。早期在模块顶层设置了全局回调,导致所有单元测试都输出大量日志,排查很久。解决:使用
setUp/tearDown重置全局回调管理器。 -
坑2:流式检测敏感词时,token 被截断。关键词“暴力”可能被分成“暴”和“力”两次回调,导致检测失败。解决:维护一个缓冲区,按逻辑边界(如标点、换行)触发检测。
-
坑3:回调中写文件忘记 flush,导致日志丢失。应用崩溃时最后几秒的日志全丢。解决:使用
logging模块的FileHandler,它自带缓冲和刷新机制。 -
坑4:异步回调方法未实现,导致同步环境报错。在同步链里使用了
AsyncCallbackHandler,结果方法未被调用。解决:严格区分同步/异步处理器。 -
坑5:回调里抛异常导致整个 Agent 崩溃。一次磁盘满,日志回调写失败抛出
OSError,Agent 直接停止。解决:所有文件操作加 try/except,并设置raise_error=False。
结语:LangChain 的回调系统是构建生产级 LLM 应用的利器,但需要开发者像对待业务代码一样认真设计和测试回调。用好它,你能获得对整个系统无与伦比的洞察力。