LangChain 1.3实战:从RAG知识库到LangGraph多智能体工作流
这次我们来看一个完整的 LangChain 1.3 系统课程,从 RAG 应用到 LangGraph 多智能体工作流,覆盖了当前最热门的 AI 应用开发技术栈。如果你正在寻找一套能真正跑通的企业级解决方案,这篇文章值得收藏。
LangChain 1.3 是目前最稳定的版本之一,特别适合构建生产环境的 RAG 系统和多智能体工作流。与早期版本相比,1.3 版本在模块化、稳定性和性能上都有显著提升。本教程将带你从零搭建完整的 RAG 知识库,然后进阶到 LangGraph 的多智能体协作系统。
最核心的特点是实战导向:每个环节都有可运行的代码示例,支持本地部署和云环境,兼容 CPU 和 GPU 推理,能够处理批量任务,并且提供完整的 API 接口设计思路。无论是个人学习还是企业项目,这套方案都能快速验证效果。
1. 核心能力速览
| 能力项 | 说明 |
|---|---|
| 技术栈 | LangChain 1.3 + LangGraph + 向量数据库 + LLM |
| 主要功能 | RAG 知识库构建、多智能体工作流、简历筛选、文档处理 |
| 硬件要求 | CPU 可运行,GPU 加速推理(推荐 8G+ 显存) |
| 部署方式 | 本地部署、Docker 容器、云服务器 |
| 接口支持 | RESTful API、Streamlit WebUI、命令行工具 |
| 批量处理 | 支持文档批量导入、多任务并行处理 |
| 适合场景 | 企业知识库、智能客服、自动化流程、AI 助手 |
2. 适用场景与使用边界
这个教程特别适合以下人群:
- 想要系统学习 LangChain 和 LangGraph 的开发者
- 需要构建企业级 RAG 系统的技术团队
- 希望实现多智能体协作应用的 AI 工程师
- 从事自动化流程开发的软件工程师
能解决的具体问题包括:
- 企业文档知识库的智能问答
- 简历自动筛选和匹配
- 多步骤复杂任务的自动化处理
- AI 智能体的协同工作流
需要注意的使用边界:
- 涉及个人隐私的数据需要脱敏处理
- 商业使用需确保数据授权合规
- 大规模部署需要考虑性能优化
- 关键业务场景需要人工审核环节
3. 环境准备与前置条件
在开始实战之前,需要准备好以下环境:
3.1 基础环境要求
- 操作系统: Windows 10/11, macOS 10.15+, Ubuntu 18.04+
- Python 版本: 3.8-3.11(推荐 3.9)
- 内存: 至少 8GB,推荐 16GB+
- 存储: 至少 10GB 可用空间
3.2 开发工具准备
# 创建虚拟环境 python -m venv langchain_env source langchain_env/bin/activate # Linux/macOS # 或 langchain_env\Scripts\activate # Windows # 安装核心依赖 pip install langchain==1.3.11 pip install langchain-community==0.3.8 pip install langgraph==0.1.03.3 向量数据库选择
根据项目需求选择合适的向量数据库:
- Chroma: 轻量级,适合学习和中小项目
- Weaviate: 功能丰富,适合生产环境
- Pinecone: 云服务,免运维
- FAISS: 本地部署,性能优秀
4. LangChain 1.3 核心概念解析
4.1 组件架构升级
LangChain 1.3 最大的变化是模块化程度更高,各个组件职责更清晰:
from langchain_core.prompts import ChatPromptTemplate from langchain_core.output_parsers import StrOutputParser from langchain_community.llms import Ollama from langchain_community.embeddings import HuggingFaceEmbeddings # 1.3 版本的标准组件使用方式 prompt = ChatPromptTemplate.from_template("回答以下问题: {question}") llm = Ollama(model="llama3.1") output_parser = StrOutputParser() chain = prompt | llm | output_parser4.2 RAG 系统核心组件
一个完整的 RAG 系统包含以下关键组件:
from langchain.text_splitter import RecursiveCharacterTextSplitter from langchain_community.vectorstores import Chroma from langchain.chains import RetrievalQA # 文档处理流水线 text_splitter = RecursiveCharacterTextSplitter( chunk_size=1000, chunk_overlap=200 ) # 向量数据库配置 embeddings = HuggingFaceEmbeddings(model_name="all-MiniLM-L6-v2") vectorstore = Chroma.from_documents(documents, embeddings) # RAG 链构建 qa_chain = RetrievalQA.from_chain_type( llm=llm, chain_type="stuff", retriever=vectorstore.as_retriever() )5. 构建企业级 RAG 知识库
5.1 文档预处理与向量化
实际项目中,文档预处理是关键的第一步:
import os from langchain_community.document_loaders import PyPDFLoader, Docx2txtLoader def load_documents(directory_path): """加载目录下的所有文档""" documents = [] for filename in os.listdir(directory_path): file_path = os.path.join(directory_path, filename) if filename.endswith('.pdf'): loader = PyPDFLoader(file_path) elif filename.endswith('.docx'): loader = Docx2txtLoader(file_path) else: continue documents.extend(loader.load()) return documents # 批量处理文档 documents = load_documents("./企业文档/") split_docs = text_splitter.split_documents(documents) # 创建向量库 vectorstore = Chroma.from_documents( documents=split_docs, embedding=embeddings, persist_directory="./vector_db/" )5.2 智能检索优化
提升检索质量的关键技巧:
from langchain.retrievers import ContextualCompressionRetriever from langchain.retrievers.document_compressors import EmbeddingsFilter # 使用重排序提升检索精度 compressor = EmbeddingsFilter(embeddings=embeddings, similarity_threshold=0.7) compression_retriever = ContextualCompressionRetriever( base_compressor=compressor, base_retriever=vectorstore.as_retriever(search_kwargs={"k": 10}) ) # 带重排序的 RAG 链 advanced_qa_chain = RetrievalQA.from_chain_type( llm=llm, chain_type="stuff", retriever=compression_retriever )6. LangGraph 多智能体工作流实战
6.1 LangGraph 核心概念
LangGraph 通过图结构定义智能体工作流:
from langgraph.graph import Graph from langgraph.prebuilt import create_react_agent # 定义智能体节点 def research_agent(state): """研究智能体:负责信息搜集""" # 实现研究逻辑 return {"research_result": "搜集到的信息"} def analysis_agent(state): """分析智能体:负责数据分析""" # 实现分析逻辑 return {"analysis_result": "分析结果"} def decision_agent(state): """决策智能体:负责最终决策""" # 实现决策逻辑 return {"final_decision": "最终决策"} # 构建工作流图 workflow = Graph() workflow.add_node("research", research_agent) workflow.add_node("analysis", analysis_agent) workflow.add_node("decision", decision_agent) # 定义边连接 workflow.add_edge("research", "analysis") workflow.add_edge("analysis", "decision")6.2 简历筛选工作流案例
实现一个完整的简历筛选多智能体系统:
from typing import Dict, Any from langchain_core.messages import HumanMessage from langgraph.graph import END, START class ResumeScreeningWorkflow: def __init__(self): self.workflow = Graph() self._build_workflow() def _parse_resume(self, state: Dict[str, Any]): """解析简历智能体""" resume_text = state["resume_text"] # 实现简历解析逻辑 return {"parsed_info": "解析后的简历信息"} def _evaluate_skills(self, state: Dict[str, Any]): """技能评估智能体""" parsed_info = state["parsed_info"] job_requirements = state["job_requirements"] # 实现技能匹配逻辑 return {"skill_match_score": 0.85} def _make_decision(self, state: Dict[str, Any]): """决策智能体""" score = state["skill_match_score"] if score > 0.8: return {"decision": "推荐面试", "confidence": score} else: return {"decision": "暂不推荐", "confidence": score} def _build_workflow(self): """构建工作流图""" self.workflow.add_node("parse", self._parse_resume) self.workflow.add_node("evaluate", self._evaluate_skills) self.workflow.add_node("decide", self._make_decision) self.workflow.add_edge(START, "parse") self.workflow.add_edge("parse", "evaluate") self.workflow.add_edge("evaluate", "decide") self.workflow.add_edge("decide", END) def run(self, resume_text: str, job_requirements: str): """运行工作流""" initial_state = { "resume_text": resume_text, "job_requirements": job_requirements } return self.workflow.invoke(initial_state)7. 系统集成与 API 部署
7.1 RESTful API 设计
使用 FastAPI 提供企业级 API 服务:
from fastapi import FastAPI, HTTPException from pydantic import BaseModel import uvicorn app = FastAPI(title="LangChain RAG API") class QueryRequest(BaseModel): question: str context: str = "" class ResumeScreeningRequest(BaseModel): resume_text: str job_requirements: str @app.post("/rag/query") async def rag_query(request: QueryRequest): """RAG 问答接口""" try: result = qa_chain.invoke({"query": request.question}) return {"answer": result["result"], "status": "success"} except Exception as e: raise HTTPException(status_code=500, detail=str(e)) @app.post("/workflow/screen-resume") async def screen_resume(request: ResumeScreeningRequest): """简历筛选工作流接口""" try: workflow = ResumeScreeningWorkflow() result = workflow.run(request.resume_text, request.job_requirements) return {"result": result, "status": "success"} except Exception as e: raise HTTPException(status_code=500, detail=str(e)) if __name__ == "__main__": uvicorn.run(app, host="0.0.0.0", port=8000)7.2 批量任务处理
实现文档批量处理功能:
import asyncio from concurrent.futures import ThreadPoolExecutor import pandas as pd class BatchProcessor: def __init__(self, max_workers=4): self.executor = ThreadPoolExecutor(max_workers=max_workers) def process_document_batch(self, document_paths: list): """批量处理文档""" results = [] def process_single_document(path): # 单个文档处理逻辑 documents = load_documents([path]) # 向量化处理 # 返回处理结果 return {"path": path, "status": "processed"} # 并行处理 futures = [ self.executor.submit(process_single_document, path) for path in document_paths ] for future in futures: try: results.append(future.result(timeout=300)) # 5分钟超时 except Exception as e: results.append({"path": path, "status": "error", "error": str(e)}) return results # 使用示例 processor = BatchProcessor() results = processor.process_document_batch(["./doc1.pdf", "./doc2.docx"])8. 性能优化与资源管理
8.1 显存和内存优化
大型语言模型部署时的资源优化策略:
import gc import torch class ResourceManager: def __init__(self): self.memory_threshold = 0.8 # 80% 内存使用阈值 def check_memory_usage(self): """检查内存使用情况""" if torch.cuda.is_available(): allocated = torch.cuda.memory_allocated() / 1024**3 # GB cached = torch.cuda.memory_reserved() / 1024**3 # GB return allocated, cached return 0, 0 def optimize_memory(self): """内存优化""" if torch.cuda.is_available(): torch.cuda.empty_cache() gc.collect() def batch_processing_with_memory_control(self, data_list, batch_size=4): """带内存控制的批量处理""" results = [] for i in range(0, len(data_list), batch_size): batch = data_list[i:i+batch_size] # 处理当前批次 batch_results = self.process_batch(batch) results.extend(batch_results) # 检查内存并优化 allocated, cached = self.check_memory_usage() if allocated > 6: # 超过 6GB self.optimize_memory() return results8.2 缓存策略实现
减少重复计算,提升响应速度:
from datetime import datetime, timedelta import hashlib import json class QueryCache: def __init__(self, ttl_hours=24): self.cache = {} self.ttl = timedelta(hours=ttl_hours) def _generate_key(self, query: str, context: str = "") -> str: """生成缓存键""" content = query + context return hashlib.md5(content.encode()).hexdigest() def get(self, query: str, context: str = ""): """获取缓存结果""" key = self._generate_key(query, context) if key in self.cache: cached_data = self.cache[key] if datetime.now() - cached_data['timestamp'] < self.ttl: return cached_data['result'] else: del self.cache[key] # 过期删除 return None def set(self, query: str, result: any, context: str = ""): """设置缓存""" key = self._generate_key(query, context) self.cache[key] = { 'result': result, 'timestamp': datetime.now() } # 在 RAG 系统中使用缓存 cache = QueryCache() def cached_rag_query(query: str, context: str = ""): """带缓存的 RAG 查询""" cached_result = cache.get(query, context) if cached_result is not None: return cached_result # 执行实际查询 result = qa_chain.invoke({"query": query}) cache.set(query, result, context) return result9. 监控与日志系统
9.1 系统监控实现
完整的监控系统帮助发现性能瓶颈:
import time import logging from prometheus_client import Counter, Histogram, start_http_server # 监控指标 QUERY_COUNTER = Counter('rag_queries_total', 'Total RAG queries') QUERY_DURATION = Histogram('rag_query_duration_seconds', 'RAG query duration') ERROR_COUNTER = Counter('rag_errors_total', 'Total RAG errors') class MonitoringSystem: def __init__(self, log_level=logging.INFO): self.logger = logging.getLogger(__name__) self.setup_logging(log_level) def setup_logging(self, level): """设置日志系统""" logging.basicConfig( level=level, format='%(asctime)s - %(name)s - %(levelname)s - %(message)s', handlers=[ logging.FileHandler('langchain_system.log'), logging.StreamHandler() ] ) @QUERY_DURATION.time() def monitor_query(self, func): """监控查询性能的装饰器""" def wrapper(*args, **kwargs): QUERY_COUNTER.inc() start_time = time.time() try: result = func(*args, **kwargs) self.logger.info(f"Query completed successfully") return result except Exception as e: ERROR_COUNTER.inc() self.logger.error(f"Query failed: {str(e)}") raise return wrapper # 启动监控服务器 start_http_server(8000) # Prometheus metrics endpoint10. 安全与合规考虑
10.1 数据安全处理
企业级部署必须考虑的安全性:
import re from typing import List class SecurityFilter: def __init__(self): self.sensitive_patterns = [ r'\b\d{4}[- ]?\d{4}[- ]?\d{4}[- ]?\d{4}\b', # 银行卡号 r'\b\d{17}[\dXx]\b', # 身份证号 r'\b\d{11}\b', # 手机号 ] def filter_sensitive_info(self, text: str) -> str: """过滤敏感信息""" filtered_text = text for pattern in self.sensitive_patterns: filtered_text = re.sub(pattern, '[REDACTED]', filtered_text) return filtered_text def validate_input(self, text: str, max_length: int = 10000) -> bool: """输入验证""" if len(text) > max_length: return False # 检查潜在的安全风险 dangerous_patterns = [ r'<script.*?>.*?</script>', r'on\w+\s*=', r'javascript:' ] for pattern in dangerous_patterns: if re.search(pattern, text, re.IGNORECASE): return False return True # 在 API 中使用安全过滤 security_filter = SecurityFilter() @app.post("/secure/rag-query") async def secure_rag_query(request: QueryRequest): """安全的 RAG 查询接口""" if not security_filter.validate_input(request.question): raise HTTPException(status_code=400, detail="Invalid input") filtered_question = security_filter.filter_sensitive_info(request.question) result = cached_rag_query(filtered_question) return {"answer": result, "status": "success"}11. 测试与验证方案
11.1 单元测试框架
确保系统稳定性的测试方案:
import unittest from unittest.mock import Mock, patch class TestRAGSystem(unittest.TestCase): def setUp(self): """测试初始化""" self.qa_chain = setup_test_chain() self.test_documents = load_test_documents() def test_basic_query(self): """基础查询测试""" result = self.qa_chain.invoke({"query": "什么是机器学习?"}) self.assertIsInstance(result, dict) self.assertIn("result", result) self.assertGreater(len(result["result"]), 10) def test_empty_query(self): """空查询测试""" with self.assertRaises(ValueError): self.qa_chain.invoke({"query": ""}) @patch('langchain_community.llms.Ollama.invoke') def test_llm_failure(self, mock_llm): """LLM 失败测试""" mock_llm.side_effect = Exception("LLM service unavailable") with self.assertRaises(Exception): self.qa_chain.invoke({"query": "测试问题"}) class TestWorkflowSystem(unittest.TestCase): def test_resume_screening_workflow(self): """简历筛选工作流测试""" workflow = ResumeScreeningWorkflow() result = workflow.run( "软件工程师简历内容...", "需要 Python 和机器学习经验" ) self.assertIn("decision", result) self.assertIn("confidence", result) if __name__ == "__main__": unittest.main()11.2 集成测试方案
端到端的系统集成测试:
import requests import json class IntegrationTests: def __init__(self, base_url="http://localhost:8000"): self.base_url = base_url def test_rag_api(self): """RAG API 集成测试""" response = requests.post( f"{self.base_url}/rag/query", json={"question": "测试问题", "context": ""} ) assert response.status_code == 200 data = response.json() assert data["status"] == "success" assert "answer" in data def test_workflow_api(self): """工作流 API 集成测试""" response = requests.post( f"{self.base_url}/workflow/screen-resume", json={ "resume_text": "测试简历内容", "job_requirements": "测试职位要求" } ) assert response.status_code == 200 data = response.json() assert data["status"] == "success" # 运行集成测试 def run_integration_tests(): tester = IntegrationTests() tester.test_rag_api() tester.test_workflow_api() print("所有集成测试通过")12. 部署与运维指南
12.1 Docker 容器化部署
生产环境推荐使用 Docker 部署:
# Dockerfile FROM python:3.9-slim WORKDIR /app # 安装系统依赖 RUN apt-get update && apt-get install -y \ gcc \ g++ \ && rm -rf /var/lib/apt/lists/* # 复制依赖文件 COPY requirements.txt . # 安装 Python 依赖 RUN pip install --no-cache-dir -r requirements.txt # 复制应用代码 COPY . . # 暴露端口 EXPOSE 8000 # 启动命令 CMD ["python", "app.py"]对应的 docker-compose.yml:
version: '3.8' services: langchain-app: build: . ports: - "8000:8000" environment: - PYTHONPATH=/app - LLM_MODEL=llama3.1 volumes: - ./data:/app/data restart: unless-stopped # 可选:向量数据库服务 chroma-db: image: chromadb/chroma ports: - "8001:8000" volumes: - chroma_data:/data restart: unless-stopped volumes: chroma_data:12.2 性能调优配置
生产环境性能优化配置:
# config.py import os class Config: # 性能配置 MAX_CONCURRENT_QUERIES = int(os.getenv('MAX_CONCURRENT_QUERIES', 10)) QUERY_TIMEOUT = int(os.getenv('QUERY_TIMEOUT', 30)) BATCH_SIZE = int(os.getenv('BATCH_SIZE', 4)) # 缓存配置 CACHE_TTL_HOURS = int(os.getenv('CACHE_TTL_HOURS', 24)) CACHE_MAX_SIZE = int(os.getenv('CACHE_MAX_SIZE', 1000)) # 模型配置 EMBEDDING_MODEL = os.getenv('EMBEDDING_MODEL', 'all-MiniLM-L6-v2') LLM_MODEL = os.getenv('LLM_MODEL', 'llama3.1') # 安全配置 MAX_INPUT_LENGTH = int(os.getenv('MAX_INPUT_LENGTH', 10000)) ENABLE_SENSITIVE_FILTER = os.getenv('ENABLE_SENSITIVE_FILTER', 'true').lower() == 'true' # 环境变量示例 """ MAX_CONCURRENT_QUERIES=20 QUERY_TIMEOUT=60 BATCH_SIZE=8 CACHE_TTL_HOURS=48 LLM_MODEL=llama3.1:8b """这套 LangChain 1.3 系统从基础概念到企业级部署提供了完整解决方案。建议先按照 RAG 知识库的步骤搭建基础系统,验证核心功能后再逐步引入 LangGraph 多智能体工作流。实际部署时重点关注资源监控和安全性配置,确保系统稳定运行。
