LangChain 源码阅读
AgentExecutor 执行循环、Tool 注册与调用、Memory 存储机制、Callback 事件系统源码深度解析
一、AgentExecutor 执行循环源码
AgentExecutor 是 LangChain 框架中驱动 Agent 持续推理与行动的核心引擎,其执行入口位于 _call() 方法。
1.1 _call() 方法主循环逻辑
# langchain/agents/agent.py (简化核心逻辑)
class AgentExecutor(Chain):
max_iterations: int = 15
early_stopping_method: str = "force"
def _call(self, inputs: Dict[str, Any]) -> Dict[str, Any]:
iterations = 0
intermediate_steps = []
# 将 memory 中的历史记录注入输入
inputs = {**inputs, **(self.memory.load_memory_variables(inputs) if self.memory else {})}
while self._should_continue(iterations):
# 1. Agent 根据当前输入和中间步骤决定下一步动作
output = self.agent.plan(intermediate_steps, **inputs)
# 2. 判断是最终输出还是需要执行工具
if isinstance(output, AgentFinish):
return self._return(output.return_values, intermediate_steps)
elif isinstance(output, AgentAction):
# 3. 执行工具调用
observation = self._take_next_step(output, intermediate_steps)
intermediate_steps.append((output, observation))
iterations += 1
else:
raise ValueError(f"Unexpected output type: {type(output)}")
# 超迭代次数后根据 early_stopping 策略处理
return self._early_stopping(intermediate_steps, **inputs)1.2 _should_continue 判断条件
def _should_continue(self, iterations: int) -> bool:
# 当前迭代次数小于最大迭代次数则继续
if iterations >= self.max_iterations:
return False
return Truemax_iterations 默认值为 15,防止 Agent 陷入无限循环。early_stopping_method 支持两种模式:"force"(直接返回已有结果)和 "generate"(让 LLM 基于中间步骤生成最终答案)。
1.3 _take_next_step 执行步骤
def _take_next_step(self, action: AgentAction, steps: List[Tuple]) -> Any:
try:
tool = self._get_tool(action.tool)
observation = tool.run(action.tool_input, verbose=self.verbose)
return observation
except Exception as e:
return f"Tool execution error: {e}"1.4 AgentFinish vs AgentAction
Agent 的 plan() 方法返回两种类型,AgentExecutor 据此判断流程走向:
| 返回类型 | 含义 | 执行结果 |
|---|---|---|
AgentFinish | Agent 认为任务完成,携带最终输出 | 直接返回 return_values |
AgentAction | Agent 需要调用某个 Tool,携带工具名和输入 | 执行 _take_next_step 并将 observation 追加到 intermediate_steps |
1.5 核心类职责对比
| 类名 | 职责 |
|---|---|
AgentExecutor | 编排 Agent 执行循环,控制迭代次数与停止策略 |
Agent | 定义推理逻辑,通过 plan() 输出 AgentAction 或 AgentFinish |
BaseSingleActionAgent | 单动作 Agent 基类,定义 plan() 接口契约 |
LLMSingleActionAgent | 使用 LLM 生成单个动作的 Agent 实现 |
二、Tool 注册与调用机制
2.1 BaseTool 抽象类结构
# langchain/tools/base.py
from pydantic import BaseModel, Field
class BaseTool(BaseModel):
name: str
description: str
args_schema: Optional[Type[BaseModel]] = None
class Config:
arbitrary_types_allowed = True
def run(self, tool_input: Union[str, Dict], **kwargs) -> Any:
return self._run(tool_input, **kwargs)
async def arun(self, tool_input: Union[str, Dict], **kwargs) -> Any:
return await self._arun(tool_input, **kwargs)
def _run(self, tool_input: str) -> str:
raise NotImplementedError
async def _arun(self, tool_input: str) -> str:
raise NotImplementedErrorBaseTool 是所有工具的基类,通过 Pydantic 进行字段校验。子类必须实现 _run(同步)和 _arun(异步)方法。
2.2 Tool 数据类 vs StructuredTool
LangChain 提供两种快速创建工具的方式:
| 类名 | 适用场景 | 参数灵活性 |
|---|---|---|
Tool | 简单工具,输入为单个字符串 | 固定,仅接受一个字符串参数 |
StructuredTool | 复杂工具,需要结构化参数 | 灵活,支持 Pydantic Schema 自动生成 |
from langchain.tools import Tool, StructuredTool
# Tool:函数式创建,输入为字符串
def search_func(query: str) -> str:
return f"Search results for: {query}"
search_tool = Tool(
name="Search",
func=search_func,
description="搜索工具,输入为查询字符串"
)
# StructuredTool:支持多参数
def booking_func(date: str, room: str) -> str:
return f"Booked {room} on {date}"
booking_tool = StructuredTool.from_function(
func=booking_func,
name="RoomBooking",
description="会议室预订,需要日期和房间名"
)2.3 Tool 参数 Schema 生成
StructuredTool 会自动从函数的类型注解和文档字符串生成 Pydantic Schema。args_schema 字段显式指定时,框架会使用该 Schema 校验输入参数:
from pydantic import BaseModel, Field
class BookingInput(BaseModel):
date: str = Field(description="预订日期,格式 YYYY-MM-DD")
room: str = Field(description="会议室名称")
duration: int = Field(default=1, description="使用时长(小时)")
booking_tool = StructuredTool.from_function(
func=booking_func,
name="RoomBooking",
description="会议室预订",
args_schema=BookingInput
)2.4 Tool 异常处理与重试
当 Tool 抛出异常时,AgentExecutor 会捕获并将其作为 observation 返回给 Agent,由 Agent 决定是否重试或更改策略。ToolException 提供更细粒度的错误分类:
from langchain.tools.base import ToolException
class RetryTool(BaseTool):
max_retries: int = 3
def _run(self, tool_input: str) -> str:
for attempt in range(self.max_retries):
try:
return self._execute(tool_input)
except ToolException as e:
if attempt == self.max_retries - 1:
raise
continue三、Memory 存储机制
3.1 类层次结构
BaseMemory
├── BaseChatMemory
│ ├── ConversationBufferMemory
│ ├── ConversationStringBufferMemory
│ └── ConversationSummaryMemory
└── VectorStoreMemory3.2 核心实现对比
| Memory 类型 | 存储方式 | 优点 | 缺点 |
|---|---|---|---|
ConversationBufferMemory | Python list 存储对话历史 | 实现简单,无损 | 随对话增长占用大量 Token |
ConversationSummaryMemory | LLM 总结压缩历史 | Token 消耗稳定 | 有信息损失,增加一次 LLM 调用 |
VectorStoreMemory | 向量数据库 + 语义检索 | 支持长期记忆,相关度检索 | 需要嵌入模型,存储成本高 |
3.3 ConversationBufferMemory 底层实现
# langchain/memory/buffer.py
class ConversationBufferMemory(BaseChatMemory):
chat_memory: BaseChatMessageHistory = Field(default_factory=ChatMessageHistory)
human_prefix: str = "Human"
ai_prefix: str = "AI"
@property
def buffer(self) -> str:
# 将 message 列表拼接为字符串
string_messages = []
for msg in self.chat_memory.messages:
if isinstance(msg, HumanMessage):
string_messages.append(f"{self.human_prefix}: {msg.content}")
elif isinstance(msg, AIMessage):
string_messages.append(f"{self.ai_prefix}: {msg.content}")
return "\n".join(string_messages)
def load_memory_variables(self, inputs: Dict[str, Any]) -> Dict[str, Any]:
return {self.memory_key: self.buffer}
def save_context(self, inputs: Dict[str, str], outputs: Dict[str, str]) -> None:
self.chat_memory.add_user_message(inputs[self.input_key])
self.chat_memory.add_ai_message(outputs[self.output_key])底层使用 ChatMessageHistory 的 messages: List[BaseMessage] 存储所有对话记录。每次 save_context 时追加 Human 和 AI 消息,load_memory_variables 时拼接为完整字符串注入 Prompt。
3.4 ConversationSummaryMemory 的 LLM 总结压缩
# langchain/memory/summary.py
class ConversationSummaryMemory(BaseChatMemory):
llm: BaseLLM
buffer: str = ""
prompt: BasePromptTemplate = SummaryPrompt()
def predict_new_summary(self, messages: List[BaseMessage], existing_summary: str) -> str:
# 调用 LLM 对当前摘要 + 新消息进行压缩总结
return self.llm.predict(self.prompt.format(
summary=existing_summary,
new_lines=self._format_messages(messages)
))
def save_context(self, inputs: Dict[str, str], outputs: Dict[str, str]) -> None:
super().save_context(inputs, outputs)
# 触发 LLM 重新总结
self.buffer = self.predict_new_summary(
self.chat_memory.messages, self.buffer
)每次新对话到来时,将现有摘要和新消息一起发送给 LLM,生成更短的总结存入 buffer。
3.5 VectorStoreMemory 的向量检索记忆
# langchain/memory/vectorstore.py
class VectorStoreMemory(BaseMemory):
vectorstore: VectorStore
retriever: BaseRetriever
memory_key: str = "memory"
def load_memory_variables(self, inputs: Dict[str, Any]) -> Dict[str, Any]:
# 从向量库中检索与当前输入相关的记忆片段
docs = self.retriever.get_relevant_documents(inputs.get("input", ""))
return {self.memory_key: "\n".join(d.page_content for d in docs)}
def save_context(self, inputs: Dict[str, str], outputs: Dict[str, str]) -> None:
# 将对话内容存入向量库
text = f"{inputs[self.input_key]}: {outputs[self.output_key]}"
self.vectorstore.add_texts([text])通过向量检索实现语义相关记忆,适合需要长期跨会话记忆的场景。
四、Callback 事件系统源码
4.1 类层次结构
BaseCallbackManager
├── CallbackManager (继承 BaseRunManager)
└── CallbackManagerForChainRun (RunManager 子类)
BaseCallbackHandler
├── StdOutCallbackHandler (控制台输出)
├── FileCallbackHandler (文件输出)
└── AsyncCallbackHandler (异步回调基类)4.2 BaseCallbackHandler 与 BaseCallbackManager
# langchain/callbacks/base.py
class BaseCallbackHandler:
"""回调处理器基类,定义各阶段事件钩子"""
def on_llm_start(self, serialized: Dict, prompts: List[str], **kwargs) -> None:
pass
def on_llm_end(self, response: LLMResult, **kwargs) -> None:
pass
def on_chain_start(self, serialized: Dict, inputs: Dict, **kwargs) -> None:
pass
def on_chain_end(self, outputs: Dict, **kwargs) -> None:
pass
def on_tool_start(self, serialized: Dict, input_str: str, **kwargs) -> None:
pass
def on_tool_end(self, output: str, **kwargs) -> None:
pass
def on_agent_action(self, action: AgentAction, **kwargs) -> None:
pass
class BaseCallbackManager:
handlers: List[BaseCallbackHandler]
inheritable_handlers: List[BaseCallbackHandler]
parent: Optional[BaseCallbackManager] = None
def add_handler(self, handler: BaseCallbackHandler, inheritable: bool = False) -> None:
if inheritable:
self.inheritable_handlers.append(handler)
else:
self.handlers.append(handler)每个钩子方法对应 LangChain 执行流程中的一个关键节点:LLM 调用开始/结束、Chain 开始/结束、Tool 调用开始/结束、Agent 执行动作等。
4.3 RunManager 的事件分发
# langchain/callbacks/manager.py
class RunManager:
"""运行时管理器,负责将事件分发给所有注册的 handler"""
handlers: List[BaseCallbackHandler]
parent_run_id: Optional[UUID]
run_id: UUID
def on_chain_start(self, serialized: Dict, inputs: Dict, **kwargs) -> None:
for handler in self.handlers:
handler.on_chain_start(serialized, inputs, run_id=self.run_id, **kwargs)
def on_llm_start(self, serialized: Dict, prompts: List[str], **kwargs) -> None:
for handler in self.handlers:
handler.on_llm_start(serialized, prompts, run_id=self.run_id, **kwargs)事件分发采用同步遍历模式,按 handler 注册顺序依次调用。每个 RunManager 实例持有独立的 run_id,用于追踪单次执行的完整链路。
4.4 inheritable tags / metadata 传递
Callback 系统支持 tags 和 metadata 的父子继承传递。当父 Chain 设置了 inheritable 的 tags 时,子 Chain、LLM、Tool 会自动继承这些标签:
# 使用示例
chain = SomeChain(
callbacks=CallbackManager([
StdOutCallbackHandler()
]),
tags=["experiment-a"],
metadata={"version": "2.1"}
)子组件通过以下逻辑继承:
def _inherit_tags(
parent_tags: List[str],
parent_inheritable_tags: List[str],
child_tags: List[str],
) -> List[str]:
"""合并父级可继承标签与子级标签"""
return list(set(parent_inheritable_tags + child_tags))4.5 FileCallbackHandler / StdOutCallbackHandler
| Handler 类 | 输出目标 | 典型用途 |
|---|---|---|
StdOutCallbackHandler | 标准输出流 | 调试时实时观察执行过程 |
FileCallbackHandler | 文件 | 记录长周期任务的详细日志 |
AsyncCallbackHandler | 异步处理 | 自定义异步回调(如写入消息队列) |
from langchain.callbacks import StdOutCallbackHandler, FileCallbackHandler
stdout_handler = StdOutCallbackHandler()
file_handler = FileCallbackHandler("agent_trace.log")
agent_executor = AgentExecutor(
agent=agent,
tools=tools,
callbacks=[stdout_handler, file_handler],
verbose=True
)4.6 Callback 事件生命周期
一次完整 Agent 执行的事件触发顺序:
on_chain_start (AgentExecutor)
└── on_agent_action (Agent 输出 AgentAction)
├── on_tool_start (调用具体 Tool)
├── on_tool_end (Tool 返回结果)
└── on_llm_start (Agent 再次调用 LLM 推理)
└── on_llm_end (LLM 返回)
└── on_agent_finish (Agent 输出 AgentFinish)
└── on_chain_end (AgentExecutor 结束)五、关键类继承关系图
5.1 Agent 体系
Chain (langchain/chains/base.py)
└── AgentExecutor (编排循环, 持有 Agent + Tools + Memory + Callbacks)
BaseSingleActionAgent
├── LLMSingleActionAgent (使用 LLM + Prompt 生成动作)
└── ZeroShotAgent (ReAct 模式 Agent, 基于 Few-shot Prompt)5.2 Tool 体系
BaseModel (pydantic)
└── BaseTool (抽象工具基类, 定义 run / arun 接口)
├── Tool (单个字符串输入, func 属性绑定函数)
├── StructuredTool (Schema 校验, 支持多参数)
└── ToolException (工具执行异常, 触发重试机制)5.3 Memory 体系
BaseMemory (内存基类)
├── BaseChatMemory (对话记忆基类, 持有一个 ChatMessageHistory)
│ ├── ConversationBufferMemory (完整列表存储)
│ ├── ConversationStringBufferMemory (字符串缓冲)
│ └── ConversationSummaryMemory (LLM 总结压缩)
└── VectorStoreMemory (向量检索记忆)5.4 Callback 体系
BaseCallbackHandler (事件钩子接口)
├── StdOutCallbackHandler (控制台日志)
├── FileCallbackHandler (文件日志)
└── AsyncCallbackHandler (异步钩子接口)
BaseCallbackManager (管理器)
└── CallbackManager (持有 handlers 列表, 实现 RunManager)
└── CallbackManagerForChainRun (Chain 级别 RunManager)
└── CallbackManagerForToolRun (Tool 级别 RunManager)
└── CallbackManagerForLLMRun (LLM 级别 RunManager)总结
LangChain 的四个核心模块协同工作流程如下:AgentExecutor 作为编排引擎驱动主循环,通过 Agent.plan() 决策下一步动作;动作涉及 Tool 调用时,通过 BaseTool.run() 执行并处理异常;对话上下文由 BaseMemory 体系持久化存储并在每次迭代前注入 Prompt;整个执行过程通过 Callback 系统向各 Handler 分发事件,实现可观测性和日志记录。理解上述源码结构是深入使用和定制 LangChain 的基础。