LangChain框架深度实战指南:从RAG管道构建到多Agent协作系统与企业级LLM应用架构全解析
LangChain框架深度实战指南:从RAG管道构建到多Agent协作系统与企业级LLM应用架构全解析
引言
LangChain作为大语言模型应用开发领域最知名的框架,在2026年依然是构建LLM应用的首选工具之一。随着AI Agent范式的兴起,LangChain从早期的Prompt chaining工具演进为一个包含LangGraph(Agent编排)、LangSmith(可观测性)、LangServe(部署服务)的完整生态系统。本文将从LangChain的核心抽象开始,系统性地讲解RAG管道构建、工具调用Agent、多Agent协作系统、记忆管理、流式输出、可观测性等关键主题,通过大量可运行的Python代码示例帮助读者掌握企业级LLM应用开发。
一、LangChain 2026生态全景
1.1 核心包结构
LangChain在2026年的包结构已经完全模块化。langchain-core包含基础抽象(Runnable、Message、PromptTemplate等),是所有其他包的依赖。langchain包含链和Agent的高级实现。langchain-community包含第三方集成(各种向量数据库、文档加载器等)。langchain-openai、langchain-anthropic等是特定模型提供商的集成包。langgraph是Agent编排框架,支持构建复杂的多Agent协作系统。langsmith是可观测性和评估平台。langserve用于将链部署为REST API。以下是典型的依赖配置:
# pyproject.toml
[project]
dependencies = [
"langchain-core>=0.3.0",
"langchain>=0.3.0",
"langchain-openai>=0.2.0",
"langchain-anthropic>=0.2.0",
"langchain-community>=0.3.0",
"langgraph>=0.2.0",
"langsmith>=0.1.0",
"chromadb>=0.5.0",
"pydantic>=2.0",
"tiktoken>=0.7.0",
]
1.2 Runnable架构
LangChain的核心抽象是Runnable接口。所有组件——Prompt、Model、Parser、Retriever、Tool——都实现了Runnable接口,可以通过统一的pipe操作符组合。这种设计使得组件之间可以自由组合,并支持批处理、流式输出和异步执行。以下是一个基本的Runnable组合示例:
from langchain_core.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI
from langchain_core.output_parsers import StrOutputParser
from langchain_core.runnables import RunnablePassthrough
# 创建组件
prompt = ChatPromptTemplate.from_messages([
("system", "你是一个专业的技术文档翻译器。将给定的英文技术文档翻译成中文,保持技术术语的准确性。"),
("user", "{input}"),
])
model = ChatOpenAI(model="gpt-4o", temperature=0.3)
parser = StrOutputParser()
# 组合成链 - 使用管道操作符
chain = prompt | model | parser
# 同步调用
result = chain.invoke({"input": "The Runnable interface allows for composing chains of components."})
print(result)
# 流式输出
for chunk in chain.stream({"input": "Streaming output is supported natively."}):
print(chunk, end="", flush=True)
# 批处理
results = chain.batch([
{"input": "Batch processing handles multiple inputs."},
{"input": "Async execution is also supported."},
])
1.3 LCEL表达式语言
LangChain Expression Language(LCEL)是LangChain的声明式链构建语法。通过LCEL,你可以用简洁的语法表达复杂的处理管道,包括并行执行、条件分支、错误处理等:
from langchain_core.runnables import RunnableParallel, RunnableLambda, RunnablePassthrough
import json
# 并行执行多个 retriever
parallel_retrieval = RunnableParallel(
code_docs=code_retriever,
api_docs=api_retriever,
user_context=RunnablePassthrough() | RunnableLambda(lambda x: x.get("user_id", "anonymous")),
)
# 带条件分支的链
def route_query(state: dict) -> str:
query = state["query"].lower()
if "code" in query or "function" in query:
return "code_search"
elif "api" in query or "endpoint" in query:
return "api_search"
else:
return "general"
chain = (
RunnableParallel(
query=RunnablePassthrough() | RunnableLambda(lambda x: x["query"]),
route=RunnableLambda(route_query),
)
| RunnableLambda(lambda x: {
"results": code_retriever.invoke(x["query"]) if x["route"] == "code_search"
else api_retriever.invoke(x["query"]) if x["route"] == "api_search"
else general_retriever.invoke(x["query"]),
"route": x["route"],
})
| prompt
| model
| parser
)
二、RAG管道深度实战
2.1 文档处理与向量化
RAG(Retrieval-Augmented Generation)是LLM应用最常见的模式。以下是一个完整的生产级RAG管道实现,包含文档加载、分块、向量化、存储和检索:
from langchain_community.document_loaders import (
DirectoryLoader,
TextLoader,
PyPDFLoader,
UnstructuredMarkdownLoader,
)
from langchain_text_splitters import RecursiveCharacterTextSplitter, Language
from langchain_openai import OpenAIEmbeddings
from langchain_community.vectorstores import Chroma
from langchain_core.documents import Document
import hashlib
from typing import Optional
class DocumentPipeline:
"""生产级文档处理管道,支持增量更新和去重。"""
def __init__(self, persist_dir: str, embeddings_model: str = "text-embedding-3-large"):
self.persist_dir = persist_dir
self.embeddings = OpenAIEmbeddings(model=embeddings_model)
self.text_splitter = RecursiveCharacterTextSplitter(
chunk_size=1000,
chunk_overlap=200,
separators=["\n\n## ", "\n\n# ", "\n\n", "\n", " ", ""],
length_function=len,
)
self.code_splitter = RecursiveCharacterTextSplitter.from_language(
language=Language.PYTHON,
chunk_size=1500,
chunk_overlap=300,
)
self.vectorstore: Optional[Chroma] = None
def load_documents(self, source_dir: str) -> list[Document]:
"""从目录加载多种格式的文档。"""
documents = []
# Markdown文档
md_loader = DirectoryLoader(
source_dir,
glob="**/*.md",
loader_cls=UnstructuredMarkdownLoader,
show_progress=True,
)
documents.extend(md_loader.load())
# Python代码文件
py_loader = DirectoryLoader(
source_dir,
glob="**/*.py",
loader_cls=TextLoader,
show_progress=True,
)
py_docs = py_loader.load()
for doc in py_docs:
doc.metadata["doc_type"] = "code"
documents.extend(py_docs)
# PDF文档
pdf_loader = DirectoryLoader(
source_dir,
glob="**/*.pdf",
loader_cls=PyPDFLoader,
show_progress=True,
)
documents.extend(pdf_loader.load())
# 纯文本
txt_loader = DirectoryLoader(
source_dir,
glob="**/*.txt",
loader_cls=TextLoader,
show_progress=True,
)
documents.extend(txt_loader.load())
return documents
def split_documents(self, documents: list[Document]) -> list[Document]:
"""智能分块:代码文件用代码分块器,其他用文本分块器。"""
code_docs = [d for d in documents if d.metadata.get("doc_type") == "code"]
text_docs = [d for d in documents if d.metadata.get("doc_type") != "code"]
split_code = self.code_splitter.split_documents(code_docs)
split_text = self.text_splitter.split_documents(text_docs)
all_splits = split_code + split_text
# 为每个分块生成唯一ID,支持增量更新
for chunk in all_splits:
content_hash = hashlib.md5(
chunk.page_content.encode()
).hexdigest()
chunk.metadata["content_hash"] = content_hash
chunk.metadata["chunk_id"] = f"{chunk.metadata.get('source', 'unknown')}_{content_hash[:12]}"
return all_splits
def build_vectorstore(self, documents: list[Document]):
"""构建或更新向量数据库。"""
splits = self.split_documents(documents)
# 检查已有向量,实现增量更新
try:
existing = Chroma(
persist_directory=self.persist_dir,
embedding_function=self.embeddings,
)
existing_ids = set(existing.get()["ids"])
new_splits = [s for s in splits if s.metadata["chunk_id"] not in existing_ids]
if new_splits:
existing.add_documents(
new_splits,
ids=[s.metadata["chunk_id"] for s in new_splits],
)
print(f"Added {len(new_splits)} new chunks (total: {len(splits)})")
self.vectorstore = existing
except Exception:
# 首次构建
self.vectorstore = Chroma.from_documents(
documents=splits,
embedding=self.embeddings,
persist_directory=self.persist_dir,
ids=[s.metadata["chunk_id"] for s in splits],
)
print(f"Built vectorstore with {len(splits)} chunks")
def get_retriever(self, k: int = 5, search_type: str = "mmr"):
"""获取检索器,支持多种搜索策略。"""
if not self.vectorstore:
raise RuntimeError("Vectorstore not built. Call build_vectorstore first.")
if search_type == "mmr":
return self.vectorstore.as_retriever(
search_type="mmr",
search_kwargs={"k": k, "fetch_k": k * 4, "lambda_mult": 0.7},
)
elif search_type == "similarity":
return self.vectorstore.as_retriever(
search_kwargs={"k": k},
)
else:
raise ValueError(f"Unknown search type: {search_type}")
# 使用示例
pipeline = DocumentPipeline(persist_dir="./vectorstore")
docs = pipeline.load_documents("./docs")
pipeline.build_vectorstore(docs)
retriever = pipeline.get_retriever(k=5, search_type="mmr")
2.2 检索增强生成管道
以下是一个包含查询重写、多路检索、重排序和答案生成的完整RAG管道:
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.runnables import RunnableParallel, RunnablePassthrough, RunnableLambda
from langchain_openai import ChatOpenAI
from langchain_cohere import CohereRerank
from langchain_core.output_parsers import StrOutputParser, JsonOutputParser
from pydantic import BaseModel, Field
from typing import Optional
class QueryRewriteResult(BaseModel):
"""查询重写结果。"""
original_query: str = Field(description="原始查询")
rewritten_queries: list[str] = Field(description="重写后的查询变体列表")
reasoning: str = Field(description="重写理由")
def create_rag_chain(
retriever,
model_name: str = "gpt-4o",
use_reranking: bool = True,
use_query_rewrite: bool = True,
):
"""创建完整的RAG管道。"""
# === 查询重写 ===
rewrite_prompt = ChatPromptTemplate.from_messages([
("system", """你是一个搜索查询优化专家。分析用户的原始查询,
生成2-3个语义不同但相关的查询变体,以提高检索覆盖率。
输出JSON格式:
{{
"original_query": "原始查询",
"rewritten_queries": ["变体1", "变体2", "变体3"],
"reasoning": "重写理由"
}}"""),
("user", "原始查询:{query}"),
])
rewrite_model = ChatOpenAI(model=model_name, temperature=0.3)
rewrite_chain = rewrite_prompt | rewrite_model | JsonOutputParser(pydantic_object=QueryRewriteResult)
# === 多路检索 ===
def multi_retrieve(state: dict) -> dict:
query = state["query"]
rewrite_result = state.get("rewrite_result")
all_queries = [query]
if rewrite_result and use_query_rewrite:
all_queries.extend(rewrite_result.get("rewritten_queries", []))
# 去重
all_queries = list(set(all_queries))
# 多路检索
all_docs = []
for q in all_queries:
docs = retriever.invoke(q)
all_docs.extend(docs)
# 去重(基于内容哈希)
seen = set()
unique_docs = []
for doc in all_docs:
h = hash(doc.page_content)
if h not in seen:
seen.add(h)
unique_docs.append(doc)
# 重排序
if use_reranking:
try:
reranker = CohereRerank(top_n=5)
unique_docs = reranker.compress_documents(unique_docs, query)
except Exception:
pass # Cohere不可用时跳过重排序
return {**state, "retrieved_docs": unique_docs[:5]}
# === 答案生成 ===
answer_prompt = ChatPromptTemplate.from_messages([
("system", """你是一个专业的技术问答助手。根据以下检索到的文档上下文回答用户问题。
规则:
1. 只基于提供的文档内容回答,不要编造信息
2. 如果文档中没有相关信息,明确说明"根据现有文档无法回答此问题"
3. 在回答末尾标注引用来源,格式:[来源: 文件名]
4. 如果多个文档信息冲突,指出冲突并说明各来源
检索到的文档:
{context}"""),
("user", "{query}"),
])
answer_model = ChatOpenAI(model=model_name, temperature=0.1)
# === 组装管道 ===
def format_docs(docs):
formatted = []
for i, doc in enumerate(docs):
source = doc.metadata.get("source", "unknown")
formatted.append(f"[文档{i+1}] 来源: {source}\n内容: {doc.page_content}\n")
return "\n".join(formatted)
if use_query_rewrite:
chain = (
RunnableParallel(
query=RunnablePassthrough() | RunnableLambda(lambda x: x["query"]),
rewrite_result=RunnableLambda(lambda x: {"query": x["query"]}) | rewrite_chain,
)
| RunnableLambda(multi_retrieve)
| RunnableParallel(
query=RunnableLambda(lambda x: x["query"]),
context=RunnableLambda(lambda x: format_docs(x["retrieved_docs"])),
)
| answer_prompt
| answer_model
| StrOutputParser()
)
else:
chain = (
RunnableParallel(
query=RunnablePassthrough() | RunnableLambda(lambda x: x["query"]),
rewrite_result=RunnableLambda(lambda x: {"rewrite_result": None}),
)
| RunnableLambda(multi_retrieve)
| RunnableParallel(
query=RunnableLambda(lambda x: x["query"]),
context=RunnableLambda(lambda x: format_docs(x["retrieved_docs"])),
)
| answer_prompt
| answer_model
| StrOutputParser()
)
return chain
# 使用示例
rag_chain = create_rag_chain(retriever, use_reranking=True, use_query_rewrite=True)
result = rag_chain.invoke({
"query": "如何在LangChain中实现自定义的Retriever?"
})
print(result)
2.3 带引用追踪的RAG
在企业场景中,追踪答案的来源引用至关重要。以下是一个带精确引用追踪的RAG实现:
from pydantic import BaseModel, Field
from langchain_core.output_parsers import PydanticOutputParser
class Citation(BaseModel):
"""答案引用来源。"""
source: str = Field(description="文档来源文件名")
chunk_id: str = Field(description="文档分块ID")
relevant_text: str = Field(description="引用的具体文本片段")
class AnswerWithCitations(BaseModel):
"""带引用的答案。"""
answer: str = Field(description="对用户问题的回答")
confidence: float = Field(description="置信度 0-1")
citations: list[Citation] = Field(description="引用来源列表")
follow_up_questions: list[str] = Field(description="建议的后续问题")
def create_cited_rag_chain(retriever, model_name="gpt-4o"):
parser = PydanticOutputParser(pydantic_object=AnswerWithCitations)
prompt = ChatPromptTemplate.from_messages([
("system", """你是一个严谨的技术问答助手。回答用户问题并标注引用来源。
{format_instructions}
可用的文档上下文:
{context}"""),
("user", "{query}"),
]).partial(format_instructions=parser.get_format_instructions())
model = ChatOpenAI(model=model_name, temperature=0)
chain = (
RunnableParallel(
query=RunnableLambda(lambda x: x["query"]),
context=RunnableLambda(lambda x: x["query"]) | retriever | RunnableLambda(format_docs_with_ids),
)
| prompt
| model
| parser
)
return chain
def format_docs_with_ids(docs):
formatted = []
for i, doc in enumerate(docs):
source = doc.metadata.get("source", "unknown")
chunk_id = doc.metadata.get("chunk_id", f"chunk_{i}")
formatted.append(f"[ID:{chunk_id}] 来源:{source}\n{doc.page_content}\n")
return "\n---\n".join(formatted)
三、工具调用Agent
3.1 ReAct Agent
LangChain的Agent可以通过工具调用来完成复杂任务。以下是一个包含多个工具的技术支持Agent:
from langchain_core.tools import tool
from langchain_openai import ChatOpenAI
from langgraph.prebuilt import create_react_agent
from langchain_core.messages import HumanMessage
import subprocess
import requests
import json
# === 定义工具 ===
@tool
def search_codebase(query: str, file_pattern: str = "*.py") -> str:
"""在代码库中搜索匹配的代码片段。
Args:
query: 搜索关键词或正则表达式
file_pattern: 文件类型过滤,如 *.py, *.ts, *.go
"""
try:
cmd = f"rg '{query}' --type-add 'custom:{file_pattern}' -t custom -n --max-count 5"
result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=10)
if result.returncode == 0:
lines = result.stdout.strip().split('\n')
return f"找到 {len(lines)} 处匹配:\n" + '\n'.join(lines[:20])
return "未找到匹配结果"
except Exception as e:
return f"搜索失败: {e}"
@tool
def run_python_code(code: str) -> str:
"""执行Python代码并返回输出。适用于数据处理、计算、文件操作等。
Args:
code: 要执行的Python代码
"""
try:
result = subprocess.run(
["python3", "-c", code],
capture_output=True, text=True, timeout=30
)
output = result.stdout
if result.stderr:
output += f"\nSTDERR: {result.stderr}"
if len(output) > 5000:
output = output[:5000] + "\n...(truncated)"
return output if output else "(no output)"
except subprocess.TimeoutExpired:
return "代码执行超时(30秒限制)"
except Exception as e:
return f"执行失败: {e}"
@tool
def check_api_endpoint(url: str, method: str = "GET") -> str:
"""检查API端点的状态和响应。
Args:
url: API端点URL
method: HTTP方法 (GET, POST, PUT, DELETE)
"""
try:
response = requests.request(
method, url, timeout=10,
headers={"Accept": "application/json"}
)
return json.dumps({
"status_code": response.status_code,
"headers": dict(response.headers),
"body": response.text[:2000],
}, indent=2, ensure_ascii=False)
except Exception as e:
return f"请求失败: {e}"
@tool
def analyze_error(error_message: str) -> str:
"""分析错误消息并提供修复建议。
Args:
error_message: 完整的错误堆栈信息
"""
analysis_prompt = f"""分析以下错误并提供修复建议:
错误信息:
{error_message}
请提供:
1. 错误类型和根本原因
2. 常见触发场景
3. 具体修复步骤
4. 预防措施
"""
model = ChatOpenAI(model="gpt-4o", temperature=0.2)
response = model.invoke(analysis_prompt)
return response.content
# === 创建Agent ===
tools = [search_codebase, run_python_code, check_api_endpoint, analyze_error]
model = ChatOpenAI(model="gpt-4o")
agent = create_react_agent(model, tools)
# 使用示例
result = agent.invoke({
"messages": [HumanMessage(content="""
请帮我完成以下任务:
1. 搜索代码库中所有的TODO注释
2. 检查 http://localhost:3000/api/health 端点是否正常
3. 如果端点返回错误,分析错误原因
""")]
})
for msg in result["messages"]:
print(f"[{msg.type}]: {msg.content[:200]}")
3.2 自定义工具
以下是一个更复杂的自定义工具示例,提供数据库查询能力并包含安全限制:
from langchain_core.tools import BaseTool
from pydantic import BaseModel, Field
import asyncpg
from typing import Optional, Type
import re
class DatabaseQueryInput(BaseModel):
sql: str = Field(description="SQL查询语句(仅支持SELECT)")
database: str = Field(default="main", description="数据库名称")
class DatabaseQueryTool(BaseTool):
"""安全的数据库查询工具,仅支持只读SQL。"""
name: str = "database_query"
description: str = """执行只读SQL查询。仅支持SELECT语句。
可以查询用户信息、订单数据、产品库存等。
不支持INSERT、UPDATE、DELETE等修改操作。"""
args_schema: Type[BaseModel] = DatabaseQueryInput
connection_string: str = ""
FORBIDDEN_KEYWORDS = [
"INSERT", "UPDATE", "DELETE", "DROP", "ALTER", "TRUNCATE",
"CREATE", "GRANT", "REVOKE", "COPY", "EXEC", "MERGE"
]
def _validate_sql(self, sql: str) -> bool:
upper = sql.strip().upper()
if not upper.startswith("SELECT") and not upper.startswith("WITH"):
return False
for kw in self.FORBIDDEN_KEYWORDS:
if re.search(r'\b' + kw + r'\b', upper):
return False
return True
def _run(self, sql: str, database: str = "main") -> str:
if not self._validate_sql(sql):
return "错误:仅支持SELECT或WITH开头的只读查询"
async def execute():
conn = await asyncpg.connect(self.connection_string)
try:
rows = await conn.fetch(sql)
if not rows:
return "查询结果为空"
columns = [col for col in rows[0].keys()]
result = [dict(row) for row in rows]
return json.dumps({
"columns": columns,
"row_count": len(result),
"rows": result[:100],
"truncated": len(result) > 100,
}, ensure_ascii=False, indent=2, default=str)
finally:
await conn.close()
import asyncio
return asyncio.run(execute())
async def _arun(self, sql: str, database: str = "main") -> str:
if not self._validate_sql(sql):
return "错误:仅支持SELECT或WITH开头的只读查询"
conn = await asyncpg.connect(self.connection_string)
try:
rows = await conn.fetch(sql)
if not rows:
return "查询结果为空"
result = [dict(row) for row in rows]
return json.dumps({
"row_count": len(result),
"rows": result[:100],
}, ensure_ascii=False, indent=2, default=str)
finally:
await conn.close()
四、LangGraph多Agent协作
4.1 多Agent架构
LangGraph是LangChain的Agent编排框架,支持构建复杂的多Agent协作系统。以下是一个软件开发团队的多Agent系统,包含产品经理Agent、开发者Agent、测试Agent和代码审查Agent:
from langgraph.graph import StateGraph, END
from langchain_core.messages import HumanMessage, AIMessage, SystemMessage
from langchain_openai import ChatOpenAI
from typing import TypedDict, Annotated, Literal
import operator
class TeamState(TypedDict):
"""团队协作状态。"""
task: str
requirements: str
code: str
test_results: str
review_feedback: str
iteration: int
status: Literal["planning", "developing", "testing", "reviewing", "done", "failed"]
messages: Annotated[list, operator.add]
def create_dev_team():
"""创建软件开发团队多Agent系统。"""
model = ChatOpenAI(model="gpt-4o", temperature=0.3)
# === 产品经理Agent ===
def product_manager(state: TeamState) -> dict:
pm_prompt = f"""你是产品经理。分析以下任务,生成详细的需求文档。
任务:{state['task']}
请输出:
1. 功能需求列表
2. 非功能需求(性能、安全等)
3. 验收标准
4. 技术建议
注意:需求要具体、可执行。"""
response = model.invoke([SystemMessage(content=pm_prompt)])
return {
"requirements": response.content,
"status": "developing",
"messages": [AIMessage(content=f"[PM] {response.content}", name="product_manager")],
}
# === 开发者Agent ===
def developer(state: TeamState) -> dict:
dev_prompt = f"""你是资深开发者。根据以下需求实现功能代码。
需求文档:
{state['requirements']}
{f"上一轮审查反馈(请修复):{state['review_feedback']}" if state.get('review_feedback') else ""}
请提供完整可运行的代码实现,包含:
1. 核心功能代码
2. 输入验证
3. 错误处理
4. 类型注解
5. 简要注释"""
response = model.invoke([SystemMessage(content=dev_prompt)])
return {
"code": response.content,
"status": "testing",
"messages": [AIMessage(content=f"[Dev] 代码已生成", name="developer")],
}
# === 测试Agent ===
def tester(state: TeamState) -> dict:
test_prompt = f"""你是测试工程师。为以下代码编写并执行测试。
代码:
{state['code']}
请提供:
1. 单元测试代码
2. 边界条件测试
3. 错误场景测试
4. 测试覆盖率分析
5. 发现的问题列表"""
response = model.invoke([SystemMessage(content=test_prompt)])
return {
"test_results": response.content,
"status": "reviewing",
"messages": [AIMessage(content=f"[QA] {response.content[:500]}", name="tester")],
}
# === 代码审查Agent ===
def reviewer(state: TeamState) -> dict:
review_prompt = f"""你是资深代码审查员。审查以下代码和测试结果。
代码:
{state['code']}
测试结果:
{state['test_results']}
审查标准:
1. 代码质量(可读性、可维护性)
2. 安全性
3. 性能
4. 测试充分性
5. 错误处理完整性
如果通过审查,回复"APPROVED"。
如果需要修改,回复"CHANGES_REQUESTED"并列出具体问题。"""
response = model.invoke([SystemMessage(content=review_prompt)])
content = response.content
if "APPROVED" in content.upper():
return {
"status": "done",
"review_feedback": "",
"messages": [AIMessage(content=f"[Reviewer] APPROVED", name="reviewer")],
}
else:
return {
"status": "developing",
"review_feedback": content,
"iteration": state.get("iteration", 0) + 1,
"messages": [AIMessage(content=f"[Reviewer] CHANGES_REQUESTED", name="reviewer")],
}
# === 路由函数 ===
def route(state: TeamState) -> str:
status = state["status"]
iteration = state.get("iteration", 0)
if status == "done":
return END
if iteration >= 3:
return END # 最多3轮迭代
if status == "developing":
return "developer"
if status == "testing":
return "tester"
if status == "reviewing":
return "reviewer"
return END
# === 构建图 ===
workflow = StateGraph(TeamState)
workflow.add_node("product_manager", product_manager)
workflow.add_node("developer", developer)
workflow.add_node("tester", tester)
workflow.add_node("reviewer", reviewer)
workflow.set_entry_point("product_manager")
workflow.add_edge("product_manager", "developer")
workflow.add_edge("developer", "tester")
workflow.add_edge("tester", "reviewer")
workflow.add_conditional_edges("reviewer", route, {
"developer": "developer",
END: END,
})
return workflow.compile()
# 使用示例
team = create_dev_team()
result = team.invoke({
"task": "实现一个用户注册API端点,支持邮箱验证、密码强度检查和重复注册防护",
"iteration": 0,
"status": "planning",
"messages": [],
})
print(f"最终状态: {result['status']}")
print(f"迭代次数: {result.get('iteration', 0)}")
print(f"\n代码:\n{result['code'][:1000]}")
4.2 人机协作Agent
在某些场景下,Agent需要在关键决策点请求人类确认。LangGraph支持通过interrupt机制实现人机协作:
from langgraph.graph import StateGraph, END
from langgraph.checkpoint.memory import MemorySaver
from langchain_core.tools import tool
from langchain_openai import ChatOpenAI
from langgraph.prebuilt import create_react_agent
from typing import TypedDict, Annotated, Optional
import operator
class HumanInLoopState(TypedDict):
task: str
plan: Optional[str]
human_approved: bool
result: Optional[str]
messages: Annotated[list, operator.add]
def create_human_in_loop_agent():
"""创建带人类确认环节的Agent。"""
model = ChatOpenAI(model="gpt-4o")
def plan_step(state: HumanInLoopState) -> dict:
prompt = f"""为以下任务制定详细执行计划:
任务:{state['task']}
请提供分步计划,包括:
1. 需要执行的操作
2. 涉及的文件
3. 预期结果
4. 风险评估"""
response = model.invoke(prompt)
return {"plan": response.content, "messages": [response]}
def execute_step(state: HumanInLoopState) -> dict:
if not state["human_approved"]:
return {"result": "计划未获批准,执行已取消"}
prompt = f"""执行以下已批准的计划:
任务:{state['task']}
计划:{state['plan']}
请逐步执行计划并提供结果。"""
response = model.invoke(prompt)
return {"result": response.content, "messages": [response]}
def should_execute(state: HumanInLoopState) -> str:
if state.get("human_approved"):
return "execute"
return "wait"
workflow = StateGraph(HumanInLoopState)
workflow.add_node("plan", plan_step)
workflow.add_node("execute", execute_step)
workflow.set_entry_point("plan")
workflow.add_conditional_edges("plan", should_execute, {
"execute": "execute",
"wait": "wait_for_approval",
})
workflow.add_node("wait_for_approval", lambda x: x)
workflow.add_edge("wait_for_approval", "execute")
workflow.add_edge("execute", END)
memory = MemorySaver()
return workflow.compile(checkpointer=memory, interrupt_before=["wait_for_approval"])
# 使用示例
agent = create_human_in_loop_agent()
config = {"configurable": {"thread_id": "task-001"}}
# 第一轮 - 生成计划
result = agent.invoke({"task": "重构认证模块,从session迁移到JWT"}, config)
print("计划已生成:")
print(result["plan"])
# 等待人类批准
human_input = input("批准此计划?(yes/no): ")
agent.update_state(config, {"human_approved": human_input.lower() == "yes"})
# 继续执行
result = agent.invoke(None, config)
print("执行结果:", result.get("result"))
五、记忆与对话状态管理
5.1 对话记忆
LangChain提供多种记忆机制来保持对话上下文。以下是一个支持长期记忆的对话系统:
from langchain_core.chat_history import BaseChatMessageHistory
from langchain_community.chat_message_histories import RedisChatMessageHistory
from langchain_core.runnables import RunnablePassthrough, RunnableLambda
from langchain_core.prompts import ChatPromptTemplate, MessagesPlaceholder
from langchain_openai import ChatOpenAI
from langchain_core.output_parsers import StrOutputParser
from typing import Optional
class ConversationManager:
"""带持久化记忆的对话管理器。"""
def __init__(self, redis_url: str = "redis://localhost:6379"):
self.redis_url = redis_url
self.model = ChatOpenAI(model="gpt-4o", temperature=0.7)
self.prompt = ChatPromptTemplate.from_messages([
("system", """你是一个智能助手,具有长期记忆能力。
以下是之前的对话历史(可能被截断):
{summary}
当前用户的个人信息和偏好:
{user_context}"""),
MessagesPlaceholder(variable_name="history"),
("user", "{input}"),
])
def get_history(self, session_id: str) -> BaseChatMessageHistory:
return RedisChatMessageHistory(
session_id=session_id,
url=self.redis_url,
key_prefix="chat:",
)
def summarize_history(self, messages: list) -> str:
"""当历史过长时生成摘要。"""
if len(messages) < 10:
return "(无历史摘要)"
recent = messages[-10:]
old = messages[:-10]
summary_prompt = f"""请总结以下对话的关键信息:
{chr(10).join([f'{m.type}: {m.content}' for m in old])}
请保留:
1. 用户的关键偏好
2. 重要的决定
3. 未解决的问题
"""
response = self.model.invoke(summary_prompt)
return response.content
def chat(self, session_id: str, user_input: str, user_context: str = "") -> str:
history = self.get_history(session_id)
messages = history.messages
summary = self.summarize_history(messages)
chain = (
RunnablePassthrough.assign(
history=RunnableLambda(lambda _: messages),
summary=RunnableLambda(lambda _: summary),
user_context=RunnableLambda(lambda _: user_context),
)
| self.prompt
| self.model
| StrOutputParser()
)
response = chain.invoke({"input": user_input})
# 保存到历史
history.add_user_message(user_input)
history.add_ai_message(response)
return response
# 使用示例
manager = ConversationManager()
print(manager.chat("user-001", "我叫张三,是一个Go后端开发者"))
print(manager.chat("user-001", "我最近在学Rust,有什么建议?"))
print(manager.chat("user-001", "根据我之前说的信息,给我推荐一个适合我的项目"))
六、LangSmith可观测性
6.1 追踪与调试
LangSmith提供LLM应用的完整可观测性。通过设置环境变量即可自动追踪所有链调用:
import os
os.environ["LANGCHAIN_TRACING_V2"] = "true"
os.environ["LANGCHAIN_API_KEY"] = "ls__your-api-key"
os.environ["LANGCHAIN_PROJECT"] = "production-rag-app"
# 所有后续的LangChain调用都会自动被追踪
# 在LangSmith UI中可以看到完整的调用链、延迟、token使用量等
6.2 评估框架
LangSmith提供评估框架来系统性地评估LLM应用质量:
from langsmith import Client
from langchain_core.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI
client = Client()
# 定义评估器
def correctness_evaluator(run, example):
"""评估答案正确性。"""
expected = example.outputs.get("expected")
actual = run.outputs.get("output")
eval_prompt = f"""判断以下回答是否正确:
问题:{example.inputs.get("question")}
期望答案:{expected}
实际答案:{actual}
输出JSON:{{"score": 0或1, "reason": "理由"}}"""
model = ChatOpenAI(model="gpt-4o", temperature=0)
response = model.invoke(eval_prompt)
import json
result = json.loads(response.content)
return {
"key": "correctness",
"score": result["score"],
"comment": result["reason"],
}
def relevance_evaluator(run, example):
"""评估答案相关性。"""
question = example.inputs.get("question")
answer = run.outputs.get("output")
eval_prompt = f"""评估回答与问题的相关性(0-5分):
问题:{question}
回答:{answer[:500]}
输出JSON:{{"score": 0到5, "reason": "理由"}}"""
model = ChatOpenAI(model="gpt-4o", temperature=0)
response = model.invoke(eval_prompt)
result = json.loads(response.content)
return {
"key": "relevance",
"score": result["score"] / 5.0,
"comment": result["reason"],
}
# 创建评估数据集
dataset_name = "rag-eval-dataset"
try:
dataset = client.create_dataset(dataset_name)
examples = [
{"inputs": {"question": "如何在Python中读取CSV文件?"}, "outputs": {"expected": "使用csv模块或pandas.read_csv()"}},
{"inputs": {"question": "什么是RAG?"}, "outputs": {"expected": "检索增强生成,结合检索和生成"}},
{"inputs": {"question": "Docker和Kubernetes的区别?"}, "outputs": {"expected": "Docker是容器运行时,K8s是容器编排平台"}},
]
for ex in examples:
client.create_example(inputs=ex["inputs"], outputs=ex["outputs"], dataset_id=dataset.id)
except Exception:
dataset = client.read_dataset(dataset_name=dataset_name)
# 运行评估
results = client.run_on_dataset(
dataset_name=dataset_name,
llm_or_chain_factory=rag_chain,
evaluation=[correctness_evaluator, relevance_evaluator],
verbose=True,
)
print(f"平均正确率: {results['aggregate_stats']['correctness']['mean']:.2%}")
print(f"平均相关性: {results['aggregate_stats']['relevance']['mean']:.2f}/5.0")
七、LangServe部署
7.1 将链部署为API
LangServe可以将LangChain链快速部署为REST API:
from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware
from langserve import add_routes
from langchain_core.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI
from langchain_core.output_parsers import StrOutputParser
app = FastAPI(title="LangChain API Server")
app.add_middleware(CORSMiddleware, allow_origins=["*"], allow_methods=["*"], allow_headers=["*"])
# 创建RAG链
rag_chain = create_rag_chain(retriever)
# 暴露为API端点
add_routes(app, rag_chain, path="/rag")
# 暴露原始LLM调用
llm = ChatOpenAI(model="gpt-4o")
add_routes(app, llm, path="/llm")
# 自定义端点
@app.post("/chat")
async def chat(message: str, session_id: str = "default"):
manager = ConversationManager()
response = manager.chat(session_id, message)
return {"response": response, "session_id": session_id}
if __name__ == "__main__":
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8000)
部署后可以通过以下方式调用:
# 流式调用
curl -N http://localhost:8000/rag/stream -H "Content-Type: application/json" \
-d '{"query": "How to implement custom retriever in LangChain?"}'
# 批量调用
curl http://localhost:8000/rag/batch -H "Content-Type: application/json" \
-d '{"inputs": [{"query": "question 1"}, {"query": "question 2"}]}'
总结
LangChain在2026年依然是构建LLM应用的首选框架,其核心优势在于丰富的组件生态、统一的Runnable抽象、LangGraph的Agent编排能力以及LangSmith的可观测性支持。本文系统性地覆盖了RAG管道构建(含查询重写、多路检索、重排序和引用追踪)、工具调用Agent、LangGraph多Agent协作系统、人机协作模式、记忆管理、可观测性评估和API部署等核心内容。关键要点包括:使用LCEL表达式语言可以简洁地构建复杂处理管道;RAG管道应包含查询重写和重排序以提升检索质量;LangGraph支持构建多Agent协作系统处理复杂工作流;LangSmith的评估框架是保证LLM应用质量的关键工具。随着AI Agent范式的成熟,LangChain + LangGraph的组合将成为构建企业级AI应用的核心技术栈。
- 点赞
- 收藏
- 关注作者
评论(0)