RAG的使用

结合文档,AI整理,作为参考

RAG使用指南:从理论到生产的完整实践手册

Date: October 23, 2025
Version: 2.0
Author: AI技术团队


目录

  1. 引言

  2. RAG技术原理与架构

  3. 主流RAG框架对比

  4. 环境搭建与配置

  5. 数据处理与准备

  6. 向量数据库选型指南

  7. RAG系统构建实践

  8. 性能优化策略

  9. 部署与监控

  10. 最佳实践与案例

  11. 2025年技术趋势

  12. 常见问题与解决方案


引言

什么是RAG

RAG(Retrieval-Augmented Generation,检索增强生成) 是一种将大型语言模型(LLM)与外部知识检索相结合的技术框架。它通过在生成回答前从外部知识库中检索相关信息,为LLM提供准确、最新的上下文,从而显著提升回答的准确性和可靠性。

RAG的核心价值

优势 描述
减少幻觉 通过外部知识验证,大幅降低LLM生成虚假信息的风险
知识时效性 能够访问最新信息,突破LLM训练数据的时间限制
领域适配性 可针对特定行业或企业构建专业知识库
成本效益 相比模型微调,实现知识更新的成本更低
可解释性 提供回答的来源引用,增强结果可信度

适用场景

  • 企业知识库问答:员工培训、政策查询、技术文档检索

  • 智能客服系统:产品咨询、故障排除、流程指导

  • 金融投研分析:市场动态、公司财报、行业报告

  • 医疗诊断支持:病历分析、医学文献、诊疗指南

  • 法律合规咨询:法条检索、案例分析、合规审查


RAG技术原理与架构

基本工作流程

RAG系统主要包含两个核心阶段:索引阶段查询阶段

索引阶段(离线处理)

1
文档加载 → 文本分块 → 向量嵌入 → 向量存储
  1. 文档加载:从各种数据源(PDF、Word、网页、数据库等)加载原始文档

  2. 文本分块:将长文档分割为适合处理的小块(Chunks)

  3. 向量嵌入:使用嵌入模型将文本块转换为高维向量

  4. 向量存储:将向量和元数据存储到向量数据库中

查询阶段(在线处理)

1
用户查询 → 查询理解 → 向量检索 → 上下文构建 → 答案生成
  1. 用户查询:用户提出问题或请求

  2. 查询理解:对查询进行预处理和优化

  3. 向量检索:在向量数据库中检索相关文档片段

  4. 上下文构建:将检索结果组织为LLM的输入上下文

  5. 答案生成:LLM结合上下文生成最终回答

技术架构演进

1. 基础RAG(Naive RAG)

  • 简单的"检索-生成"两阶段架构

  • 适用于简单问答场景

  • 实现难度低,部署快速

2. 高级RAG(Advanced RAG)

  • 增加查询重写、结果重排序等优化环节

  • 支持多轮对话和上下文记忆

  • 提升复杂问题的处理能力

3. 模块化RAG(Modular RAG)

  • 采用可插拔的模块化设计

  • 支持多模态数据处理

  • 与Agent技术深度融合


主流RAG框架对比

LangChain

特点

  • 功能全面,生态丰富

  • 支持多种LLM和向量数据库

  • 提供声明式的链构建方式

优势

  • 社区活跃,文档完善

  • 支持复杂的工作流编排

  • 与主流工具集成良好

适用场景

  • 快速原型开发

  • 复杂业务逻辑的RAG应用

  • 需要多工具协作的场景

代码示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
from langchain_community.document_loaders import WebBaseLoader
from langchain_text_splitters import RecursiveCharacterTextSplitter
from langchain_chroma import Chroma
from langchain_openai import OpenAIEmbeddings, ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.runnables import RunnablePassthrough
from langchain_core.output_parsers import StrOutputParser

# 1. 加载文档
loader = WebBaseLoader("https://example.com/docs")
docs = loader.load()

# 2. 文本分块
text_splitter = RecursiveCharacterTextSplitter(
chunk_size=1000,
chunk_overlap=200
)
splits = text_splitter.split_documents(docs)

# 3. 构建向量存储
vectorstore = Chroma.from_documents(
documents=splits,
embedding=OpenAIEmbeddings()
)

# 4. 创建检索器
retriever = vectorstore.as_retriever()

# 5. 设置提示模板
prompt = ChatPromptTemplate.from_template("""
基于以下上下文回答用户问题:
{context}

用户问题:{question}
""")

# 6. 构建RAG链
llm = ChatOpenAI(model="gpt-3.5-turbo")
rag_chain = (
{"context": retriever | format_docs, "question": RunnablePassthrough()}
| prompt
| llm
| StrOutputParser()
)

# 7. 执行查询
response = rag_chain.invoke("你的问题是什么?")

LlamaIndex

特点

  • 专为RAG优化的框架

  • 简化的数据索引流程

  • 强大的查询引擎

优势

  • 开箱即用的解决方案

  • 智能的文档处理能力

  • 支持多种查询模式

适用场景

  • 企业级知识库

  • 文档密集型应用

  • 需要高级查询功能的场景

代码示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
from llama_index.core import VectorStoreIndex, SimpleDirectoryReader
from llama_index.embeddings.openai import OpenAIEmbedding
from llama_index.llms.openai import OpenAI

# 1. 加载文档
documents = SimpleDirectoryReader("./data").load_data()

# 2. 配置模型
llm = OpenAI(model="gpt-3.5-turbo")
embed_model = OpenAIEmbedding(model="text-embedding-3-small")

# 3. 构建索引
index = VectorStoreIndex.from_documents(
documents,
llm=llm,
embed_model=embed_model
)

# 4. 创建查询引擎
query_engine = index.as_query_engine()

# 5. 执行查询
response = query_engine.query("你的问题是什么?")
print(response)

Haystack

特点

  • 模块化设计,高度可定制

  • 支持生产级部署

  • 强大的管道编排能力

优势

  • 企业级稳定性

  • 丰富的评估工具

  • 支持多模态处理

适用场景

  • 生产环境部署

  • 需要高度定制的场景

  • 企业级应用

框架选择建议

评估维度 LangChain LlamaIndex Haystack
开发速度 ⭐⭐⭐⭐⭐ ⭐⭐⭐⭐⭐ ⭐⭐⭐⭐
功能丰富度 ⭐⭐⭐⭐⭐ ⭐⭐⭐⭐ ⭐⭐⭐⭐⭐
生产就绪性 ⭐⭐⭐⭐ ⭐⭐⭐⭐ ⭐⭐⭐⭐⭐
学习曲线 ⭐⭐⭐ ⭐⭐ ⭐⭐⭐⭐
社区支持 ⭐⭐⭐⭐⭐ ⭐⭐⭐⭐ ⭐⭐⭐⭐

环境搭建与配置

系统要求

推荐配置

  • 操作系统:Linux(Ubuntu 20.04+)或macOS

  • Python版本:Python 3.8+

  • 内存:最低8GB,推荐16GB+

  • 存储:根据知识库大小而定,推荐SSD

基础环境安装

1
2
3
4
5
6
7
8
9
10
11
12
# 创建虚拟环境
python -m venv rag_env

# 激活虚拟环境
# Linux/macOS
source rag_env/bin/activate
# Windows
rag_env\Scripts\activate

# 安装核心依赖
pip install langchain llama-index openai chromadb
pip install pypdf python-docx requests beautifulsoup4

API密钥配置

1
2
3
4
5
6
7
8
9
10
# OpenAI API密钥
export OPENAI_API_KEY="your-api-key"

# Pinecone API密钥(如果使用)
export PINECONE_API_KEY="your-api-key"
export PINECONE_ENV="your-environment"

# Milvus配置(如果使用)
export MILVUS_HOST="localhost"
export MILVUS_PORT="19530"

配置文件示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
# config.yaml
llm:
provider: "openai"
model: "gpt-3.5-turbo"
temperature: 0.7
max_tokens: 1024

embedding:
provider: "openai"
model: "text-embedding-3-small"
dimensions: 1536

vector_store:
provider: "chromadb"
persist_directory: "./chroma_db"
collection_name: "knowledge_base"

chunking:
chunk_size: 1000
chunk_overlap: 200
separator: "\n\n"

retrieval:
top_k: 5
similarity_threshold: 0.7

数据处理与准备

数据源类型

1. 文档文件

  • PDF文档:技术手册、报告、论文

  • Word文档:政策文件、流程文档

  • Markdown文件:技术文档、博客文章

  • TXT文件:日志文件、纯文本数据

2. 结构化数据

  • CSV/Excel:表格数据、统计信息

  • JSON:API响应、配置文件

  • 数据库:关系型数据库、NoSQL数据库

