跳转至

回调系统 (Callbacks)

⚙️ LangChain 的回调系统解决了什么问题?为什么需要它?

在 LangChain 的早期版本中,想要监控一个链的运行状态,通常只能修改链的内部源码,比如在 _call 方法里硬插入 print 语句。这种做法极其脆弱:一旦框架升级,修改就会丢失;而且无法根据环境(开发/生产)灵活开启或关闭。回调系统正是为了用一种无侵入、可组合的方式解决这些横切关注点而生的。

🔍 回调系统本质上是一个事件驱动的观察者模式实现。 它在链、LLM、工具、Agent 等关键执行路径上预埋了数十个“钩子”。开发者只需实现一个或多个回调处理器,并将它们注册到链上,就能在这些事件发生时自动收到通知并执行自定义逻辑,而完全不需要触碰框架核心代码。

🎯 它具体解决的核心问题包括:

  1. 全链路可观测性:在一次对话中,LLM 被调用了几次?每次 Prompt 是什么?返回了什么?消耗了多少 Token?Agent 的每一步推理(Thought/Action/Observation)是怎样的?这些信息通过回调可以被完整地记录下来,形成类似分布式追踪的调用链。每个事件都带有 run_idparent_run_id,让你能串联起从用户请求到最终回答的全部过程。

  2. 无侵入的日志与调试:开发复杂 Agent 时,最头疼的就是“模型为什么选了工具 A 而不是 B?”、“为什么在这一步卡死了?”。通过 on_agent_actionon_tool_end 回调,你可以实时打印出 Agent 的思维过程,快速定位问题,而不需要在整个 Agent 代码里到处打补丁。

  3. 成本控制与资源管理:大模型调用费用高昂。你可以在 on_llm_end 回调中实时累加 Token 用量,并在达到预算阈值时抛出异常或发送告警,实现自动熔断。这在自动化任务(如 AutoGPT)中至关重要。

  4. 流式输出的实现:on_llm_new_token 回调是前端实现“打字机效果”的基础。它每产生一个新 Token 就会触发一次,服务端可以将这些 Token 通过 WebSocket 或 SSE 推送给前端。

  5. 安全与合规审计:在某些行业,所有 AI 交互必须存档。回调可以在不修改业务链的前提下,将所有对话原文、工具调用结果加密后存入审计系统。

  6. A/B 测试与实验追踪:你可以通过回调将不同链、不同模型的运行数据发送到 MLflow 或 LangSmith,对比效果。

💡 为什么不能只用日志?

普通日志只能记录零散的字符串,无法携带结构化的上下文(如运行 ID、父子关系、标签)。回调系统提供了标准化的上下文对象,使得监控工具能够自动构建出完整的调用拓扑图,这是日志做不到的。


📋 列出至少 8 种 LangChain 支持的回调事件类型。

LangChain 的回调事件覆盖了所有关键组件。以下是我在实际项目中使用最频繁的 12 种,按组件分类:

🔗 链 (Chain) 级别

  1. on_chain_start:任何链开始执行时触发。可获取链的名称、输入。

  2. on_chain_end:链执行完毕时触发。可获取输出。

  3. on_chain_error:链执行过程中发生异常时触发。可获取异常对象。

🧠 LLM 级别

  1. on_llm_start:LLM 开始生成前触发。可获取 Prompt、模型名称。

  2. on_llm_end:LLM 生成完成后触发。可获取完整的生成结果及 Token 消耗。

  3. on_llm_error:LLM 调用出错时触发。

  4. on_llm_new_token:流式生成时,每产生一个新 Token 就触发一次。仅在流式模式有效。

🛠️ 工具 (Tool) 级别

  1. on_tool_start:工具被调用前触发。可获取工具名称和输入参数。

  2. on_tool_end:工具执行完毕后触发。可获取返回值。

  3. on_tool_error:工具执行出错时触发。

🤖 Agent 级别

  1. on_agent_action:Agent 决定调用某个工具时触发。可获取 Action 对象(工具名、输入、日志)。

  2. on_agent_finish:Agent 输出最终答案时触发。可获取 AgentFinish 对象。

📝 其他

  1. on_text:某些链产生文本输出时触发,用于兼容旧版。

  2. 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_starton_llm_end 可能不在同一个线程,run_id 是关键。必须使用 kwargs.get("run_id") 来关联。

  • 部分 LLM 实现(如 Anthropic)的 Token 统计可能不在 llm_output 中,而在 generationsgeneration_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_actionon_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 是类似这样的字符串:

