LangChain 模型层深度剖析:抽象、并发、容错与本地推理¶
LangChain 的模型层是其与 LLM 交互的核心。理解 BaseLLM 与 BaseChatModel 的本质区别,掌握如何封装自定义模型、控制输出、实现并发调用和容错回退,是在生产环境中高效使用 LangChain 的关键。以下逐一深入。
LangChain 中 BaseLLM 和 BaseChatModel 的根本区别是什么?什么时候该用哪个?¶
LangChain 将语言模型分为两大抽象:BaseLLM 和 BaseChatModel。它们底层都继承自 BaseLanguageModel,但面向的交互范式截然不同。
📜 BaseLLM:面向传统的“文本补全”模型(如 GPT-3、text-davinci-003、早期的 Cohere Generate API)。
-
输入和输出都是纯字符串。
-
你的 Prompt 直接是一个字符串,模型返回一个字符串。
-
没有内置的对话角色管理(如 system、user、assistant 分离)。
-
典型使用场景:单轮文本生成(摘要、翻译、代码补全)、批量处理无上下文的任务。
💬 BaseChatModel:面向现代的“对话”模型(如 GPT-3.5-Turbo、GPT-4、Claude、Llama 3)。
-
输入是一个
List[BaseMessage],包含SystemMessage、HumanMessage、AIMessage等。 -
输出是一个
BaseMessage(通常是AIMessage),支持content和tool_calls等字段。 -
原生支持多轮对话、角色扮演、工具调用(Function Calling)。
-
典型使用场景:聊天机器人、Agent、任何需要多轮交互的复杂应用。
🧭 什么时候该用哪个?
-
如果你的任务很简单,且模型只支持文本补全 API(如某些开源模型),可以使用
BaseLLM。 -
几乎所有现代 LLM 都提供 Chat API,因此新项目应默认使用
BaseChatModel。 -
即使你的任务是单轮问答,也可以使用
BaseChatModel,在 Prompt 中只放一条HumanMessage,效果通常比文本补全更好,因为对话模型经过了指令微调。 -
BaseLLM在新版 LangChain 中逐渐边缘化,官方重心在BaseChatModel上。
💡 设计差异的本质:BaseChatModel 是 BaseLLM 的超集,它通过消息结构携带了更丰富的语义信息(角色、工具调用),使得模型能理解更复杂的指令。LangChain 选择将两者分开,是为了向后兼容历史模型,同时为新一代模型提供更清晰的接口。
如何封装一个自定义的 LLM?需要重写哪些核心方法?¶
当你想接入一个 LangChain 官方未集成的模型(例如自研模型、内部服务),可以继承 BaseLLM 或 BaseChatModel 来实现自定义适配器。
🛠️ 以封装自定义 ChatModel 为例,你需要至少重写以下方法:
_generate或_agenerate(底层生成逻辑)- 输入:
messages(List[BaseMessage])、stop(停止词列表)、run_manager(回调管理器)等。 - 输出:
ChatResult对象,包含generations(List[ChatGeneration]),每个ChatGeneration包含message(AIMessage)和可选的generation_info。 -
你需要在此方法中实现模型的实际调用(HTTP 请求、本地推理等),并将返回的文本包装为
AIMessage。 -
_stream或_astream(流式生成逻辑,可选但推荐) - 返回一个生成器,每次 yield 一个
ChatGenerationChunk,包含message(AIMessageChunk)。 -
支持流式输出可以显著提升用户体验。
-
_llm_type属性(返回一个字符串,标识模型类型) -
返回自定义的名称,如
"my-custom-llm"。 -
_identifying_params属性(可选,用于日志和追踪) - 返回一个字典,包含模型的标识参数。
📜 示例伪代码:
from langchain_core.language_models.chat_models import BaseChatModel
from langchain_core.messages import AIMessage, BaseMessage
from langchain_core.outputs import ChatResult, ChatGeneration
from typing import List, Optional, Any
class MyCustomChatModel(BaseChatModel):
model_name: str = "my-model"
api_url: str = "http://localhost:8000/chat"
def _generate(self, messages: List[BaseMessage], stop: Optional[List[str]] = None, run_manager=None, **kwargs) -> ChatResult:
# 1. 将 LangChain 消息列表转换为 API 所需的 JSON 格式
payload = self._convert_messages(messages)
# 2. 调用模型 API
import requests
response = requests.post(self.api_url, json=payload)
response.raise_for_status()
text = response.json()["content"]
# 3. 包装为 AIMessage
message = AIMessage(content=text)
generation = ChatGeneration(message=message)
return ChatResult(generations=[generation])
def _stream(self, messages, stop=None, run_manager=None, **kwargs):
# 流式实现略
pass
@property
def _llm_type(self):
return "my-custom-chat-model"
🧪 测试与集成:
-
实现后,你的自定义类可以直接在 LCEL 中使用:
chain = prompt | MyCustomChatModel() | parser。 -
确保处理异常(网络超时、模型错误),考虑重试和回退。
-
如果要支持异步,重写
_agenerate和_astream方法。
解释 invoke、ainvoke、stream、astream 四个方法的调用时机和返回值差异。¶
这四个方法是 Runnable 接口的核心,定义了与组件的交互方式。它们都基于底层的 _generate 和 _stream 实现。
| 方法 | 调用时机 | 输入 | 返回值 | 适用场景 |
|---|---|---|---|---|
| invoke | 同步单次调用 | 消息列表 List[BaseMessage] | AIMessage | 不需要流式,请求-响应模式 |
| ainvoke | 异步单次调用 | 同上 | AIMessage | 异步环境(FastAPI、asyncio),不阻塞事件循环 |
| stream | 同步流式调用 | 同上 | Iterator[AIMessageChunk] | 需要逐 Token 实时展示(聊天界面) |
| astream | 异步流式调用 | 同上 | AsyncIterator[AIMessageChunk] | 异步流式,高并发服务端推送 |
🔍 细节差异:
-
stream/astream返回的是 AIMessageChunk 的迭代器。每个 Chunk 包含部分内容,最后可能需要通过StrOutputParser等解析器处理,也可以直接累加。 -
在 LCEL 中,如果链中任何组件支持
stream,且你调用链的stream方法,整个链会自动变为流式,数据逐级传递。 -
ainvoke需要 Python 的asyncio环境。如果你的应用基于 FastAPI,使用异步方法可以避免阻塞其他请求,提升并发。
在 LangChain 中,如何根据模型是否支持 Function Calling 来动态切换 Agent 策略?¶
不同的模型对 Function Calling 的支持程度不同。例如 GPT-4 原生支持,而一些开源模型需要用 ReAct 风格的 Prompt 模拟工具调用。LangChain 提供了灵活的机制来动态选择策略。
🎯 方案一:使用 bind_tools() 判断
对于支持 Function Calling 的模型,调用 bind_tools() 会生成正确的 API 参数。对于不支持的模型,bind_tools() 可能只是将工具定义转换为 Prompt 文本。你可以通过检查模型的 supports_tool_calls 属性来判断。
📜 示例:
from langchain_openai import ChatOpenAI
from langchain_ollama import ChatOllama
def create_agent(model, tools):
if hasattr(model, "bind_tools"): # 支持 Function Calling
model_with_tools = model.bind_tools(tools)
# 使用 OpenAI Functions Agent
return create_openai_functions_agent(model_with_tools, tools)
else:
# 使用 ReAct Agent(通过 Prompt 模拟工具调用)
prompt = hub.pull("hwchase17/react")
return create_react_agent(model, tools, prompt)
🔀 方案二:使用 RunnableBranch 动态路由
在运行时根据模型类型选择不同的 Agent 执行路径:
from langchain_core.runnables import RunnableBranch
def is_openai_model(model):
return isinstance(model, ChatOpenAI)
branch = RunnableBranch(
(lambda x: is_openai_model(x["model"]), openai_agent_chain),
react_agent_chain # 默认
)
这样,当部署环境切换模型时,Agent 策略自动适配,无需手动修改代码。
怎样在 LangChain 中控制 LLM 的输出 token 数量?有哪些参数可以设置?¶
控制输出 token 数量对于成本控制和响应长度管理至关重要。LangChain 提供了多种方式:
🎛️ 核心参数:
-
max_tokens(或max_new_tokens):限制生成的最大 token 数。在ChatOpenAI等模型中,可直接在初始化时传入。 -
temperature:控制输出的随机性。较低的温度(0.0~0.3)使输出更确定,token 使用更高效;较高的温度(0.7~1.0)使输出更多样,可能消耗更多 token。 -
top_p:核采样,与temperature配合使用,控制输出的多样性。
📜 示例:
⚙️ 在 Prompt 中约束:
你可以在 System Prompt 中明确要求模型“用不超过 50 个字回答”或“用一句话总结”。虽然模型不一定严格遵守,但能显著减少输出长度。
📊 运行时监控:
通过回调或 LangSmith 追踪实际的 Token 消耗,如果发现超支,可以在应用层截断、替换为更短的模型响应,或者切换至更便宜的模型。
如果模型 API 返回了超出 max_tokens 限制的错误,你在 LangChain 中如何处理和重试?¶
当模型返回类似 context_length_exceeded 或 maximum context length 的错误时,意味着输入 Token 超出模型上下文窗口。LangChain 没有内置自动处理,需要手动实现重试与降级策略。
🛠️ 处理流程:
-
捕获异常:在自定义的
Callback或RunnableLambda中捕获OpenAIError(或相应提供商的异常)。 -
重试机制:使用
with_fallbacks()或手动实现try/except。重试时,可以采取以下降级措施: - 自动截断历史消息(保留最近 K 轮)。
- 切换至更大的上下文窗口模型(如
gpt-4-128k)。 -
压缩或摘要历史对话(使用
ConversationSummaryMemory生成摘要,替换原始历史)。 -
使用 LangChain 的
trim_messages工具:在 Prompt 模板中,可以使用MessagesPlaceholder配合trim_messages辅助函数,自动裁剪超出窗口的消息。
📜 示例:
from langchain_core.runnables import RunnableLambda
def handle_context_error(input, config):
try:
return model.invoke(input)
except ContextWindowExceededError:
# 裁减消息:只保留最近 5 轮
trimmed = input[-10:] # 假设每轮2条消息
return model.invoke(trimmed)
safe_chain = prompt | RunnableLambda(handle_context_error) | parser
🧠 更好的方案:在构建 Prompt 时,使用 trim_messages 限制总 Token 数,防患于未然:
from langchain_core.messages import trim_messages
trimmed = trim_messages(messages, max_tokens=3000, token_counter=model.get_num_tokens)
谈谈 LangChain 对“模型回退”(fallback)的设计,例如 GPT-4 失败后切换 GPT-3.5。¶
🛡️ LCEL 的 with_fallbacks 是实现模型回退的核心机制。它允许为一个 Runnable 指定一个或多个备用 Runnable,当前者失败时自动尝试备用。
📜 基础用法:
gpt4 = ChatOpenAI(model="gpt-4o")
gpt35 = ChatOpenAI(model="gpt-3.5-turbo")
robust_model = gpt4.with_fallbacks([gpt35])
chain = prompt | robust_model | parser
⚙️ 执行逻辑:
-
链首先使用
gpt4。 -
如果
gpt4抛出任何异常(如openai.APIError、openai.APITimeoutError、openai.RateLimitError),LangChain 会自动捕获并尝试gpt35。 -
如果
gpt35也失败,异常向上传播。
🔧 高级用法:
-
可以为不同的异常类型指定不同的备用模型(例如,
RateLimitError切换至另一个 API Key 的gpt4,而不是降级模型)。 -
可以嵌套回退,形成链条:
gpt4.with_fallbacks([gpt35.with_fallbacks([local_model])])。 -
结合回调,在回退发生时发送告警通知。
🌩️ 局限性:
-
with_fallbacks不重试原模型,直接切换。若需要同一模型的重试(如网络抖动),应在模型初始化时配置max_retries。 -
回退时不会自动调整 Prompt,如果失败是上下文长度导致,降级模型可能仍会失败。需要配合上下文截断策略。
在 LangChain 中,如何使用本地 HuggingFace 模型进行推理?与 API 模型有何不同?¶
🖥️ 使用本地模型的核心优势:数据隐私、零 API 成本、离线可用、可完全定制。
📦 集成方式:
- 通过
HuggingFacePipeline(适合本地服务器)
from transformers import pipeline, AutoModelForCausalLM, AutoTokenizer
from langchain_huggingface import HuggingFacePipeline
model = AutoModelForCausalLM.from_pretrained("gpt2")
tokenizer = AutoTokenizer.from_pretrained("gpt2")
pipe = pipeline("text-generation", model=model, tokenizer=tokenizer, max_new_tokens=256)
llm = HuggingFacePipeline(pipeline=pipe)
chain = llm | StrOutputParser()
-
通过
HuggingFaceEndpoint(适合 HuggingFace Hub 上的推理端点) 不需要本地 GPU,调用 HuggingFace 的在线推理服务。 -
通过 Ollama 或 vLLM 的 OpenAI 兼容接口 启动 vLLM 后,使用
ChatOpenAI并指定base_url指向本地服务,完全兼容 OpenAI 的调用方式。
🆚 与 API 模型的差异:
| 维度 | 本地模型 | API 模型 |
|---|---|---|
| 延迟 | 取决于硬件,可能较高 | 通常更低(云端GPU集群) |
| 并发能力 | 受限于本地资源 | 弹性扩展 |
| 功能支持 | 部分模型不支持 Function Calling | 通常完全支持 |
| 运维成本 | 需要自己管理硬件和模型更新 | 零运维 |
| 隐私 | 数据不出域 | 数据发送至第三方 |
⚖️ 建议:开发调试阶段用 API 模型快速迭代,对隐私敏感或需要高频调用的场景,部署本地模型进行推理。可以结合 LangChain 的回退机制,在本地模型忙时自动切换云端 API。
写一段代码,演示如何用 LangChain 并发调用多个 LLM,然后汇总结果。¶
在 LCEL 中,并发调用多个 LLM 主要依赖 RunnableParallel。它接受一个字典,键是输出字段名,值是一个 Runnable,框架会并发执行所有分支,最后将结果合并为一个字典。
📜 示例:同时让 GPT-4 和 Claude 回答同一个问题,然后用另一个模型汇总:
from langchain_core.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI
from langchain_anthropic import ChatAnthropic
from langchain_core.runnables import RunnableParallel
# 定义用于回答的 Prompt
answer_prompt = ChatPromptTemplate.from_template("问题:{question}\n请用中文简洁回答。")
# 模型1
gpt4 = ChatOpenAI(model="gpt-4o", temperature=0)
# 模型2
claude = ChatAnthropic(model="claude-3-sonnet", temperature=0)
# 并发调用链:两个模型独立回答
parallel_answers = RunnableParallel(
gpt_answer=answer_prompt | gpt4,
claude_answer=answer_prompt | claude
)
# 汇总 Prompt
summary_prompt = ChatPromptTemplate.from_template("""
你是一个评审专家。请综合以下两个回答,给出一个更全面、准确的最终答案。
GPT-4 的回答:{gpt_answer}
Claude 的回答:{claude_answer}
最终答案:
""")
# 汇总模型
summarizer = ChatOpenAI(model="gpt-4o", temperature=0)
# 完整链:并发获取两个答案 → 汇总
full_chain = parallel_answers | summary_prompt | summarizer
# 调用
result = full_chain.invoke({"question": "什么是量子计算?"})
print(result.content)
⚡ 执行细节:
-
RunnableParallel会异步执行gpt_answer和claude_answer分支,总耗时约等于最慢的单个模型调用时间。 -
如果你需要流式输出最终结果,可以调用
full_chain.stream(),框架会自动处理并发分支的流式聚合。 -
对于更复杂的并发场景(如并行调用多个工具),可以使用
RunnableMap或asyncio.gather结合ainvoke。
🎛️ 扩展:动态并发模型列表
你可以根据运行时配置动态生成并行分支:
models = {"gpt4": gpt4, "claude": claude, "llama": local_model}
branches = {name: answer_prompt | model for name, model in models.items()}
parallel = RunnableParallel(**branches)
这种模式在“多模型投票”、“集成学习”等场景中非常实用,可以显著提升答案的准确性和鲁棒性。
解释 LangChain 中 model_name 和 model_kwargs 的作用及它们最终如何传递给底层 API。¶
🏷️ model_name 是 BaseChatModel 或 BaseLLM 的必填参数,用于指定底层具体的模型标识。例如 "gpt-4o"、"claude-3-opus"。当构建 ChatOpenAI(model="gpt-4o") 时,这个字符串会被原样传递给 OpenAI SDK 的 model 参数,告诉 API 使用哪个模型版本。对于 HuggingFace 本地模型,model_name 则对应 Hub 上的仓库 ID(如 "meta-llama/Llama-2-7b-chat-hf")。
⚙️ model_kwargs 是一个字典,用于传递超出 LangChain 标准参数之外的额外参数。LangChain 的标准参数(如 temperature、max_tokens、top_p)会被显式声明为 Pydantic 字段,而任何底层 API 支持但 LangChain 未封装的参数,都可以通过 model_kwargs 传递。
例如,OpenAI 的 logit_bias 可以控制特定 token 的出现概率,LangChain 并未提供直接的 Pydantic 字段。此时你可以:
model = ChatOpenAI(
model="gpt-4",
temperature=0.3,
model_kwargs={"logit_bias": {1234: 100}, "user": "abc"}
)
在 ChatOpenAI 内部,这些 model_kwargs 最终会被合并到发给 OpenAI API 的请求体 **kwargs 中。
📦 传递机制:LangChain 在调用底层 SDK 时,会将 Pydantic 字段(如 temperature)和 model_kwargs 合并为一个字典,并过滤掉 None 值,然后解包传给 SDK 的 create 方法。例如,OpenAI 的调用最终等价于:
client.chat.completions.create(
model="gpt-4",
temperature=0.3,
logit_bias={1234: 100},
user="abc",
messages=[...]
)
⚠️ 注意事项:
-
model_kwargs中的参数会覆盖同名的标准参数,可能导致意外行为,需谨慎使用。 -
不同提供商的 SDK 支持的额外参数不同,使用前需查阅对应文档。
-
如果你需要动态改变
model_name或model_kwargs,可以在运行时通过bind()方法或重新创建模型实例实现。
如何在 LangChain 中实现模型的动态选择(如根据输入内容复杂度选择不同规模的模型)?¶
🧠 动态模型选择可以根据任务复杂度、用户等级、成本预算等条件,自动切换不同能力/价格的模型。例如:简单问题用 gpt-3.5-turbo,复杂推理用 gpt-4o。
🔀 实现方案一:使用 RunnableBranch 进行条件路由
RunnableBranch 根据输入条件选择不同的 Runnable 分支。它接受一个 (condition, runnable) 对的列表,顺序评估,首个满足条件的路径被执行。
from langchain_core.runnables import RunnableBranch, RunnableLambda
def is_complex(input_dict):
# 判断逻辑:问题长度、关键词、历史复杂度评分等
return len(input_dict["question"]) > 100 or "复杂" in input_dict["question"]
gpt35 = ChatOpenAI(model="gpt-3.5-turbo")
gpt4 = ChatOpenAI(model="gpt-4o")
branch = RunnableBranch(
(lambda x: is_complex(x), prompt | gpt4 | parser),
prompt | gpt35 | parser # 默认简单分支
)
result = branch.invoke({"question": "..."})
🔀 实现方案二:使用 RunnableLambda 自定义函数
更灵活的方式是在自定义函数中根据输入选择模型:
def choose_model(input_dict):
if is_complex(input_dict):
model = gpt4
else:
model = gpt35
chain = prompt | model | parser
return chain.invoke(input_dict)
dynamic_chain = RunnableLambda(choose_model)
🔀 实现方案三:结合 Router Chain(旧式)或 LangGraph
在更复杂的多步 Agent 中,可以由一个“分类器” LLM 先判断任务类型,然后路由到不同的链。这可以通过 LangGraph 构建一个状态机实现。
⚙️ 性能与成本优化:
-
可以将模型实例缓存起来(例如通过
functools.lru_cache),避免每次创建新实例。 -
结合
with_fallbacks,当复杂模型不可用时自动降级到简单模型。 -
在回调中记录每次决策的依据和模型选择,便于后续分析成本与效果。
LangChain 的 Chat Model 如何管理不同角色的消息(system, human, ai, function)?用代码说明。¶
💬 LangChain 使用 BaseMessage 的子类来区分消息角色,从而构建结构化的对话历史。这些消息类型构成了 ChatPromptTemplate 和模型输入的基石。
📋 核心消息类型:
| 消息类型 | 角色 | 用途 | 示例 |
|---|---|---|---|
| SystemMessage | system | 设定模型行为、角色、规则 | “你是一个有帮助的助手” |
| HumanMessage | user | 用户的输入 | “今天天气如何?” |
| AIMessage | assistant | 模型的回复 | “今天晴天,25℃” |
| FunctionMessage | function | 工具执行的结果 | “查询结果:25℃” |
| ToolMessage | tool | 工具调用的结果(新版) | 与 FunctionMessage 类似,但绑定 tool_call_id |
📜 在 Prompt 模板中使用:
from langchain_core.prompts import ChatPromptTemplate, MessagesPlaceholder
prompt = ChatPromptTemplate.from_messages([
("system", "你是一个{role},用{language}回答。"),
MessagesPlaceholder(variable_name="history"), # 历史消息占位
("human", "{question}")
])
# 构建消息列表
from langchain_core.messages import SystemMessage, HumanMessage, AIMessage
messages = prompt.format_messages(
role="助手", language="中文", question="你好",
history=[
HumanMessage(content="我叫小明"),
AIMessage(content="你好小明,有什么可以帮你?")
]
)
# messages 是一个 List[BaseMessage],可直接传给 model.invoke(messages)
🔄 模型如何处理:
-
ChatOpenAI等内部会将SystemMessage映射为role: "system"的字典。 -
HumanMessage→role: "user"。 -
AIMessage→role: "assistant"(如果包含tool_calls,还会附带tool_calls字段)。 -
ToolMessage→role: "tool",并附加tool_call_id。
这种设计使得对话历史可以跨模型提供商保持统一,开发者无需关心底层 API 的消息格式差异。
🔧 动态管理消息:
-
使用
trim_messages可以自动裁剪历史消息,以适配模型的上下文窗口。 -
MessagesPlaceholder允许在 Prompt 模板中灵活插入可变长度的消息列表。
当使用流式输出时,LangChain 是如何把 token 逐个回调给调用方的?底层用了什么机制?¶
🌊 流式输出的核心机制是 Python 的生成器(Generator)和异步生成器(AsyncGenerator)。
🔧 底层实现原理:
-
模型层的
_stream方法:每个 ChatModel 可以实现_stream(同步)或_astream(异步)。这个方法返回一个Iterator[ChatGenerationChunk],每个 Chunk 包含一个AIMessageChunk(内容为增量文本)。 -
LCEL 的自动流式传播:当你在 LCEL 链上调用
stream()时,框架会检查链中每个组件是否支持流式。如果上游是流式的,下游组件(如StrOutputParser)会逐个处理每个 Chunk,而不是等待完整响应。 -
从 SDK 获取流式数据:以 OpenAI 为例,
ChatOpenAI._astream内部调用client.chat.completions.create(stream=True),这会返回一个事件流。框架遍历事件流,每次收到choices[0].delta.content不为空时,就 yield 一个新的AIMessageChunk。 -
回调(Callbacks)的参与:在流式过程中,回调系统会触发
on_llm_new_token事件,开发者可以在这里实现实时日志、Token 计数、前端推送等。
📜 使用示例:
chain = prompt | model | StrOutputParser()
# 流式调用
for chunk in chain.stream({"question": "介绍量子计算"}):
print(chunk, end="", flush=True) # 逐 token 打印
⚙️ 底层伪代码(简化):
class ChatOpenAI(BaseChatModel):
async def _astream(self, messages, **kwargs):
response = await client.chat.completions.create(
model=self.model_name, messages=messages, stream=True, **kwargs
)
async for chunk in response:
if chunk.choices[0].delta.content:
yield ChatGenerationChunk(
message=AIMessageChunk(content=chunk.choices[0].delta.content)
)
📊 回调机制:
在 _astream 中,每次 yield 之前,LangChain 会触发 run_manager.on_llm_new_token(token, chunk=...)。如果你注册了自定义回调,就能拿到每个 token 及其上下文,用于构建打字机效果、审计日志等。
如果你要封装一个支持 System Prompt 的自定义 Chat Model,应该继承哪个基类?¶
🧬 应该继承 BaseChatModel。
理由:
-
BaseChatModel原生支持消息列表作为输入,其中SystemMessage就是消息列表的一部分。 -
它内部处理了消息格式转换、工具调用绑定、流式支持等,开发者只需实现
_generate和(可选的)_stream。 -
不需要去处理
SystemMessage的特殊逻辑,因为它只是messages参数中的一个元素。你只需将整个消息列表转发给底层模型 API 即可。
📜 示例:
from langchain_core.language_models.chat_models import BaseChatModel
from langchain_core.messages import AIMessage
from langchain_core.outputs import ChatResult, ChatGeneration
class MySystemPromptModel(BaseChatModel):
model_name: str = "my-model"
api_url: str = "http://localhost:8000/chat"
def _generate(self, messages, stop=None, run_manager=None, **kwargs):
# messages 中已包含 SystemMessage, HumanMessage 等,直接转发
payload = self._convert_messages(messages)
response = requests.post(self.api_url, json=payload)
text = response.json()["content"]
return ChatResult(generations=[ChatGeneration(message=AIMessage(content=text))])
@property
def _llm_type(self):
return "my-system-prompt-model"
💡 注意事项:
-
如果你的底层 API 不支持 System Prompt 或需要将其放在特定位置,你可以在
_generate内部对messages进行重排或合并。 -
如果模型需要将 System Prompt 作为单独参数(如某些开源模型的
system_prompt字段),你可以覆盖_generate,从messages中提取SystemMessage,然后单独处理。
讨论一下 LangChain 模型调用中的 Token 计数:如何获取每次请求的实际 token 消耗?¶
💰 Token 计数是成本控制和配额管理的核心指标。LangChain 提供了多种方式获取每次 LLM 调用的实际 Token 消耗。
📊 方式一:通过 LLMResult 或 AIMessage.response_metadata
在模型返回的结果中,底层 API 通常包含 token_usage 字段。LangChain 会将其提取到 AIMessage 的 response_metadata 中。
response = model.invoke(messages)
print(response.response_metadata.get("token_usage"))
# 输出: {'prompt_tokens': 150, 'completion_tokens': 80, 'total_tokens': 230}
对于非聊天模型,LLMResult.llm_output 中也包含 token_usage。
📊 方式二:通过回调(Callback)记录
自定义一个 BaseCallbackHandler,在 on_llm_end 中提取 Token 用量并记录到日志或数据库。
from langchain_core.callbacks import BaseCallbackHandler
class TokenCounterCallback(BaseCallbackHandler):
def on_llm_end(self, response, **kwargs):
usage = response.llm_output.get("token_usage", {})
print(f"本次调用消耗 Token: {usage}")
📊 方式三:LangSmith 自动统计
当启用 LangSmith 追踪时,每次调用的 Token 用量和成本会被自动记录,可以在 Web 界面中查看聚合统计。
📊 方式四:使用 get_num_tokens 预测
对于需要提前限制输入长度的场景,可以使用模型的 get_num_tokens(text) 方法,调用底层 Tokenizer 估算 Token 数。但注意这仅是估算,实际消耗可能因格式包装而略有偏差。
⚠️ 注意:并非所有模型都返回 Token 用量(尤其是一些本地模型),此时需要自己实现 Token 计数逻辑。
在 LangChain 中,如何为模型调用设置超时时间?超过后如何处理?¶
⏱️ 设置超时可以防止因为网络延迟、模型过载等原因导致的无限等待。
🛠️ 方式一:在模型初始化时设置 request_timeout 参数(适用于 ChatOpenAI 等)
这个参数会传递给底层 HTTP 客户端(如 httpx 或 requests),当响应时间超过设定值时抛出 ReadTimeout 异常。
🛠️ 方式二:通过 model_kwargs 传递底层参数
某些提供商的超时参数可能不同,可以通过 model_kwargs 传递。
🛠️ 方式三:使用 with_fallbacks 结合超时异常处理
你可以自定义一个 Runnable,在其中捕获超时异常并尝试备用模型。
from langchain_core.runnables import RunnableLambda
def call_with_timeout(input):
try:
return model.invoke(input)
except (Timeout, ConnectionError):
return fallback_model.invoke(input)
robust_chain = prompt | RunnableLambda(call_with_timeout) | parser
🛠️ 方式四:异步环境中的超时控制
在异步代码中,可以使用 asyncio.wait_for 包裹 ainvoke,设置超时后抛出 asyncio.TimeoutError。
import asyncio
try:
result = await asyncio.wait_for(model.ainvoke(messages), timeout=10.0)
except asyncio.TimeoutError:
result = await fallback_model.ainvoke(messages)
⚙️ 超过超时后的处理建议:
-
重试(特别是网络抖动):使用
max_retries参数配合指数退避。 -
降级(fallback):切换至更快的模型或本地模型。
-
返回缓存结果:如果相同请求之前成功过,可以考虑返回缓存。
-
上报告警:记录超时次数,用于监控和优化。
你如何实现对某个 LLM 的请求重试(例如遇到 Rate Limit 时指数退避)?¶
🔄 指数退避(Exponential Backoff) 是处理临时性错误(如 429 Rate Limit)的标准策略。
🛠️ 方案一:利用底层 SDK 的重试机制
大多数 LLM SDK(如 openai)已经内置了重试逻辑。LangChain 的 ChatOpenAI 接受 max_retries 参数,它会传递给 OpenAI 客户端。
🛠️ 方案二:自定义重试 Runnable
你可以使用 tenacity 库创建一个带重试逻辑的 Runnable。
from tenacity import retry, stop_after_attempt, wait_exponential
from langchain_core.runnables import RunnableLambda
@retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=30))
def call_model_with_retry(input):
return model.invoke(input)
retry_chain = prompt | RunnableLambda(call_model_with_retry) | parser
🛠️ 方案三:使用 LangChain 的 with_retry 功能(部分版本提供)
LangChain 的 Runnable 提供了 with_retry 方法(较新版本),可以配置重试次数和退避策略。
robust_model = model.with_retry(
retry_if_exception_type=(openai.RateLimitError,),
stop_after_attempt=3,
wait_exponential_multiplier=1,
wait_exponential_max=30
)
📊 监控重试:在回调中记录重试次数,当重试频繁时发出告警,可能是配额不足或模型选择不当。
⚠️ 注意事项:
-
对于幂等的请求(如
temperature=0的相同 Prompt),重试是安全的。 -
对于非幂等请求,重试可能导致重复消费(如扣费多次),需要确保服务端的幂等性。
-
在
stream模式下,重试处理较复杂,通常建议在应用层手动处理。
在微服务中,你倾向于将 LangChain 模型调用放在同步还是异步环境?为什么?¶
⚡ 强烈推荐异步环境。原因如下:
-
高并发 I/O 密集型:LLM 调用本质上是网络 I/O 等待(等待 API 响应)。同步模式下,每个请求会占用一个线程,线程切换和内存开销大,且容易受限于线程池大小。异步模式下,协程在等待网络 I/O 时主动让出控制权,单线程即可处理上千并发连接,资源利用率极高。
-
与 FastAPI/Starlette 天然契合:现代 Python Web 框架几乎都基于
asyncio。使用异步 LangChain(ainvoke、astream)可以避免阻塞事件循环,最大化服务吞吐量。 -
流式输出:异步流式(
astream)能更高效地实现 Server-Sent Events (SSE) 推送,避免线程阻塞,支持大量客户端同时接收实时数据。
🧩 何时可能用同步?
-
简单的离线批处理脚本。
-
内部工具的短调用。
-
当使用的某些库不支持异步且无法替换时(但此时可以通过
asyncio.to_thread将同步代码桥接到异步环境)。
📦 微服务架构示例:
from fastapi import FastAPI
import asyncio
app = FastAPI()
@app.post("/ask")
async def ask_question(question: str):
# 异步调用链
result = await chain.ainvoke({"question": question})
return {"answer": result}
💡 最佳实践:在构建 LangChain 服务时,始终使用异步方法(ainvoke, astream),并配置 asyncio 的事件循环策略为 uvloop 以获得额外性能提升。
如何通过 LangChain 调用 Azure OpenAI Service?和直接调用 OpenAI 有何不同?¶
☁️ Azure OpenAI Service 是微软提供的企业级 OpenAI 模型服务,与原生 OpenAI API 的主要区别在于:
-
认证方式:使用 API Key 或 Azure Active Directory 认证。
-
端点定制:每个部署的资源有自己的基础 URL(
https://{your-resource-name}.``openai.azure.com)。 -
模型部署名称:模型通过“部署”来管理,API 调用时使用部署名而非模型名(如
gpt-4-deployment)。 -
API 版本控制:需要在请求中指定 API 版本(如
2024-02-01)。 -
内容过滤:默认启用 Azure 的内容安全过滤。
-
企业功能:支持 VNet/Private Link、托管身份、合规认证等。
🔌 在 LangChain 中使用 Azure OpenAI:
LangChain 提供了专用的 AzureChatOpenAI 类(在 langchain-openai 包中),封装了上述差异。
from langchain_openai import AzureChatOpenAI
model = AzureChatOpenAI(
azure_deployment="gpt-4-deployment", # 你的部署名
api_version="2024-02-01", # API 版本
azure_endpoint="https://my-resource.openai.azure.com", # 你的端点
api_key="YOUR_API_KEY", # 或者使用 azure_ad_token
temperature=0.3
)
或者使用环境变量:
export AZURE_OPENAI_API_KEY="..."
export AZURE_OPENAI_ENDPOINT="https://my-resource.openai.azure.com"
export OPENAI_API_VERSION="2024-02-01"
然后代码中仅需指定部署名:
🧩 与直接调用 OpenAI 的代码差异:
| 维度 | OpenAI | Azure OpenAI |
|---|---|---|
| 类 | ChatOpenAI | AzureChatOpenAI |
| 模型指定 | model="gpt-4" | azure_deployment="my-gpt-4" |
| 端点 | 固定 api.openai.com | 自定义 azure_endpoint |
| API 版本 | 无需指定 | 必须指定 api_version |
| 认证 | 仅 API Key | API Key 或 Azure AD Token |
| 内容安全 | 无 | 默认启用,可配置 |
⚙️ 注意事项:
-
Azure OpenAI 的速率限制是按部署和区域分配的,与 OpenAI 的全局限制不同。
-
当使用
with_fallbacks时,可以将 Azure 和 OpenAI 的实例组合,实现跨云容灾。 -
如果使用自托管模型,Azure 还提供吞吐量预配(PTU),适合生产环境。
如果让你设计一个多模型路由(根据语言、领域自动选择 LLM),用 LangChain 如何实现?¶
🌐 多模型路由的核心是根据输入特征(语言、领域、复杂度、成本要求等)将请求分发到最合适的 LLM。这可以通过 LangChain 的 RunnableBranch、RunnableLambda 或自定义路由逻辑实现。
🛠️ 设计步骤:
- 定义模型池:创建不同能力的模型实例。例如:
gpt4:通用复杂推理。claude:长文档处理。local_code_model:代码生成。-
chinese_llm:中文优化模型。 -
路由规则:编写判断函数,根据输入内容决定路由。例如:
- 检测语言:
is_chinese(text)→chinese_llm - 检测领域(通过关键词或分类模型):
is_code_related(text)→local_code_model - 检测复杂度:文本长度 > 2000 →
claude -
默认:
gpt4 -
实现路由链:
方案一:使用 RunnableBranch 顺序匹配
from langchain_core.runnables import RunnableBranch
def is_chinese(input_dict):
# 简单的语言检测逻辑
return any('\u4e00' <= ch <= '\u9fff' for ch in input_dict["question"])
def is_code(input_dict):
return "def " in input_dict["question"] or "class " in input_dict["question"]
branch = RunnableBranch(
(lambda x: is_chinese(x), chinese_chain),
(lambda x: is_code(x), code_chain),
default_chain # gpt4
)
result = branch.invoke({"question": "..."})
方案二:使用一个“路由器” LLM 进行意图分类
先用一个轻量、便宜的模型分析用户输入,输出一个标签(如“翻译”、“代码”、“通用”),然后根据标签选择模型。
router_prompt = PromptTemplate.from_template("""
判断以下用户请求的类型,只输出一个词:翻译、代码、通用。
请求:{question}
类型:""")
router_chain = router_prompt | cheap_model | StrOutputParser()
def route_by_intent(input_dict):
intent = router_chain.invoke(input_dict)
if intent == "翻译":
return translation_chain.invoke(input_dict)
elif intent == "代码":
return code_chain.invoke(input_dict)
else:
return default_chain.invoke(input_dict)
full_chain = RunnableLambda(route_by_intent)
方案三:结合 LangGraph 实现带状态的复杂路由
对于需要多步判断或需要上下文的路由,可以用 LangGraph 构建一个状态图,在节点中执行分类和模型调用。
📊 生产化考量:
-
路由逻辑要尽量轻量,避免路由本身消耗大量 Token。
-
记录每次路由决策和模型调用,用于成本分析和 A/B 测试。
-
可以为路由配置回退:如果目标模型失败,自动降级到默认模型。
-
动态路由表可以存储在配置中心(如 Redis),实现热更新。