3. 在线数据

  • 网页内容:新闻、博客、官方网站

  • API接口:实时数据、第三方服务

  • 邮件:邮件通信、通知

文档加载器

LangChain加载器示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
from langchain_community.document_loaders import (
PyPDFLoader,
Docx2txtLoader,
TextLoader,
WebBaseLoader,
CSVLoader
)

# PDF加载
pdf_loader = PyPDFLoader("document.pdf")
pdf_docs = pdf_loader.load()

# Word加载
docx_loader = Docx2txtLoader("document.docx")
docx_docs = docx_loader.load()

# 网页加载
web_loader = WebBaseLoader("https://example.com/article")
web_docs = web_loader.load()

# CSV加载
csv_loader = CSVLoader("data.csv")
csv_docs = csv_loader.load()

LlamaIndex加载器示例

1
2
3
4
5
6
7
8
9
from llama_index.core import SimpleDirectoryReader, SimpleWebPageReader

# 目录加载
dir_reader = SimpleDirectoryReader("./documents")
dir_docs = dir_reader.load_data()

# 网页加载
web_reader = SimpleWebPageReader()
web_docs = web_reader.load_data(["https://example.com/page1", "https://example.com/page2"])

文本分块策略

分块原则

  1. 语义完整性:保持段落、句子的完整性

  2. 大小适中:根据LLM上下文窗口调整

  3. 重叠保留:避免关键信息被截断

  4. 类型适配:根据文档类型选择策略

分块技术对比

分块策略 优点 缺点 适用场景
固定大小分块 简单高效 可能破坏语义 通用场景
句子分块 语义完整 块大小不均 结构化文档
段落分块 逻辑清晰 大块处理困难 文章类文档
递归分块 智能适应 计算复杂 混合类型文档

分块代码示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
from langchain_text_splitters import (
RecursiveCharacterTextSplitter,
SentenceSplitter,
TokenTextSplitter
)

# 递归字符分块(推荐)
recursive_splitter = RecursiveCharacterTextSplitter(
chunk_size=1000,
chunk_overlap=200,
length_function=len,
is_separator_regex=False
)

# 句子分块
sentence_splitter = SentenceSplitter(
chunk_size=1000,
chunk_overlap=100
)

# Token分块
token_splitter = TokenTextSplitter(
chunk_size=512,
chunk_overlap=64
)

# 分块处理
chunks = recursive_splitter.split_text(long_document_text)

元数据处理

元数据类型

  • 基础信息:文件名、路径、大小、创建时间

  • 内容信息:标题、作者、摘要、关键词

  • 结构信息:页码、章节、段落编号

  • 自定义标签:分类、优先级、来源

元数据添加示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
from langchain_core.documents import Document

# 为文档添加元数据
document = Document(
page_content="文档内容...",
metadata={
"source": "manual.pdf",
"page": 42,
"category": "technical",
"priority": "high",
"update_time": "2025-10-23"
}
)

# 批量处理
documents_with_metadata = []
for i, doc in enumerate(raw_documents):
documents_with_metadata.append(Document(
page_content=doc.page_content,
metadata={
"source": doc.metadata["source"],
"chunk_id": f"chunk_{i}",
"length": len(doc.page_content),
"processed_time": datetime.now().isoformat()
}
))

向量数据库选型指南

主流向量数据库对比

1. Milvus

特点

  • 开源分布式向量数据库

  • 支持大规模向量检索

  • 丰富的索引类型

优势

  • 高吞吐量和低延迟

  • 支持水平扩展

  • 完善的企业级特性

适用场景

  • 大规模生产环境

  • 需要分布式部署

  • 高并发查询场景

部署方式

1
2
3
# Docker Compose部署
wget https://github.com/milvus-io/milvus/releases/download/v2.4.0/milvus-standalone-docker-compose.yml
docker-compose -f milvus-standalone-docker-compose.yml up -d

2. Pinecone

特点

  • 全托管向量数据库服务

  • 开箱即用,无需运维

  • 支持实时更新

优势

  • 零运维成本

  • 高可用性

  • 简单易用的API

适用场景

  • 快速原型开发

  • 中小规模应用

  • 不想管理基础设施的团队

使用示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
import pinecone
from langchain_community.vectorstores import Pinecone

# 初始化Pinecone
pinecone.init(
api_key="your-api-key",
environment="your-environment"
)

# 创建索引
if "knowledge-base" not in pinecone.list_indexes():
pinecone.create_index(
"knowledge-base",
dimension=1536,
metric="cosine"
)

# 构建向量存储
vectorstore = Pinecone.from_documents(
documents=documents,
embedding=embeddings,
index_name="knowledge-base"
)

3. Chroma

特点

  • 轻量级开源向量数据库

  • 适合开发和测试

  • 支持持久化存储

优势

  • 简单易用

  • 无需额外部署

  • 适合本地开发

适用场景

  • 开发和测试环境

  • 小规模数据集

  • 快速原型验证

4. Weaviate

特点

  • 混合搜索能力(向量+关键词)

  • 内置AI功能

  • 支持GraphQL查询

优势

  • 强大的搜索功能

  • 丰富的查询语言

  • 支持知识图谱构建

选型决策框架

评估维度

1. 数据规模

  • 小规模(<10万向量):Chroma、FAISS

  • 中等规模(10万-1000万):Pinecone、Weaviate

  • 大规模(>1000万):Milvus、Zilliz Cloud

2. 性能要求

  • 低延迟:Milvus、Pinecone

  • 高吞吐量:Milvus、Weaviate

  • 实时更新:Pinecone、Weaviate

3. 运维复杂度

  • 零运维:Pinecone(托管服务)

  • 简单运维:Chroma、Weaviate

  • 专业运维:Milvus

4. 成本预算

  • 免费:Chroma、FAISS(开源)

  • 低成本:Weaviate、自托管Milvus

  • 按需付费:Pinecone、Zilliz Cloud

索引类型选择

HNSW索引(推荐)

1
2
3
4
5
6
7
8
9
# Milvus HNSW索引配置
index_params = {
"index_type": "HNSW",
"metric_type": "L2",
"params": {
"M": 16,
"efConstruction": 200
}
}

IVF_FLAT索引

1
2
3
4
5
6
7
8
9
10
11
12
# FAISS IVF索引配置
import faiss

dimension = 1536
nlist = 100 # 聚类数量

index = faiss.IndexIVFFlat(
faiss.IndexFlatL2(dimension),
dimension,
nlist,
faiss.METRIC_L2
)

RAG系统构建实践

基础RAG实现

使用LangChain构建

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
import os
from langchain_community.document_loaders import PyPDFLoader
from langchain_text_splitters import RecursiveCharacterTextSplitter
from langchain_chroma import Chroma
from langchain_openai import OpenAIEmbeddings, ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.runnables import RunnablePassthrough
from langchain_core.output_parsers import StrOutputParser

class BasicRAGSystem:
def __init__(self, config):
self.config = config
self.embeddings = OpenAIEmbeddings(model=config["embedding"]["model"])
self.llm = ChatOpenAI(
model=config["llm"]["model"],
temperature=config["llm"]["temperature"],
max_tokens=config["llm"]["max_tokens"]
)
self.vectorstore = None
self.rag_chain = None

def load_and_process_documents(self, file_paths):
"""加载和处理文档"""
all_docs = []

for file_path in file_paths:
# 根据文件类型选择加载器
if file_path.endswith(".pdf"):
loader = PyPDFLoader(file_path)
elif file_path.endswith(".docx"):
from langchain_community.document_loaders import Docx2txtLoader
loader = Docx2txtLoader(file_path)
else:
loader = TextLoader(file_path)

docs = loader.load()
all_docs.extend(docs)

# 文本分块
text_splitter = RecursiveCharacterTextSplitter(
chunk_size=self.config["chunking"]["chunk_size"],
chunk_overlap=self.config["chunking"]["chunk_overlap"],
separator=self.config["chunking"]["separator"]
)

return text_splitter.split_documents(all_docs)

def build_vector_store(self, documents):
"""构建向量存储"""
self.vectorstore = Chroma.from_documents(
documents=documents,
embedding=self.embeddings,
persist_directory=self.config["vector_store"]["persist_directory"],
collection_name=self.config["vector_store"]["collection_name"]
)
self.vectorstore.persist()

def create_rag_chain(self):
"""创建RAG链"""
if not self.vectorstore:
raise ValueError("Vector store not initialized. Call build_vector_store first.")

retriever = self.vectorstore.as_retriever(
search_kwargs={
"k": self.config["retrieval"]["top_k"],
"score_threshold": self.config["retrieval"]["similarity_threshold"]
}
)

