LangChain 1.3实战:从RAG知识库到LangGraph多智能体工作流

📅 发布时间:2026/7/30 12:43:45
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.04Python 版本: 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 langchain1.3.11 pip install langchain-community0.3.8 pip install langgraph0.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(modelllama3.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_size1000, chunk_overlap200 ) # 向量数据库配置 embeddings HuggingFaceEmbeddings(model_nameall-MiniLM-L6-v2) vectorstore Chroma.from_documents(documents, embeddings) # RAG 链构建 qa_chain RetrievalQA.from_chain_type( llmllm, chain_typestuff, retrievervectorstore.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( documentssplit_docs, embeddingembeddings, persist_directory./vector_db/ )5.2 智能检索优化提升检索质量的关键技巧from langchain.retrievers import ContextualCompressionRetriever from langchain.retrievers.document_compressors import EmbeddingsFilter # 使用重排序提升检索精度 compressor EmbeddingsFilter(embeddingsembeddings, similarity_threshold0.7) compression_retriever ContextualCompressionRetriever( base_compressorcompressor, base_retrievervectorstore.as_retriever(search_kwargs{k: 10}) ) # 带重排序的 RAG 链 advanced_qa_chain RetrievalQA.from_chain_type( llmllm, chain_typestuff, retrievercompression_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(titleLangChain 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_code500, detailstr(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_code500, detailstr(e)) if __name__ __main__: uvicorn.run(app, host0.0.0.0, port8000)7.2 批量任务处理实现文档批量处理功能import asyncio from concurrent.futures import ThreadPoolExecutor import pandas as pd class BatchProcessor: def __init__(self, max_workers4): self.executor ThreadPoolExecutor(max_workersmax_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(timeout300)) # 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_size4): 带内存控制的批量处理 results [] for i in range(0, len(data_list), batch_size): batch data_list[i:ibatch_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_hours24): self.cache {} self.ttl timedelta(hoursttl_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_levellogging.INFO): self.logger logging.getLogger(__name__) self.setup_logging(log_level) def setup_logging(self, level): 设置日志系统 logging.basicConfig( levellevel, 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(fQuery completed successfully) return result except Exception as e: ERROR_COUNTER.inc() self.logger.error(fQuery 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 [ rscript.*?.*?/script, ron\w\s*, rjavascript: ] 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_code400, detailInvalid 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_urlhttp://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.ymlversion: 3.8 services: langchain-app: build: . ports: - 8000:8000 environment: - PYTHONPATH/app - LLM_MODELllama3.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_QUERIES20 QUERY_TIMEOUT60 BATCH_SIZE8 CACHE_TTL_HOURS48 LLM_MODELllama3.1:8b 这套 LangChain 1.3 系统从基础概念到企业级部署提供了完整解决方案。建议先按照 RAG 知识库的步骤搭建基础系统验证核心功能后再逐步引入 LangGraph 多智能体工作流。实际部署时重点关注资源监控和安全性配置确保系统稳定运行。