Thought: I need to search for the weather.
Action: Search
Action Input: "New York weather"

通过记录每一步,你可以完整回溯 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 内部维护了一个回调处理器列表。当事件发生时,它会按顺序遍历这个列表,依次调用每个处理器的相应方法。

责任链模式体现在:每个处理器都有机会处理事件,但不能中断后续处理器的执行(除非抛出异常)。这与经典的责任链(可终止)略有不同,更类似于“广播”模式。

执行顺序:

  1. 全局回调(global_callbacks)排在前面。

  2. 局部回调(local_callbacks)排在后面。

  3. 在同一个列表内部,按照注册的顺序执行。

如果某个处理器的方法抛出异常,默认行为是停止执行后续处理器,并将异常向上传播。你可以通过 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_idparent_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、工具的所有事件并上报。

启用方法:

  1. 注册 LangSmith 账号,获取 API Key。

  2. 设置环境变量:

export LANGCHAIN_TRACING_V2=true
export LANGCHAIN_API_KEY=your_api_key
export LANGCHAIN_PROJECT=your_project_name   # 可选
  1. 运行任何 LangChain 链或 Agent,执行轨迹会自动出现在 LangSmith 控制台。

工作原理 当 LANGCHAIN_TRACING_V2 设置为 true 时,LangChain 在初始化全局回调管理器时,会自动向其中注入 LangChainTracer 回调处理器。这个处理器内部实现了所有回调事件的方法,将每个事件及其上下文(run_id, parent_run_id, 输入输出等)序列化并异步发送到 LangSmith 后端。

自定义上报 如果你需要自定义一些字段(比如为每次运行添加用户 ID、实验标签),可以在链调用时通过 metadatatags 参数传递,它们会被 LangSmith 记录。

chain.run("Hello", config={"metadata": {"user_id": "123"}})

注意

  • 该功能需要在能访问外网的环境中使用。

  • 它会产生一定的网络延迟,但基本不影响主链性能,因为上报是异步非阻塞的。

  • 可以在代码中动态关闭,例如测试环境不发送。


💥 在回调中,如果某个 handler 抛出异常,会影响主链的执行吗?如何隔离?

默认行为:会影响。BaseCallbackManager 在处理事件时,会按顺序调用每个 handler 的方法。如果任何一个 handler 抛出异常,异常会立即向上传播,导致当前正在执行的链或 Agent 失败。

这意味着一个“脆弱”的回调(比如写日志时磁盘满了)可能导致整个业务流程中断。在生产环境中这是不可接受的。

隔离方法:

  1. 在 handler 内部捕获异常 这是最推荐的做法。每个 handler 应该自己保证健壮性,将异常控制在内部,不影响主链。
def on_llm_end(self, response, **kwargs):
    try:
        self._write_to_db(response)
    except Exception as e:
        # 记录错误,但不抛出
        print(f"回调异常: {e}")
  1. 使用全局配置 raise_error=False 在创建 CallbackManager 时,可以设置 raise_error=False。这样,当某个 handler 抛出异常时,管理器会静默捕获并继续执行后续 handler。但这也会隐藏潜在的错误,不推荐作为唯一手段。
manager = CallbackManager(handlers=[...], raise_error=False)
  1. 将高风险回调包装为“安全回调” 编写一个通用的装饰器或包装类,在调用 handler 方法时自动加 try/except。

  2. 使用异步回调和非阻塞机制 对于耗时或高风险操作,应该放在异步回调中执行,或者将任务提交给线程池/消息队列,彻底解耦。

最佳实践:永远不要在回调处理器中抛出未捕获的异常。每个 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

chain.invoke({"input": "..."}, config={"callbacks": [MyHandler()]})

只有这次调用的所有子步骤会继承该配置。

注意:回调有“继承”特性。当你在父链上设置了回调,子链和 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

result = my_chain.invoke(
    {"question": "..."},
    config={"callbacks": [MyHandler()]}
)

这两种方式效果相同,都会在这次调用中激活回调,且不会影响其他并发调用。

原理: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 拖慢主链。

  • 设置合理的采样率,高流量时只记录部分请求,防止日志爆炸。


⏱️ 如何防止回调中执行耗时操作(如网络请求)拖慢主链?

回调在主线程中同步执行(除非使用异步回调且在异步环境中),如果有耗时操作,会直接增加用户感知的延迟。