# 文档格式化函数
def format_docs(docs):
return "\n\n---\n\n".join([f"文档来源: {doc.metadata.get('source', '未知')}\n内容: {doc.page_content}" for doc in docs])

# 提示模板
prompt = ChatPromptTemplate.from_template("""
你是一个专业的AI助手,擅长基于提供的文档内容回答问题。

## 回答规则:
1. 严格基于提供的上下文信息回答,不要编造内容
2. 如果上下文没有相关信息,明确说明"根据提供的文档,无法回答该问题"
3. 引用具体的文档来源,并保持回答的准确性和专业性
4. 组织语言清晰,分点说明复杂问题

## 上下文信息:
{context}

## 用户问题:
{question}

## 回答格式:
回答: [你的回答内容]
""")

# 构建RAG链
self.rag_chain = (
{"context": retriever | format_docs, "question": RunnablePassthrough()}
| prompt
| self.llm
| StrOutputParser()
)

def query(self, question):
"""执行查询"""
if not self.rag_chain:
raise ValueError("RAG chain not created. Call create_rag_chain first.")

return self.rag_chain.invoke(question)

# 使用示例
if __name__ == "__main__":
config = {
"llm": {
"model": "gpt-3.5-turbo",
"temperature": 0.3,
"max_tokens": 1024
},
"embedding": {
"model": "text-embedding-3-small"
},
"vector_store": {
"persist_directory": "./chroma_db",
"collection_name": "knowledge_base"
},
"chunking": {
"chunk_size": 1000,
"chunk_overlap": 200,
"separator": "\n\n"
},
"retrieval": {
"top_k": 5,
"similarity_threshold": 0.7
}
}

# 初始化RAG系统
rag_system = BasicRAGSystem(config)

# 加载文档
documents = rag_system.load_and_process_documents([
"technical_manual.pdf",
"policy_guide.docx",
"faq.txt"
])

# 构建向量存储
rag_system.build_vector_store(documents)

# 创建RAG链
rag_system.create_rag_chain()

# 执行查询
response = rag_system.query("如何配置系统参数?")
print(response)

使用LlamaIndex构建

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
from llama_index.core import (
VectorStoreIndex,
SimpleDirectoryReader,
StorageContext,
load_index_from_storage
)
from llama_index.embeddings.openai import OpenAIEmbedding
from llama_index.llms.openai import OpenAI
from llama_index.core.node_parser import SentenceSplitter
from llama_index.core.prompts import PromptTemplate

class LlamaIndexRAG:
def __init__(self, config):
self.config = config
self.llm = OpenAI(
model=config["llm"]["model"],
temperature=config["llm"]["temperature"],
max_tokens=config["llm"]["max_tokens"]
)
self.embed_model = OpenAIEmbedding(model=config["embedding"]["model"])
self.index = None

def load_documents(self, data_dir):
"""加载文档"""
reader = SimpleDirectoryReader(data_dir)
documents = reader.load_data()
return documents

def build_index(self, documents, persist_dir="./storage"):
"""构建索引"""
# 文本分块
node_parser = SentenceSplitter(
chunk_size=self.config["chunking"]["chunk_size"],
chunk_overlap=self.config["chunking"]["chunk_overlap"]
)

# 构建索引
self.index = VectorStoreIndex.from_documents(
documents,
llm=self.llm,
embed_model=self.embed_model,
node_parser=node_parser
)

# 持久化索引
self.index.storage_context.persist(persist_dir)

def load_existing_index(self, persist_dir="./storage"):
"""加载已存在的索引"""
storage_context = StorageContext.from_defaults(persist_dir=persist_dir)
self.index = load_index_from_storage(
storage_context,
llm=self.llm,
embed_model=self.embed_model
)

def create_query_engine(self):
"""创建查询引擎"""
if not self.index:
raise ValueError("Index not initialized. Call build_index or load_existing_index first.")

# 自定义提示模板
qa_prompt_tmpl = """
你是一个专业的AI助手,基于提供的文档内容回答问题。

## 回答规则:
1. 严格基于提供的上下文信息回答,不要编造内容
2. 如果上下文没有相关信息,明确说明"根据提供的文档,无法回答该问题"
3. 引用具体的文档来源,并保持回答的准确性和专业性
4. 组织语言清晰,分点说明复杂问题

## 上下文信息:
{context_str}

## 用户问题:
{query_str}

## 回答格式:
回答: [你的回答内容]
"""

qa_prompt = PromptTemplate(qa_prompt_tmpl)

return self.index.as_query_engine(
similarity_top_k=self.config["retrieval"]["top_k"],
text_qa_template=qa_prompt
)

def query(self, question):
"""执行查询"""
query_engine = self.create_query_engine()
response = query_engine.query(question)
return str(response)

# 使用示例
if __name__ == "__main__":
config = {
"llm": {
"model": "gpt-3.5-turbo",
"temperature": 0.3,
"max_tokens": 1024
},
"embedding": {
"model": "text-embedding-3-small"
},
"chunking": {
"chunk_size": 1000,
"chunk_overlap": 200
},
"retrieval": {
"top_k": 5
}
}

rag_system = LlamaIndexRAG(config)

# 加载文档并构建索引
documents = rag_system.load_documents("./documents")
rag_system.build_index(documents)

# 执行查询
response = rag_system.query("如何配置系统参数?")
print(response)

高级RAG功能

多轮对话支持

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
from langchain_core.messages import HumanMessage, AIMessage
from langchain_core.chat_history import InMemoryChatMessageHistory
from langchain_core.runnables.history import RunnableWithMessageHistory

class ConversationalRAG:
def __init__(self, basic_rag):
self.basic_rag = basic_rag
self.store = {}

def get_session_history(self, session_id: str):
"""获取会话历史"""
if session_id not in self.store:
self.store[session_id] = InMemoryChatMessageHistory()
return self.store[session_id]

def create_conversational_chain(self):
"""创建对话式RAG链"""
# 对话提示模板
conversational_prompt = ChatPromptTemplate.from_messages([
("system", """
你是一个专业的AI助手,基于提供的文档内容回答问题。

## 回答规则:
1. 严格基于提供的上下文信息回答,不要编造内容
2. 如果上下文没有相关信息,明确说明"根据提供的文档,无法回答该问题"
3. 引用具体的文档来源,并保持回答的准确性和专业性
4. 组织语言清晰,分点说明复杂问题
5. 考虑对话历史,保持回答的连贯性
"""),
("history", "{history}"),
("human", "基于以下上下文回答问题:\n{context}\n\n用户问题:{question}")
])

# 文档格式化函数
def format_docs(docs):
return "\n\n---\n\n".join([f"文档来源: {doc.metadata.get('source', '未知')}\n内容: {doc.page_content}" for doc in docs])

# 构建基础链
base_chain = (
{"context": self.basic_rag.vectorstore.as_retriever() | format_docs,
"question": RunnablePassthrough(),
"history": RunnablePassthrough()}
| conversational_prompt
| self.basic_rag.llm
| StrOutputParser()
)

# 添加对话历史支持
self.conversational_rag_chain = RunnableWithMessageHistory(
base_chain,
self.get_session_history,
input_messages_key="question",
history_messages_key="history",
)

def chat(self, session_id, question):
"""对话式查询"""
if not hasattr(self, 'conversational_rag_chain'):
self.create_conversational_chain()

response = self.conversational_rag_chain.invoke(
{"question": question},
config={"configurable": {"session_id": session_id}}
)

return response

查询重写优化

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
class QueryRewriter:
def __init__(self, llm):
self.llm = llm

def rewrite_query(self, question, history=None):
"""重写查询以提升检索效果"""
history_text = ""
if history:
history_text = "\n".join([f"用户: {h.content}" if isinstance(h, HumanMessage)
else f"助手: {h.content}" for h in history])

prompt = ChatPromptTemplate.from_template("""
你是一个查询优化专家,需要将用户的问题重写为更适合检索的形式。

## 任务:
1. 分析用户问题,识别核心意图
2. 扩展相关关键词和同义词
3. 重写为更精确、更完整的查询
4. 如果有对话历史,考虑上下文信息

## 对话历史(如果有):
{history}

## 用户原始问题:
{question}

## 输出格式:
优化后的查询: [重写后的查询]
""")

response = self.llm.invoke(prompt.format(
history=history_text,
question=question
))

# 提取优化后的查询
return response.content.split("优化后的查询: ")[-1].strip()

# 在RAG系统中集成查询重写
class EnhancedRAGSystem(BasicRAGSystem):
def __init__(self, config):
super().__init__(config)
self.query_rewriter = QueryRewriter(self.llm)

def query(self, question, session_id=None):
"""增强版查询,包含查询重写"""
# 查询重写
rewritten_question = self.query_rewriter.rewrite_query(question)
print(f"原始查询: {question}")
print(f"优化查询: {rewritten_question}")

