LangChain框架深度实战指南:从RAG管道构建到多Agent协作系统与企业级LLM应用架构全解析

举报
江南清风起 发表于 2026/08/21 22:27:33 2026/08/21
【摘要】 LangChain框架深度实战指南:从RAG管道构建到多Agent协作系统与企业级LLM应用架构全解析 引言LangChain作为大语言模型应用开发领域最知名的框架,在2026年依然是构建LLM应用的首选工具之一。随着AI Agent范式的兴起,LangChain从早期的Prompt chaining工具演进为一个包含LangGraph(Agent编排)、LangSmith(可观测性)、L...

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应用的核心技术栈。

【声明】本内容来自华为云开发者社区博主,不代表华为云及华为云开发者社区的观点和立场。转载时必须标注文章的来源(华为云社区)、文章链接、文章作者等基本信息,否则作者和本社区有权追究责任。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱: cloudbbs@huaweicloud.com
  • 点赞
  • 收藏
  • 关注作者

评论(0

0/1000
抱歉,系统识别当前为高风险访问,暂不支持该操作

全部回复

上滑加载中

设置昵称

在此一键设置昵称,即可参与社区互动!

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。