防护策略:

  1. 使用异步回调(AsyncCallbackHandler):如果你的应用是基于 asyncio 的,实现异步回调,耗时操作使用 await,不会阻塞事件循环。

  2. 将耗时任务提交到线程池或任务队列:在同步回调中,使用 concurrent.futures.ThreadPoolExecutorasyncio.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)
  1. 写入内存队列,异步消费:回调将事件放入本地的 queue.Queue,另一个后台线程不断从队列取数据并批量发送。

  2. 设置回调超时:如果使用线程池,可以用 future.result(timeout=...) 来限制最大等待,但回调方法本身不应等待结果。

  3. 降低回调频率:对于高频事件(如 on_llm_new_token),只在满足特定条件(如每 10 个 token 或遇到换行)时才触发耗时逻辑。

总结:永远不要在回调方法体内做同步的网络 I/O 或磁盘写入。如果必须做,一定要“甩”出去。


🧪 你如何测试自定义回调处理器的正确性?

回调处理器的测试属于单元测试范畴,核心是验证“在特定事件发生时,处理器是否执行了预期的动作”。

方法:

  1. 直接实例化处理器,手动调用事件方法:
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
  1. 使用 mock 的 LLM 和链,在集成环境中触发回调:利用 FakeLLMunittest.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()
  1. 测试异常处理:模拟回调内部抛出异常,验证你的 try/except 是否生效,以及主链是否被影响。

  2. 验证副作用:比如检查日志文件是否生成、内存统计是否正确、外部 API 是否被调用。

建议:回调处理器应该尽量保持“纯函数”特性,依赖注入外部资源(如文件路径、API 客户端),这样测试时容易替换。


📦 什么是“BaseCallbackManager”?它和列表式的回调处理有什么区别?

BaseCallbackManager 是 LangChain 回调系统的核心管理器。它不是一个简单的回调列表,而是一个具有事件分发、异常处理和上下文传播能力的容器。

主要区别:

  • 事件分发:它实现了所有 on_* 方法,当这些方法被调用时,会遍历内部的 handler 列表,依次调用每个 handler 的对应方法。你可以把它看作一个“代理”。

  • 继承与合并:它支持通过 copy() 方法创建子管理器,并继承父管理器的回调。LCEL 中的链式调用会使用这一机制,保证回调在调用树中正确传播。

  • 元数据管理:除了 handler 列表,它还携带了 tagsmetadatainheritable_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 读取 .iniyaml 文件。

  • 通过环境变量 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 中的最佳实践,以及你踩过的坑。

最佳实践:

  1. 分离关注点:日志、监控、安全检测使用独立的回调处理器,不要混在一个大类里。

  2. 防御性编程:每个回调方法都要用 try/except 包裹,确保不会中断主链。

  3. 异步优先:在异步应用中使用异步回调,避免阻塞事件循环。对同步回调中的耗时操作,务必提交到线程池。

  4. 环境区分:通过配置动态加载回调,开发环境详细记录,生产环境只保留关键指标。

  5. 携带上下文:充分利用 metadatatags 传递用户 ID、会话 ID 等信息,方便追踪。

  6. 使用 LangSmith 作为标准监控,自定义回调作为补充。

  7. 流式回调要处理截断:敏感词检测等要注意 token 切割问题。

  8. 测试回调:对每个回调编写单元测试,验证正常和异常场景。

我踩过的坑:

  • 坑1:全局回调污染测试。早期在模块顶层设置了全局回调,导致所有单元测试都输出大量日志,排查很久。解决:使用 setUp/tearDown 重置全局回调管理器。

  • 坑2:流式检测敏感词时,token 被截断。关键词“暴力”可能被分成“暴”和“力”两次回调,导致检测失败。解决:维护一个缓冲区,按逻辑边界(如标点、换行)触发检测。

  • 坑3:回调中写文件忘记 flush,导致日志丢失。应用崩溃时最后几秒的日志全丢。解决:使用 logging 模块的 FileHandler,它自带缓冲和刷新机制。

  • 坑4:异步回调方法未实现,导致同步环境报错。在同步链里使用了 AsyncCallbackHandler,结果方法未被调用。解决:严格区分同步/异步处理器。

  • 坑5:回调里抛异常导致整个 Agent 崩溃。一次磁盘满,日志回调写失败抛出 OSError,Agent 直接停止。解决:所有文件操作加 try/except,并设置 raise_error=False

结语:LangChain 的回调系统是构建生产级 LLM 应用的利器,但需要开发者像对待业务代码一样认真设计和测试回调。用好它,你能获得对整个系统无与伦比的洞察力。