# 执行优化后的查询
return super().query(rewritten_question)

性能优化策略

检索性能优化

1. 索引优化

HNSW索引调优

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
# Milvus HNSW索引优化参数
index_params = {
"index_type": "HNSW",
"metric_type": "L2",
"params": {
"M": 32, # 每个节点的最大邻居数
"efConstruction": 400, # 构建时的搜索范围
"ef": 100 # 查询时的搜索范围
}
}

# 动态调整参数
def optimize_index_params(data_size):
if data_size < 10000:
return {"M": 16, "efConstruction": 200, "ef": 50}
elif data_size < 100000:
return {"M": 32, "efConstruction": 400, "ef": 100}
else:
return {"M": 64, "efConstruction": 800, "ef": 200}

混合搜索策略

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
from langchain.retrievers import EnsembleRetriever
from langchain_community.retrievers import BM25Retriever

class HybridRetriever:
def __init__(self, vectorstore, documents):
# 向量检索器
self.vector_retriever = vectorstore.as_retriever()

# BM25检索器
self.bm25_retriever = BM25Retriever.from_documents(documents)
self.bm25_retriever.k = 5

# 集成检索器
self.ensemble_retriever = EnsembleRetriever(
retrievers=[self.vector_retriever, self.bm25_retriever],
weights=[0.7, 0.3] # 权重分配
)

def retrieve(self, query):
return self.ensemble_retriever.get_relevant_documents(query)

2. 查询优化

查询过滤与路由

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
class QueryRouter:
def __init__(self, llm, vectorstore, keyword_store):
self.llm = llm
self.vectorstore = vectorstore
self.keyword_store = keyword_store

def route_query(self, query):
"""根据查询类型选择合适的检索器"""
# 分析查询类型
query_type = self._classify_query(query)

if query_type == "factual":
# 事实性查询:使用向量检索
return self.vectorstore.as_retriever().get_relevant_documents(query)
elif query_type == "keyword":
# 关键词查询:使用BM25
return self.keyword_store.search(query)
else:
# 混合查询:使用集成检索
vector_results = self.vectorstore.as_retriever().get_relevant_documents(query)
keyword_results = self.keyword_store.search(query)
return self._merge_results(vector_results, keyword_results)

def _classify_query(self, query):
"""分类查询类型"""
prompt = ChatPromptTemplate.from_template("""
将用户查询分类为以下类型之一:
- factual: 需要具体事实信息的查询
- keyword: 适合关键词搜索的查询
- conversational: 对话式查询

查询: {query}
分类:
""")

response = self.llm.invoke(prompt.format(query=query))
return response.content.strip().lower()

3. 缓存策略

多级缓存实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
from functools import lru_cache
from datetime import datetime, timedelta
import hashlib

class RAGCache:
def __init__(self, ttl=3600):
self.ttl = ttl # 缓存有效期(秒)
self.memory_cache = {}
self.disk_cache = {}

def get_cache_key(self, query, params):
"""生成缓存key"""
cache_input = f"{query}_{str(sorted(params.items()))}"
return hashlib.md5(cache_input.encode()).hexdigest()

def get(self, query, params):
"""获取缓存"""
key = self.get_cache_key(query, params)

# 检查内存缓存
if key in self.memory_cache:
cache_entry = self.memory_cache[key]
if datetime.now() - cache_entry["timestamp"] < timedelta(seconds=self.ttl):
return cache_entry["data"]

# 检查磁盘缓存(如果实现)
if key in self.disk_cache:
cache_entry = self.disk_cache[key]
if datetime.now() - cache_entry["timestamp"] < timedelta(seconds=self.ttl):
return cache_entry["data"]

return None

def set(self, query, params, data):
"""设置缓存"""
key = self.get_cache_key(query, params)

# 更新内存缓存
self.memory_cache[key] = {
"data": data,
"timestamp": datetime.now()
}

# 更新磁盘缓存(如果实现)
self.disk_cache[key] = {
"data": data,
"timestamp": datetime.now()
}

def clear_expired(self):
"""清理过期缓存"""
now = datetime.now()
self.memory_cache = {k: v for k, v in self.memory_cache.items()
if now - v["timestamp"] < timedelta(seconds=self.ttl)}
self.disk_cache = {k: v for k, v in self.disk_cache.items()
if now - v["timestamp"] < timedelta(seconds=self.ttl)}

生成性能优化

1. 模型优化

量化与蒸馏

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
from langchain_openai import ChatOpenAI
from transformers import AutoModelForCausalLM, AutoTokenizer, BitsAndBytesConfig

class ModelOptimizer:
@staticmethod
def load_optimized_model(model_name, quantize=True):
"""加载优化的模型"""
if "gpt" in model_name or "claude" in model_name:
# 闭源模型使用API
return ChatOpenAI(model=model_name, temperature=0.3)
else:
# 开源模型使用本地部署
if quantize:
# 4位量化配置
bnb_config = BitsAndBytesConfig(
load_in_4bit=True,
bnb_4bit_use_double_quant=True,
bnb_4bit_quant_type="nf4",
bnb_4bit_compute_dtype=torch.bfloat16
)

model = AutoModelForCausalLM.from_pretrained(
model_name,
quantization_config=bnb_config,
device_map="auto"
)
else:
model = AutoModelForCausalLM.from_pretrained(
model_name,
device_map="auto"
)

tokenizer = AutoTokenizer.from_pretrained(model_name)
return model, tokenizer

2. 上下文管理

动态上下文窗口

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
class ContextManager:
def __init__(self, max_context_length=8000):
self.max_context_length = max_context_length

def optimize_context(self, documents, query):
"""优化上下文长度"""
# 计算可用上下文空间
query_length = len(query)
available_space = self.max_context_length - query_length - 1000 # 预留空间

# 按相关性排序
sorted_docs = sorted(documents, key=lambda x: x.metadata.get('score', 0), reverse=True)

optimized_context = []
current_length = 0

for doc in sorted_docs:
doc_length = len(doc.page_content)

if current_length + doc_length <= available_space:
optimized_context.append(doc)
current_length += doc_length
else:
# 如果还有空间,添加文档摘要
if available_space - current_length > 200:
summary = self._summarize_document(doc, available_space - current_length)
optimized_context.append(Document(
page_content=summary,
metadata={**doc.metadata, "summary": True}
))
break

return optimized_context

def _summarize_document(self, document, max_length):
"""生成文档摘要"""
# 使用LLM生成摘要
prompt = f"请将以下内容总结为不超过{max_length}字符的摘要:\n\n{document.page_content}"
# 调用LLM生成摘要...
return summary

系统级优化

1. 异步处理

异步RAG实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
import asyncio
from langchain_community.document_loaders import AsyncChromiumLoader
from langchain_community.document_transformers import BeautifulSoupTransformer

class AsyncRAGSystem:
def __init__(self, config):
self.config = config
self.embeddings = OpenAIEmbeddings(model=config["embedding"]["model"])
self.llm = ChatOpenAI(model=config["llm"]["model"])

async def load_documents_async(self, urls):
"""异步加载网页文档"""
loader = AsyncChromiumLoader(urls)
docs = await loader.aload()

# 转换HTML为文本
bs_transformer = BeautifulSoupTransformer()
docs_transformed = bs_transformer.transform_documents(
docs, tags_to_extract=["p", "h1", "h2", "h3"]
)

return docs_transformed

async def process_documents_batch(self, file_paths, batch_size=5):
"""批量处理文档"""
tasks = []
for i in range(0, len(file_paths), batch_size):
batch = file_paths[i:i+batch_size]
task = self._process_batch(batch)
tasks.append(task)

results = await asyncio.gather(*tasks)
return [doc for batch_result in results for doc in batch_result]

async def _process_batch(self, batch):
"""处理单个批次"""
# 实现批次处理逻辑
pass

2. 分布式部署

Kubernetes部署配置

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
# rag-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: rag-service
spec:
replicas: 3
selector:
matchLabels:
app: rag-service
template:
metadata:
labels:
app: rag-service
spec:
containers:
- name: rag-api
image: rag-service:latest
ports:
- containerPort: 8000
resources:
limits:
cpu: "2"
memory: "4Gi"
requests:
cpu: "1"
memory: "2Gi"
env:
- name: OPENAI_API_KEY
valueFrom:
secretKeyRef:
name: api-keys
key: openai-api-key
- name: VECTOR_STORE_HOST
value: "vector-store-service"
livenessProbe:
httpGet:
path: /health
port: 8000
initialDelaySeconds: 30
periodSeconds: 10
readinessProbe:
httpGet:
path: /ready
port: 8000
initialDelaySeconds: 5
periodSeconds: 5

部署与监控

部署架构

1. 单机部署(开发环境)

1
2
3
4
5
6
7
┌─────────────────────────────────────┐
│ RAG Service │
│ ┌─────────┐ ┌─────────┐ ┌──────┐ │
│ │ LLM │ │VectorDB │ │ API │ │
│ │Module │ │(Chroma) │ │Server│ │
│ └─────────┘ └─────────┘ └──────┘ │
└─────────────────────────────────────┘

2. 分布式部署(生产环境)

1
2
3
4
5
6
7
8
9
10
┌─────────────┐    ┌─────────────┐    ┌─────────────┐
│ Load │ │ RAG API │ │ VectorDB │
│ Balancer │───>│ Servers │───>│ Cluster │
└─────────────┘ └─────────────┘ └─────────────┘
│ │ │
│ │ │
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ LLM │ │ Monitoring │ │ Backup & │
│ Service │<───│ & Logging │<───│ Restore │
└─────────────┘ └─────────────┘ └─────────────┘

API服务实现

FastAPI服务

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
from fastapi import FastAPI, HTTPException, Depends, Query
from fastapi.middleware.cors import CORSMiddleware
from pydantic import BaseModel
from typing import List, Optional, Dict
import logging
import time
from contextlib import asynccontextmanager

# 配置日志
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s - %(name)s - %(levelname)s - %(message)s"
)
logger = logging.getLogger(__name__)

class RAGService:
"""RAG服务管理类"""
def __init__(self):
self.rag_system = None
self.is_ready = False

async def initialize(self):
"""初始化RAG系统"""
try:
# 加载配置
config = self._load_config()

# 初始化RAG系统
self.rag_system = BasicRAGSystem(config)

# 加载现有索引或构建新索引
if os.path.exists(config["vector_store"]["persist_directory"]):
logger.info("加载现有向量存储...")
self.rag_system.load_existing_vector_store()
else:
logger.info("构建新的向量存储...")
documents = self.rag_system.load_and_process_documents(
config["data"]["document_paths"]
)
self.rag_system.build_vector_store(documents)

self.rag_system.create_rag_chain()
self.is_ready = True
logger.info("RAG服务初始化完成")

except Exception as e:
logger.error(f"RAG服务初始化失败: {str(e)}", exc_info=True)
raise

# 创建全局服务实例
rag_service = RAGService()

@asynccontextmanager
async def lifespan(app: FastAPI):
"""应用生命周期管理"""
# 启动时初始化
await rag_service.initialize()
yield
# 关闭时清理
logger.info("RAG服务关闭")

# 创建FastAPI应用
app = FastAPI(
title="RAG API Service",
description="检索增强生成(RAG)服务API",
version="1.0.0",
lifespan=lifespan
)

# 配置CORS
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)

# 数据模型
class QueryRequest(BaseModel):
question: str
session_id: Optional[str] = None
top_k: Optional[int] = 5
temperature: Optional[float] = 0.3

class QueryResponse(BaseModel):
answer: str
sources: List[Dict]
response_time: float
session_id: Optional[str] = None

class HealthCheckResponse(BaseModel):
status: str
ready: bool
version: str

# API端点
@app.get("/health", response_model=HealthCheckResponse, tags=["系统"])
async def health_check():
"""健康检查"""
return {
"status": "running",
"ready": rag_service.is_ready,
"version": "1.0.0"
}

@app.post("/query", response_model=QueryResponse, tags=["RAG"])
async def query_rag(request: QueryRequest):
"""RAG查询接口"""
if not rag_service.is_ready:
raise HTTPException(status_code=503, detail="RAG服务尚未准备就绪")

start_time = time.time()

try:
# 执行查询
response = rag_service.rag_system.query(
request.question,
session_id=request.session_id,
top_k=request.top_k,
temperature=request.temperature
)

# 解析结果
answer = response["answer"]
sources = response["sources"]

response_time = time.time() - start_time

logger.info(f"查询完成 - 问题: {request.question[:50]}... "
f"响应时间: {response_time:.2f}s")

return {
"answer": answer,
"sources": sources,
"response_time": response_time,
"session_id": request.session_id
}

except Exception as e:
logger.error(f"查询处理失败: {str(e)}", exc_info=True)
raise HTTPException(status_code=500, detail=str(e))

@app.post("/ingest", tags=["数据管理"])
async def ingest_documents(file_paths: List[str]):
"""文档摄入接口"""
if not rag_service.is_ready:
raise HTTPException(status_code=503, detail="RAG服务尚未准备就绪")

try:
documents = rag_service.rag_system.load_and_process_documents(file_paths)
rag_service.rag_system.add_documents(documents)

return {
"status": "success",
"message": f"成功摄入 {len(documents)} 个文档块",
"file_paths": file_paths
}

except Exception as e:
logger.error(f"文档摄入失败: {str(e)}", exc_info=True)
raise HTTPException(status_code=500, detail=str(e))

if __name__ == "__main__":
import uvicorn
uvicorn.run("rag_api:app", host="0.0.0.0", port=8000, reload=True)

监控系统

Prometheus监控指标

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
from prometheus_client import Counter, Histogram, Gauge, generate_latest
from prometheus_client.core import CollectorRegistry

class RAGMetrics:
def __init__(self):
self.registry = CollectorRegistry()

# 计数器
self.query_counter = Counter(
"rag_queries_total",
"Total number of RAG queries",
["status", "model"],
registry=self.registry
)

self.document_counter = Counter(
"rag_documents_total",
"Total number of documents processed",
["action", "source_type"],
registry=self.registry
)

# 直方图
self.response_time_histogram = Histogram(
"rag_response_time_seconds",
"RAG query response time",
["model"],
registry=self.registry,
buckets=[0.1, 0.5, 1.0, 2.0, 5.0, 10.0, 30.0]
)

self.document_size_histogram = Histogram(
"rag_document_size_bytes",
"Size of processed documents",
["document_type"],
registry=self.registry
)

# 仪表
self.active_sessions_gauge = Gauge(
"rag_active_sessions",
"Number of active RAG sessions",
registry=self.registry
)

self.vector_count_gauge = Gauge(
"rag_vector_count",
"Number of vectors in storage",
["collection"],
registry=self.registry
)

def record_query(self, status="success", model="gpt-3.5-turbo", response_time=0):
"""记录查询指标"""
self.query_counter.labels(status=status, model=model).inc()
self.response_time_histogram.labels(model=model).observe(response_time)

def record_document(self, action="ingest", source_type="pdf", size=0):
"""记录文档处理指标"""
self.document_counter.labels(action=action, source_type=source_type).inc()
self.document_size_histogram.labels(document_type=source_type).observe(size)

def update_vector_count(self, count, collection="default"):
"""更新向量计数"""
self.vector_count_gauge.labels(collection=collection).set(count)

日志系统配置

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
import logging
from logging.handlers import RotatingFileHandler
import json
from datetime import datetime

class JSONFormatter(logging.Formatter):
"""JSON格式日志"""
def format(self, record):
log_entry = {
"timestamp": datetime.utcnow().isoformat(),
"level": record.levelname,
"logger": record.name,
"message": record.getMessage(),
"module": record.module,
"function": record.funcName,
"line": record.lineno
}

# 添加额外字段
if hasattr(record, "user_id"):
log_entry["user_id"] = record.user_id
if hasattr(record, "session_id"):
log_entry["session_id"] = record.session_id
if hasattr(record, "query"):
log_entry["query"] = record.query
if hasattr(record, "response_time"):
log_entry["response_time"] = record.response_time

return json.dumps(log_entry)

def setup_logging(log_file="rag_service.log", max_bytes=10*1024*1024, backup_count=5):
"""配置日志系统"""
logger = logging.getLogger()
logger.setLevel(logging.INFO)

# 移除现有处理器
for handler in logger.handlers[:]:
logger.removeHandler(handler)

# 控制台处理器
console_handler = logging.StreamHandler()
console_handler.setFormatter(JSONFormatter())
logger.addHandler(console_handler)

# 文件处理器(轮转)
file_handler = RotatingFileHandler(
log_file,
maxBytes=max_bytes,
backupCount=backup_count,
encoding="utf-8"
)
file_handler.setFormatter(JSONFormatter())
logger.addHandler(file_handler)

return logger

部署脚本

Docker部署

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
# Dockerfile
FROM python:3.11-slim

WORKDIR /app

# 安装系统依赖
RUN apt-get update && apt-get install -y --no-install-recommends \
build-essential \
&& rm -rf /var/lib/apt/lists/*

# 安装Python依赖
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

# 复制应用代码
COPY . .

# 创建数据目录
RUN mkdir -p /app/data /app/chroma_db /app/logs

# 环境变量
ENV PYTHONUNBUFFERED=1
ENV PYTHONDONTWRITEBYTECODE=1
ENV LOG_LEVEL=INFO

# 暴露端口
EXPOSE 8000

# 健康检查
HEALTHCHECK --interval=30s --timeout=10s --start-period=60s --retries=3 \
CMD curl -f http://localhost:8000/health || exit 1

# 启动命令
CMD ["uvicorn", "rag_api:app", "--host", "0.0.0.0", "--port", "8000", "--workers", "4"]
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
# docker-compose.yml
version: '3.8'

services:
rag-service:
build: .
ports:
- "8000:8000"
environment:
- OPENAI_API_KEY=${OPENAI_API_KEY}
- VECTOR_STORE_TYPE=chroma
- CHROMA_PERSIST_DIRECTORY=/app/chroma_db
- LOG_LEVEL=INFO
volumes:
- ./data:/app/data
- ./chroma_db:/app/chroma_db
- ./logs:/app/logs
restart: unless-stopped
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:8000/health"]
interval: 30s
timeout: 10s
retries: 3
start_period: 60s

prometheus:
image: prom/prometheus:latest
ports:
- "9090:9090"
volumes:
- ./prometheus.yml:/etc/prometheus/prometheus.yml
- prometheus_data:/prometheus
command:
- '--config.file=/etc/prometheus/prometheus.yml'
- '--storage.tsdb.path=/prometheus'
- '--web.console.libraries=/etc/prometheus/console_libraries'
- '--web.console.templates=/etc/prometheus/consoles'
- '--web.enable-lifecycle'

grafana:
image: grafana/grafana:latest
ports:
- "3000:3000"
environment:
- GF_SECURITY_ADMIN_PASSWORD=admin
volumes:
- grafana_data:/var/lib/grafana
depends_on:
- prometheus

volumes:
prometheus_data:
grafana_data:

最佳实践与案例

企业知识库案例

案例背景

某大型制造企业需要构建智能知识库系统,帮助员工快速获取技术文档、政策文件和操作指南。

技术方案

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
class EnterpriseKnowledgeBase:
def __init__(self):
self.config = {
"llm": {
"model": "gpt-4-turbo",
"temperature": 0.2,
"max_tokens": 2048
},
"embedding": {
"model": "text-embedding-3-large"
},
"vector_store": {
"provider": "milvus",
"host": "milvus-cluster",
"port": 19530,
"collection_name": "enterprise_knowledge"
},
"chunking": {
"chunk_size": 1500,
"chunk_overlap": 200
},
"retrieval": {
"top_k": 8,
"similarity_threshold": 0.65
}
}

self.rag_system = EnhancedRAGSystem(self.config)
self.setup_domain_specific_features()

def setup_domain_specific_features(self):
"""设置领域特定功能"""
# 添加制造业专业术语词典
self.term_dictionary = self._load_manufacturing_terms()

# 设置文档分类器
self.document_classifier = self._create_document_classifier()

def _load_manufacturing_terms(self):
"""加载制造业专业术语"""
# 从文件或数据库加载专业术语
return {
"CNC": "计算机数控",
"PLC": "可编程逻辑控制器",
"CAD": "计算机辅助设计",
# 更多术语...
}

def process_technical_document(self, file_path):
"""处理技术文档"""
# 1. 文档类型识别
doc_type = self.document_classifier.classify(file_path)

# 2. 根据文档类型选择处理策略
if doc_type == "maintenance_manual":
return self._process_maintenance_manual(file_path)
elif doc_type == "engineering_drawing":
return self._process_engineering_drawing(file_path)
elif doc_type == "quality_standard":
return self._process_quality_standard(file_path)
else:
return self._process_generic_document(file_path)

def _process_maintenance_manual(self, file_path):
"""处理维护手册"""
# 特殊的分块策略:按维修步骤分块
maintenance_splitter = MaintenanceManualSplitter()
documents = maintenance_splitter.split(file_path)

# 添加维护相关元数据
for doc in documents:
doc.metadata.update({
"document_type": "maintenance",
"priority": "high",
"requires_validation": True
})

return documents

效果评估

  • 知识检索准确率:从65%提升至92%

  • 员工培训时间:减少40%

  • 问题解决效率:提升60%

  • 文档更新周期:从月级缩短至日级

智能客服系统案例

系统架构

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
┌─────────────┐    ┌─────────────┐    ┌─────────────┐
│ 用户 │ │ 前端界面 │ │ API网关 │
│ 交互 │───>│ │───>│ │
└─────────────┘ └─────────────┘ └─────────────┘


┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ 人工坐席 │<───│ 对话管理 │<───│ RAG服务 │
│ 系统 │ │ 模块 │ │ │
└─────────────┘ └─────────────┘ └─────────────┘


┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ 客户数据 │ │ 产品知识 │ │ 历史对话 │
│ 库 │<───│ 库 │<───│ 库 │
└─────────────┘ └─────────────┘ └─────────────┘

核心功能实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
class IntelligentCustomerService:
def __init__(self):
self.rag_system = RAGSystem()
self.dialogue_manager = DialogueManager()
self.customer_profile_service = CustomerProfileService()
self.escalation_service = EscalationService()

async def handle_customer_query(self, customer_id, query, session_id=None):
"""处理客户查询"""
# 1. 获取客户信息
customer_info = await self.customer_profile_service.get_customer_info(customer_id)

# 2. 获取对话历史
dialogue_history = self.dialogue_manager.get_history(session_id)

# 3. 个性化查询处理
personalized_query = self._personalize_query(query, customer_info, dialogue_history)

# 4. RAG查询
rag_response = self.rag_system.query(
personalized_query,
session_id=session_id,
customer_info=customer_info
)

# 5. 对话状态管理
self.dialogue_manager.update_history(
session_id,
user_message=query,
ai_response=rag_response["answer"]
)

# 6. 升级判断
if self._should_escalate(rag_response, dialogue_history):
escalation_result = await self.escalation_service.escalate_to_human_agent(
customer_id,
query,
rag_response,
dialogue_history
)
return escalation_result

return rag_response

def _personalize_query(self, query, customer_info, history):
"""个性化查询"""
# 根据客户信息和历史对话优化查询
personalization_prompt = f"""
基于客户信息和对话历史,优化查询以提升相关性:

客户信息:
{customer_info}

对话历史:
{history}

原始查询:{query}

优化后的查询:
"""

return self.rag_system.llm.generate(personalization_prompt)

def _should_escalate(self, response, history):
"""判断是否需要升级到人工坐席"""
# 基于多种因素判断是否需要人工干预
confidence_score = response.get("confidence_score", 0)
query_complexity = self._assess_query_complexity(history[-1]["user_message"])
customer_value = self._get_customer_value_score(response["customer_info"])

# 升级规则
if confidence_score < 0.7:
return True
if query_complexity > 0.8 and customer_value > 0.7:
return True
if len(history) > 3 and not self._is_progress_being_made(history):
return True

return False

关键性能指标

  • 自助解决率:从35%提升至78%

  • 平均响应时间:从2分钟缩短至15秒

  • 客户满意度:从72%提升至91%

  • 人工坐席工作量:减少55%

金融投研分析案例

数据处理流程

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
class FinancialResearchAssistant:
def __init__(self):
self.market_data_provider = MarketDataProvider()
self.news_analyzer = NewsAnalyzer()
self.report_generator = ReportGenerator()
self.rag_system = FinancialRAGSystem()

async def analyze_stock(self, ticker, time_period="1Y", analysis_type="comprehensive"):
"""分析股票"""
# 1. 获取实时市场数据
market_data = await self.market_data_provider.get_stock_data(
ticker, time_period=time_period
)

# 2. 获取相关新闻
news_articles = await self.news_analyzer.get_relevant_news(
ticker, limit=20, time_period=time_period
)

# 3. 获取财报数据
financial_reports = await self._get_financial_reports(ticker)

# 4. 整合所有数据
analysis_context = self._integrate_data(
market_data, news_articles, financial_reports
)

# 5. 执行RAG分析
analysis_query = self._generate_analysis_query(ticker, analysis_type)
rag_analysis = self.rag_system.analyze(analysis_query, analysis_context)

# 6. 生成专业报告
report = self.report_generator.generate_stock_report(
ticker, rag_analysis, market_data, news_articles
)

return report

def _generate_analysis_query(self, ticker, analysis_type):
"""生成分析查询"""
query_templates = {
"comprehensive": f"""
请对{ticker}进行全面分析,包括:
1. 基本面分析(营收、利润、现金流等)
2. 技术面分析(价格走势、交易量、技术指标)
3. 行业对比分析
4. 风险评估
5. 投资建议

基于提供的所有数据,给出专业的投资分析报告。
""",

"risk_assessment": f"""
请对{ticker}进行风险评估,重点分析:
1. 市场风险(系统性风险敞口)
2. 行业风险(竞争格局变化)
3. 公司特定风险(财务风险、管理风险)
4. 宏观经济风险敏感性
5. 风险缓解措施建议
""",

"valuation": f"""
请对{ticker}进行估值分析,包括:
1. 绝对估值(DCF模型)
2. 相对估值(PE、PB、EV/EBITDA等)
3. 与历史估值对比
4. 与同行业对比
5. 目标价格区间
"""
}

return query_templates.get(analysis_type, query_templates["comprehensive"])

def _integrate_data(self, market_data, news, reports):
"""整合多种数据源"""
# 处理市场数据
market_summary = self._summarize_market_data(market_data)

# 分析新闻情感
news_summary = self.news_analyzer.analyze_news_sentiment(news)

# 提取财报关键指标
financial_summary = self._extract_financial_highlights(reports)

return {
"market_data": market_summary,
"news_analysis": news_summary,
"financial_reports": financial_summary,
"timestamp": datetime.now().isoformat()
}

2025年技术趋势

1. 多模态RAG

技术原理

多模态RAG将文本、图像、音频、视频等多种模态数据整合到RAG框架中,提供更全面的信息检索和生成能力。

实现方案

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
class MultimodalRAGSystem:
def __init__(self):
self.text_embedding = TextEmbeddingModel()
self.image_embedding = ImageEmbeddingModel()
self.audio_embedding = AudioEmbeddingModel()
self.multimodal_vectorstore = MultimodalVectorStore()
self.mllm = MultimodalLLM()

def process_multimodal_document(self, document_path):
"""处理多模态文档"""
documents = []

if document_path.endswith((".jpg", ".png", ".gif")):
# 处理图像
image_embedding = self.image_embedding.embed_image(document_path)
image_caption = self._generate_image_caption(document_path)

documents.append(MultimodalDocument(
content_type="image",
content=document_path,
embedding=image_embedding,
metadata={
"caption": image_caption,
"file_type": os.path.splitext(document_path)[1]
}
))

elif document_path.endswith((".mp3", ".wav", ".flac")):
# 处理音频
audio_text = self._transcribe_audio(document_path)
audio_embedding = self.audio_embedding.embed_audio(document_path)

documents.append(MultimodalDocument(
content_type="audio",
content=audio_text,
embedding=audio_embedding,
metadata={
"transcript": audio_text,
"duration": self._get_audio_duration(document_path)
}
))

elif document_path.endswith((".mp4", ".avi", ".mov")):
# 处理视频
video_frames = self._extract_video_frames(document_path)
video_audio = self._extract_video_audio(document_path)

# 处理视频帧
frame_embeddings = [self.image_embedding.embed_image(frame)
for frame in video_frames[:5]] # 取前5帧

# 处理音频
audio_text = self._transcribe_audio(video_audio)
audio_embedding = self.audio_embedding.embed_audio(video_audio)

# 融合嵌入
combined_embedding = self._fuse_embeddings(frame_embeddings + [audio_embedding])

documents.append(MultimodalDocument(
content_type="video",
content=audio_text,
embedding=combined_embedding,
metadata={
"transcript": audio_text,
"frame_count": len(video_frames),
"duration": self._get_video_duration(document_path)
}
))

return documents

def query_multimodal(self, query, modalities=["text", "image", "audio", "video"]):
"""多模态查询"""
# 生成多模态查询嵌入
query_embeddings = {}

if "text" in modalities:
query_embeddings["text"] = self.text_embedding.embed_text(query)
if "image" in modalities:
# 如果查询包含图像,生成图像嵌入
pass
if "audio" in modalities:
# 如果查询包含音频,生成音频嵌入
pass

# 多模态检索
results = self.multimodal_vectorstore.search(
query_embeddings,
modalities=modalities,
top_k=10
)

# 多模态上下文构建
multimodal_context = self._build_multimodal_context(results)

# 多模态生成
response = self.mllm.generate(query, multimodal_context)

return response

2. Agentic RAG

智能代理架构

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
class RAGAgent:
def __init__(self):
self.rag_system = AdvancedRAGSystem()
self.tool_manager = ToolManager()
self.task_planner = TaskPlanner()
self.memory_manager = MemoryManager()
self.reflection_engine = ReflectionEngine()

async def process_query(self, query, session_id=None):
"""处理复杂查询"""
# 1. 任务理解和规划
task_plan = self.task_planner.plan_task(query)

# 2. 执行任务计划
execution_results = []

for step in task_plan["steps"]:
if step["type"] == "rag_query":
# RAG查询步骤
result = await self._execute_rag_step(step, session_id)
elif step["type"] == "tool_call":
# 工具调用步骤
result = await self._execute_tool_step(step)
elif step["type"] == "reasoning":
# 推理步骤
result = self._execute_reasoning_step(step)
else:
# 默认RAG查询
result = await self._execute_rag_step(step, session_id)

execution_results.append(result)

# 检查是否需要调整计划
if self._needs_plan_adjustment(result, step):
task_plan = self.task_planner.adjust_plan(task_plan, result, step)
break

# 3. 结果整合
final_response = self._integrate_results(execution_results, task_plan)

# 4. 反思和学习
self.reflection_engine.learn_from_execution(
query, task_plan, execution_results, final_response
)

return final_response

async def _execute_rag_step(self, step, session_id):
"""执行RAG查询步骤"""
# 增强查询
enhanced_query = self._enhance_query(step["query"], session_id)

# 执行查询
response = self.rag_system.query(
enhanced_query,
session_id=session_id,
**step.get("parameters", {})
)

# 添加反思信息
response["reflection"] = self.reflection_engine.reflect_on_rag_response(
step["query"], response
)

return response

def _enhance_query(self, query, session_id):
"""增强查询"""
# 获取对话历史
history = self.memory_manager.get_conversation_history(session_id)

# 查询重写
rewritten_query = self.rag_system.query_rewriter.rewrite_query(query, history)

# 查询扩展
expanded_queries = self.rag_system.query_expander.expand_query(rewritten_query)

return {
"original": query,
"rewritten": rewritten_query,
"expanded": expanded_queries
}

3. 实时RAG

实时数据处理

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
class RealTimeRAGSystem:
def __init__(self):
self.static_vectorstore = StaticVectorStore()
self.real_time_vectorstore = RealTimeVectorStore()
self.data_stream_processor = DataStreamProcessor()
self.cache_manager = CacheManager()

async def initialize(self):
"""初始化系统"""
# 启动实时数据处理器
self.data_stream_processor.start()

# 订阅实时数据源
self.data_stream_processor.subscribe("news", self._process_news_stream)
self.data_stream_processor.subscribe("market_data", self._process_market_stream)
self.data_stream_processor.subscribe("social_media", self._process_social_media_stream)

async def _process_news_stream(self, news_item):
"""处理新闻流"""
# 实时处理新闻
processed_news = self._process_news_item(news_item)

# 更新实时向量存储
await self.real_time_vectorstore.update(
documents=[processed_news],
ttl=3600*24*7 # 新闻保留7天
)

async def query(self, question, include_real_time=True, real_time_weight=0.3):
"""实时查询"""
start_time = time.time()

# 1. 静态数据检索
static_results = self.static_vectorstore.search(question, top_k=10)

# 2. 实时数据检索
real_time_results = []
if include_real_time:
real_time_results = await self.real_time_vectorstore.search(question, top_k=5)

# 3. 结果融合
fused_results = self._fuse_results(
static_results,
real_time_results,
real_time_weight=real_time_weight
)

# 4. 生成回答
response = self._generate_response(question, fused_results)

# 5. 缓存结果
self.cache_manager.cache_response(
question, response, ttl=300 # 缓存5分钟
)

response_time = time.time() - start_time
return {
"answer": response,
"sources": fused_results,
"response_time": response_time,
"real_time_included": include_real_time
}

def _fuse_results(self, static_results, real_time_results, real_time_weight=0.3):
"""融合静态和实时结果"""
# 为实时结果添加时间衰减权重
current_time = time.time()
weighted_real_time = []

for result in real_time_results:
time_diff = current_time - result["timestamp"]
time_decay = max(0.1, 1 - (time_diff / 3600)) # 1小时内权重从1衰减到0.1

weighted_score = result["score"] * time_decay * real_time_weight
weighted_real_time.append({
**result,
"weighted_score": weighted_score,
"source_type": "real_time"
})

# 为静态结果添加权重
weighted_static = [{
**result,
"weighted_score": result["score"] * (1 - real_time_weight),
"source_type": "static"
} for result in static_results]

# 合并并排序
all_results = weighted_static + weighted_real_time
return sorted(all_results, key=lambda x: x["weighted_score"], reverse=True)

常见问题与解决方案

检索相关问题

问题1:检索结果相关性低

症状:返回的文档与查询问题相关性不高

解决方案

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
def optimize_retrieval():
# 1. 优化嵌入模型
better_embedding = OpenAIEmbeddings(model="text-embedding-3-large")

# 2. 调整分块策略
optimized_splitter = RecursiveCharacterTextSplitter(
chunk_size=800,
chunk_overlap=100,
separators=["\n\n", "\n", ". ", " ", ""]
)

# 3. 改进检索参数
retriever = vectorstore.as_retriever(
search_kwargs={
"k": 10,
"score_threshold": 0.6,
"fetch_k": 50 # 先获取更多结果再过滤
}
)

# 4. 添加重排序
from langchain.retrievers import ContextualCompressionRetriever
from langchain.retrievers.document_compressors import CrossEncoderReranker
from langchain_community.cross_encoders import HuggingFaceCrossEncoder

cross_encoder = HuggingFaceCrossEncoder(model_name="cross-encoder/ms-marco-MiniLM-L-6-v2")
reranker = CrossEncoderReranker(model=cross_encoder, top_n=5)

compression_retriever = ContextualCompressionRetriever(
base_compressor=reranker,
base_retriever=retriever
)

return compression_retriever

问题2:检索速度慢

症状:查询响应时间过长

解决方案

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
def optimize_performance():
# 1. 索引优化
optimized_index_params = {
"index_type": "HNSW",
"metric_type": "L2",
"params": {
"M": 32,
"efConstruction": 400,
"ef": 100
}
}

# 2. 缓存策略
cache = LRUCache(maxsize=1000, ttl=300) # 5分钟缓存

# 3. 异步处理
async def async_retrieve(query):
loop = asyncio.get_event_loop()
return await loop.run_in_executor(None, retriever.get_relevant_documents, query)

# 4. 批处理优化
def batch_process_documents(documents, batch_size=100):
for i in range(0, len(documents), batch_size):
batch = documents[i:i+batch_size]
vectorstore.add_documents(batch)

# 5. 硬件优化
# 使用GPU加速嵌入计算
# 升级向量数据库硬件配置

生成相关问题

问题3:答案存在幻觉

症状:模型生成的内容与提供的上下文不符

解决方案

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
def reduce_hallucinations():
# 1. 改进提示工程
anti_hallucination_prompt = ChatPromptTemplate.from_template("""
基于以下上下文信息回答用户问题:
{context}

## 严格遵守以下规则:
1. 只使用提供的上下文信息回答,不要编造任何内容
2. 如果上下文没有相关信息,明确说明"根据提供的文档,无法回答该问题"
3. 对于不确定的信息,使用"可能"、"似乎"等措辞
4. 引用具体的文档来源和页码

用户问题:{question}

回答:
""")

# 2. 事实核查
def fact_check(answer, context):
fact_check_prompt = f"""
检查以下回答是否完全基于提供的上下文:
上下文:{context}
回答:{answer}
回答是否完全基于上下文?如果不是,指出哪些部分不是:
"""

check_result = llm.invoke(fact_check_prompt)
return "完全基于上下文" in check_result.content

# 3. 来源追踪
def add_source_citations(answer, sources):
"""为回答添加来源引用"""
citation_map = build_citation_map(answer, sources)
# 添加引用标记...
return cited_answer

# 4. 置信度评分
def calculate_confidence(answer, context):
"""计算回答的置信度"""
confidence_prompt = f"""
基于上下文计算回答的置信度(0-100):
上下文:{context}
回答:{answer}
置信度:
"""

confidence_score = llm.invoke(confidence_prompt)
return int(confidence_score.content.strip())

问题4:回答格式不规范

症状:生成的回答格式混乱,缺乏结构化

解决方案

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
def improve_formatting():
# 1. 结构化输出提示
structured_prompt = ChatPromptTemplate.from_template("""
基于以下上下文回答用户问题:
{context}

## 输出格式要求:
- 使用Markdown格式
- 对于复杂问题,使用分点列表
- 重要信息加粗显示
- 添加相关文档来源

## 回答结构:
### 回答
[你的详细回答]

### 相关文档
- [文档标题](文档链接) - 页码:X

用户问题:{question}
""")

# 2. 输出解析器
from langchain_core.output_parsers import ResponseSchema, StructuredOutputParser

response_schemas = [
ResponseSchema(name="answer", description="详细回答内容"),
ResponseSchema(name="sources", description="相关文档来源列表"),
ResponseSchema(name="confidence", description="回答置信度(0-100)")
]

output_parser = StructuredOutputParser.from_response_schemas(response_schemas)
format_instructions = output_parser.get_format_instructions()

# 3. 格式验证
def validate_format(answer):
"""验证回答格式是否正确"""
try:
parsed = output_parser.parse(answer)
return True, parsed
except Exception as e:
return False, str(e)

系统相关问题

问题5:内存使用过高

症状:系统内存占用持续增长,可能导致OOM

解决方案

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
def optimize_memory_usage():
# 1. 分批处理
def process_in_batches(items, batch_size=100):
for i in range(0, len(items), batch_size):
batch = items[i:i+batch_size]
yield process_batch(batch)

# 2. 垃圾回收优化
import gc

def optimize_gc():
gc.collect()
gc.set_threshold(700, 10, 10)

# 3. 数据类型优化
def optimize_data_types(documents):
optimized_docs = []
for doc in documents:
# 压缩文本内容
compressed_content = compress_text(doc.page_content)
# 优化元数据
optimized_metadata = {k: str(v) for k, v in doc.metadata.items()}

optimized_docs.append(Document(
page_content=compressed_content,
metadata=optimized_metadata
))
return optimized_docs

# 4. 缓存管理
class MemoryEfficientCache:
def __init__(self, max_size=1000, compression_level=1):
self.cache = {}
self.max_size = max_size
self.compression_level = compression_level

def set(self, key, value):
if len(self.cache) >= self.max_size:
# LRU淘汰
oldest_key = next(iter(self.cache.keys()))
del self.cache[oldest_key]

# 压缩存储
compressed_value = compress_pickle(value, self.compression_level)
self.cache[key] = compressed_value

def get(self, key):
if key in self.cache:
return decompress_pickle(self.cache[key])
return None

问题6:并发处理能力不足

症状:高并发场景下系统响应缓慢或超时

解决方案

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
def improve_concurrency():
# 1. 异步架构
async def async_rag_pipeline(query):
# 并行执行检索和其他操作
retrieval_task = asyncio.create_task(retrieve_documents(query))
query_analysis_task = asyncio.create_task(analyze_query(query))

# 等待所有任务完成
documents, query_analysis = await asyncio.gather(
retrieval_task, query_analysis_task
)

# 生成回答
return await generate_answer(query, documents, query_analysis)

# 2. 连接池管理
class ConnectionPool:
def __init__(self, max_connections=10):
self.pool = asyncio.Queue(maxsize=max_connections)
self._init_pool()

async def _init_pool(self):
for _ in range(self.pool.maxsize):
connection = await create_connection()
await self.pool.put(connection)

async def acquire(self):
return await self.pool.get()

async def release(self, connection):
await self.pool.put(connection)

# 3. 限流控制
from functools import wraps
from asyncio import Semaphore

class RateLimiter:
def __init__(self, max_concurrent=50):
self.semaphore = Semaphore(max_concurrent)

def __call__(self, func):
@wraps(func)
async def wrapper(*args, **kwargs):
async with self.semaphore:
return await func(*args, **kwargs)
return wrapper

# 4. 负载均衡
class LoadBalancer:
def __init__(self, servers):
self.servers = servers
self.current = 0

def get_server(self):
server = self.servers[self.current]
self.current = (self.current + 1) % len(self.servers)
return server

附录

相关资源链接

术语表

术语 英文 描述
检索增强生成 Retrieval-Augmented Generation (RAG) 将检索和生成相结合的AI技术
向量嵌入 Vector Embedding 将文本转换为数值向量的过程
向量数据库 Vector Database 专门存储和检索向量数据的数据库
文本分块 Text Chunking 将长文档分割为小片段的过程
提示工程 Prompt Engineering 设计和优化提示以引导模型输出的技术
语义检索 Semantic Retrieval 基于语义相似性的信息检索方法
混合搜索 Hybrid Search 结合向量搜索和关键词搜索的检索方法

版本更新记录

版本 日期 更新内容
2.0 2025-10-23 增加2025年技术趋势,更新最佳实践案例
1.5 2025-08-15 完善部署和监控章节,增加性能优化策略
1.0 2025-06-01 初始版本发布,包含基础RAG实现指南