从零搭建企业级 RAG 系统:文档切分、Embedding 选型到检索重排的全链路实践
一、RAG 系统定位与业务场景分析
1.1 趣玩搭 RAG 的核心使用场景
基于阶段1解析,趣玩搭的 RAG 系统服务于智能客服,需要解决以下业务场景:
1.2 当前痛点(阶段1已识别)
痛点1:知识分散 → 大模型幻觉严重(准确率仅65%)
痛点2:文档更新滞后 → 用户获取过期信息
痛点3:初创团队成本敏感 → 不能用昂贵的商用方案
痛点4:业务增长快 → 知识库需要快速扩容
1.3 RAG 系统目标指标
二、RAG 全链路架构设计
2.1 系统架构总览
2.2 数据流转全链路时序图
三、文档切分策略(Chunking Strategy)
3.1 趣玩搭文档特征分析
3.2 分层切分策略
针对趣玩搭不同文档类型,采用分层切分而非单一策略:
Layer 1: 文档级预处理
├── Markdown → 按 ## 标题拆分为Section
├── PDF → PyMuPDF提取 → 按段落拆分
└── FAQ → 按 Q&A 对直接作为独立Chunk
Layer 2: Section级精切分(RecursiveCharacterTextSplitter)
├── chunk_size = 512 tokens(趣玩搭文档偏短,512足够)
├── chunk_overlap = 64 tokens(12.5%重叠,保留上下文)
└── separators = ["\n## ", "\n### ", "\n\n", "\n", "。", ";"]
Layer 3: 元数据增强
├── 父文档标题(Parent Title)
├── 文档类别(category: 活动规则/退款/FAQ/...)
├── 更新时间(updated_at)
└── 权限标签(role: user/host/admin)
3.3 切分参数选择依据
3.4 特殊处理:Parent Document Retrieval
对于操作手册类长文档,采用父子文档策略:
父文档(完整Section,用于LLM上下文)
├── 子文档1(切分Chunk,用于向量检索)
├── 子文档2
└── 子文档3
检索流程:
1. 用子文档Chunk做向量检索(精准定位)
2. 命中后取其父文档作为LLM上下文(完整语境)
3.5 核心切分代码
from langchain.text_splitter import RecursiveCharacterTextSplitter
from langchain.schema import Document
import hashlib
from datetime import datetime
from typing import List, Dict
class FunPlayDocumentProcessor:
"""趣玩搭文档处理器 —— 针对业务文档特征定制"""
def __init__(self):
# 主切分器:适配趣玩搭短文本业务文档
self.text_splitter = RecursiveCharacterTextSplitter(
chunk_size=512, # 趣玩搭文档偏短,512 tokens 足够
chunk_overlap=64, # 12.5% 重叠,保留条件语句上下文
separators=[
"\n## ", # Markdown 二级标题(最高优先级)
"\n### ", # Markdown 三级标题
"\n\n", # 段落分隔
"\n", # 行分隔
"。", # 中文句号(适配中文文档)
";", # 中文分号
" ", # 空格
"" # 兜底:字符级分割
],
length_function=len,
is_separator_regex=False,
)
# FAQ 专用切分器(更小粒度)
self.faq_splitter = RecursiveCharacterTextSplitter(
chunk_size=256,
chunk_overlap=32,
separators=["\n\n", "\n", "。"],
)
def process_document(
self,
content: str,
doc_type: str, # "rule" | "faq" | "manual" | "policy"
category: str, # "活动规则" | "退款政策" | "操作手册" | "会员权益"
source: str, # 文件名
role_access: str = "user" # "user" | "host" | "admin"
) -> List[Document]:
"""
文档处理主入口
根据文档类型选择不同的切分策略
"""
if doc_type == "faq":
chunks = self._process_faq(content)
elif doc_type == "manual":
chunks = self._process_manual_with_parent(content)
else:
chunks = self.text_splitter.split_text(content)
# 构建带元数据的 Document 对象
documents = []
for i, chunk in enumerate(chunks):
doc = Document(
page_content=chunk,
metadata={
"doc_id": self._generate_doc_id(source, i),
"source": source,
"category": category,
"doc_type": doc_type,
"role_access": role_access,
"chunk_index": i,
"total_chunks": len(chunks),
"updated_at": datetime.now().isoformat(),
"char_count": len(chunk),
}
)
documents.append(doc)
return documents
def _process_faq(self, content: str) -> List[str]:
"""FAQ 文档:按 Q&A 对拆分"""
pairs = []
current_q, current_a = "", ""
for line in content.split("\n"):
line = line.strip()
if line.startswith("Q:") or line.startswith("Q:") or line.startswith("问:"):
if current_q and current_a:
pairs.append(f"{current_q}\n{current_a}")
current_q = line
current_a = ""
elif line.startswith("A:") or line.startswith("A:") or line.startswith("答:"):
current_a = line
elif current_a:
current_a += "\n" + line
if current_q and current_a:
pairs.append(f"{current_q}\n{current_a}")
# 对过长的 QA 对再做二次切分
result = []
for pair in pairs:
if len(pair) > 256:
result.extend(self.faq_splitter.split_text(pair))
else:
result.append(pair)
return result
def _process_manual_with_parent(self, content: str) -> List[str]:
"""
操作手册:Parent Document Retrieval 策略
先按大Section切分(父文档),再细切(子文档)
子文档用于检索,父文档用于LLM上下文
"""
# 按二级标题切分为 Section(父文档)
sections = content.split("\n## ")
child_chunks = []
for section in sections:
if len(section.strip()) < 10:
continue
# 对每个 Section 做细粒度切分
chunks = self.text_splitter.split_text(section)
for chunk in chunks:
# 子文档中记录父文档的完整内容引用
child_chunks.append(chunk)
return child_chunks
@staticmethod
def _generate_doc_id(source: str, index: int) -> str:
"""生成唯一文档ID"""
raw = f"{source}_{index}_{datetime.now().isoformat()}"
return hashlib.md5(raw.encode()).hexdigest()[:16]
# ========== 使用示例 ==========
if __name__ == "__main__":
processor = FunPlayDocumentProcessor()
# 处理退款政策文档
refund_policy = """
## 退款规则
### 活动开始前退款
1. 活动开始前24小时以上申请退款,全额退还票款。
2. 活动开始前12-24小时申请退款,扣除票款的20%作为手续费。
3. 活动开始前12小时内申请退款,扣除票款的50%作为手续费。
### 活动开始后退款
活动开始后原则上不予退款,特殊情况(如活动取消、场地问题)由平台客服处理。
### 退款到账时间
退款将原路退回,微信支付1-3个工作日到账,钱包余额即时到账。
"""
docs = processor.process_document(
content=refund_policy,
doc_type="policy",
category="退款政策",
source="退款规则v3.md",
role_access="user"
)
for doc in docs:
print(f"[Chunk {doc.metadata['chunk_index']}] "
f"长度={doc.metadata['char_count']} "
f"类别={doc.metadata['category']}")
print(f" 内容预览: {doc.page_content[:80]}...")
print()
四、Embedding 模型选型
4.1 选型对比(针对趣玩搭场景:中文为主、成本敏感、本地部署优先)
4.2 选型决策:维持 bge-m3(✅ 最优选择)
理由:
中文能力最强:MTEB中文基准得分71.4,碾压竞品,趣玩搭文档全部是中文
原生支持混合检索:bge-m3同时输出Dense + Sparse向量,天然适配Hybrid Search,无需额外BM25索引
本地部署零成本:初创公司成本敏感,本地部署不产生API调用费用
数据隐私:用户数据(问题/活动信息)不外传,符合社交平台数据合规要求
多粒度能力:支持不同层级的语义匹配,适合趣玩搭多类型文档场景
部署建议:
# bge-m3 部署配置(趣玩搭推荐)
deployment:
model: BAAI/bge-m3
runtime: sentence-transformers # 或用 FlagEmbedding 原生库
device: cuda:0 # 有GPU用GPU,无GPU可用CPU(速度约慢3x)
batch_size: 32 # 批量处理文档时的batch大小
max_length: 512 # 与chunk_size对齐
normalize: true # L2归一化,便于余弦相似度计算
# 资源需求
resources:
gpu_memory: ~4GB (float16) # 单张4GB显卡即可
cpu_fallback: ~6GB RAM # 无GPU时CPU模式
qps_capacity: ~200/s (GPU) # 远超趣玩搭需求
4.3 Embedding 服务代码
from FlagEmbedding import BGEM3FlagModel
from typing import List, Dict
import numpy as np
class FunPlayEmbeddingService:
"""趣玩搭 Embedding 服务 —— 基于 bge-m3"""
def __init__(self, model_path: str = "BAAI/bge-m3", use_gpu: bool = True):
self.model = BGEM3FlagModel(
model_path,
use_fp16=use_gpu, # GPU 用 FP16 加速
device="cuda:0" if use_gpu else "cpu"
)
self.max_length = 512
def encode_documents(self, texts: List[str]) -> Dict[str, np.ndarray]:
"""
文档向量化(同时输出 Dense + Sparse 向量)
用于写入 Qdrant 时的双路索引
"""
output = self.model.encode(
texts,
batch_size=32,
max_length=self.max_length,
return_dense=True,
return_sparse=True, # bge-m3 特有:稀疏向量
return_colbert_vecs=False, # ColBERT 向量暂不需要
)
return {
"dense": output["dense_vecs"], # shape: (N, 1024)
"sparse": output["lexical_weights"] # 稀疏权重字典
}
def encode_query(self, query: str) -> Dict[str, np.ndarray]:
"""
Query 向量化
单条处理,低延迟优先
"""
output = self.model.encode(
[query],
batch_size=1,
max_length=self.max_length,
return_dense=True,
return_sparse=True,
return_colbert_vecs=False,
)
return {
"dense": output["dense_vecs"][0], # shape: (1024,)
"sparse": output["lexical_weights"][0] # 稀疏权重
}
# ========== 使用示例 ==========
if __name__ == "__main__":
service = FunPlayEmbeddingService(use_gpu=True)
# 文档向量化
docs = [
"活动开始前24小时以上申请退款,全额退还票款。",
"候补排队成功后,系统会自动发送微信通知。",
"金牌会员享受活动票价9折优惠。"
]
doc_vectors = service.encode_documents(docs)
print(f"Dense 向量维度: {doc_vectors['dense'].shape}") # (3, 1024)
# Query 向量化
query_vec = service.encode_query("退票怎么退?")
print(f"Query 向量维度: {query_vec['dense'].shape}") # (1024,)
五、Vector Database 选型
5.1 选型对比(针对趣玩搭:小规模、成本敏感、已选Qdrant)
5.2 选型决策:维持 Qdrant(✅ 最优选择)
核心理由(匹配趣玩搭实际):
不选其他方案的原因:
Milvus:功能强大但重(依赖 etcd + MinIO),趣玩搭规模不需要
Chroma:适合原型开发,生产环境稳定性不足,无混合检索
Pinecone:SaaS 托管,数据出境 + 月费贵,初创不适合
Weaviate:不原生支持稀疏向量,与 bge-m3 配合不如 Qdrant
5.3 Qdrant 部署与集合配置
# docker-compose.yml —— Qdrant 部署(趣玩搭生产环境)
version: '3.8'
services:
qdrant:
image: qdrant/qdrant:v1.9.0
container_name: funplay-qdrant
ports:
- "6333:6333" # REST API
- "6334:6334" # gRPC(Spring AI 用此端口)
volumes:
- ./qdrant_data:/qdrant/storage # 数据持久化
- ./qdrant_snapshots:/qdrant/snapshots # 快照备份
environment:
- QDRANT__STORAGE__STORAGE_PATH=/qdrant/storage
- QDRANT__SERVICE__GRPC_PORT=6334
- QDRANT__STORAGE__ON_DISK_PAYLOAD=true # Payload 存磁盘,省内存
deploy:
resources:
limits:
memory: 4G
cpus: "2.0"
restart: always
from qdrant_client import QdrantClient
from qdrant_client.models import (
VectorParams, SparseVectorParams, Distance,
PointStruct, SparseVector, NamedVector,
Filter, FieldCondition, MatchValue,
CreateCollection, OptimizersConfigDiff,
)
class FunPlayVectorStore:
"""趣玩搭向量存储管理 —— Qdrant"""
COLLECTION_NAME = "funplay_knowledge_base"
def __init__(self, host: str = "localhost", port: int = 6333):
self.client = QdrantClient(host=host, port=port)
def create_collection(self):
"""
创建集合 —— 支持 Dense + Sparse 双路向量
配合 bge-m3 的混合检索能力
"""
self.client.create_collection(
collection_name=self.COLLECTION_NAME,
vectors_config={
"dense": VectorParams(
size=1024, # bge-m3 Dense 向量维度
distance=Distance.COSINE,
on_disk=False, # 数据量小,全内存加速
)
},
sparse_vectors_config={
"sparse": SparseVectorParams(
index={"on_disk": False} # 稀疏索引也放内存
)
},
optimizers_config=OptimizersConfigDiff(
indexing_threshold=1000, # 超过1000条才建HNSW索引
memmap_threshold=50000, # 超过5万条才用mmap
),
)
# 创建 Payload 索引(加速元数据过滤)
self.client.create_payload_index(
collection_name=self.COLLECTION_NAME,
field_name="category",
field_schema="keyword", # 类别过滤
)
self.client.create_payload_index(
collection_name=self.COLLECTION_NAME,
field_name="role_access",
field_schema="keyword", # 权限过滤
)
self.client.create_payload_index(
collection_name=self.COLLECTION_NAME,
field_name="doc_type",
field_schema="keyword", # 文档类型过滤
)
print(f"Collection '{self.COLLECTION_NAME}' created with hybrid vector support.")
def upsert_documents(self, documents: list, dense_vectors, sparse_vectors):
"""
批量写入文档向量
同时写入 Dense 和 Sparse 双路向量 + 元数据
"""
points = []
for i, doc in enumerate(documents):
# 构建稀疏向量
sparse_indices = list(sparse_vectors[i].keys())
sparse_values = list(sparse_vectors[i].values())
point = PointStruct(
id=hash(doc.metadata["doc_id"]) & 0x7FFFFFFFFFFFFFFF, # 正整数ID
vector={
"dense": dense_vectors[i].tolist(),
"sparse": SparseVector(
indices=sparse_indices,
values=sparse_values,
)
},
payload={
"content": doc.page_content, # 原文内容
"doc_id": doc.metadata["doc_id"],
"source": doc.metadata["source"],
"category": doc.metadata["category"],
"doc_type": doc.metadata["doc_type"],
"role_access": doc.metadata["role_access"],
"chunk_index": doc.metadata["chunk_index"],
"updated_at": doc.metadata["updated_at"],
}
)
points.append(point)
self.client.upsert(
collection_name=self.COLLECTION_NAME,
points=points,
)
print(f"Upserted {len(points)} documents.")
def hybrid_search(
self,
dense_vector,
sparse_vector,
category: str = None,
role_access: str = "user",
top_k: int = 20,
) -> list:
"""
混合检索:Dense语义检索 + Sparse关键词检索
使用 Qdrant 原生 Query API
"""
# 构建过滤条件
must_conditions = [
FieldCondition(key="role_access", match=MatchValue(value=role_access))
]
if category:
must_conditions.append(
FieldCondition(key="category", match=MatchValue(value=category))
)
query_filter = Filter(must=must_conditions)
# Dense 检索
dense_results = self.client.search(
collection_name=self.COLLECTION_NAME,
query_vector=("dense", dense_vector.tolist()),
query_filter=query_filter,
limit=top_k,
with_payload=True,
)
# Sparse 检索
sparse_indices = list(sparse_vector.keys())
sparse_values = list(sparse_vector.values())
sparse_results = self.client.search(
collection_name=self.COLLECTION_NAME,
query_vector=NamedVector(
name="sparse",
vector=SparseVector(indices=sparse_indices, values=sparse_values)
),
query_filter=query_filter,
limit=top_k,
with_payload=True,
)
return dense_results, sparse_results
六、检索策略:Hybrid Search + Rerank
6.1 混合检索流程
graph LR
A[用户Query] --> B[Query改写]
B --> C[bge-m3 Encode]
C --> D1[Dense Vector<br/>1024d]
C --> D2[Sparse Vector<br/>词级权重]
D1 --> E1[Qdrant Dense Search<br/>HNSW ANN<br/>Top-20]
D2 --> E2[Qdrant Sparse Search<br/>倒排索引<br/>Top-20]
E1 --> F[RRF 融合排序<br/>Reciprocal Rank Fusion]
E2 --> F
F --> G[候选集 Top-20]
G --> H[bge-reranker-v2-m3<br/>交叉编码重排序]
H --> I[最终结果 Top-5]
style E1 fill:#e3f2fd
style E2 fill:#fff3e0
style F fill:#f3e5f5
style H fill:#fce4ec
6.2 为什么趣玩搭需要混合检索?
6.3 RRF 融合排序算法
from typing import List, Dict, Tuple
from collections import defaultdict
def reciprocal_rank_fusion(
dense_results: list,
sparse_results: list,
k: int = 60, # RRF常数,标准值60
dense_weight: float = 0.6, # 语义检索权重(趣玩搭语义查询占比高)
sparse_weight: float = 0.4, # 关键词检索权重
) -> List[Tuple[str, float, dict]]:
"""
RRF(Reciprocal Rank Fusion)融合排序
趣玩搭权重配置理由:
- dense_weight=0.6:用户提问以自然语言为主("怎么退票""能不能退")
- sparse_weight=0.4:部分查询包含精确术语("金牌会员""候补")
"""
scores = defaultdict(float)
payloads = {}
# Dense 结果计算 RRF 分数
for rank, result in enumerate(dense_results):
doc_id = result.payload["doc_id"]
scores[doc_id] += dense_weight * (1.0 / (k + rank + 1))
payloads[doc_id] = result.payload
# Sparse 结果计算 RRF 分数
for rank, result in enumerate(sparse_results):
doc_id = result.payload["doc_id"]
scores[doc_id] += sparse_weight * (1.0 / (k + rank + 1))
if doc_id not in payloads:
payloads[doc_id] = result.payload
# 按融合分数降序排列
ranked = sorted(scores.items(), key=lambda x: x[1], reverse=True)
results = [
(doc_id, score, payloads[doc_id])
for doc_id, score in ranked
]
return results
6.4 Query 改写(适配趣玩搭口语化查询)
class FunPlayQueryRewriter:
"""
Query 改写器
趣玩搭用户画像:18-35岁年轻人,提问口语化、简短
典型查询:"咋退票""人不够咋办""vip有啥用"
需要将口语化查询改写为知识库可匹配的标准表述
"""
# 趣玩搭业务同义词表
SYNONYM_MAP = {
"退票": ["退款", "取消订单", "退费", "不想去了"],
"候补": ["排队", "等位", "候补票", "等票"],
"会员": ["VIP", "vip", "金牌", "银牌", "等级"],
"主理人": ["组织者", "发起人", "俱乐部"],
"活动": ["局", "组局", "拼局", "场次"],
"票": ["门票", "入场券", "名额"],
}
def rewrite(self, query: str) -> str:
"""
基础改写:同义词扩展
将口语表述映射到知识库术语
"""
expanded = query
for standard_term, synonyms in self.SYNONYM_MAP.items():
for syn in synonyms:
if syn in query and standard_term not in query:
expanded += f" {standard_term}"
break
return expanded
def generate_hypothetical_answer(self, query: str) -> str:
"""
HyDE(Hypothetical Document Embedding)策略
让LLM生成一个假设性答案,用这个答案去检索
适用于模糊查询场景
注意:此方法会增加一次LLM调用延迟,仅在复杂查询时启用
"""
# 实际实现中调用 DeepSeek/Qwen 生成假设答案
# 这里展示 Prompt 模板
prompt = f"""你是趣玩搭平台的客服助手。
请针对以下用户问题,生成一个简短的假设性答案(50字以内),
这个答案应该包含可能出现在知识库中的关键术语。
用户问题:{query}
假设性答案:"""
return prompt # 实际调用 LLM 获取答案
七、重排序(Reranker)
7.1 Reranker 选型对比
7.2 选型决策:bge-reranker-v2-m3
理由:
与 bge-m3 同系列,配合效果最优(BAAI 团队联合优化)
本地部署零成本,与趣玩搭成本敏感策略一致
中文重排序能力领先
7.3 Reranker 服务代码
from FlagEmbedding import FlagReranker
from typing import List, Tuple
class FunPlayReranker:
"""趣玩搭重排序服务 —— bge-reranker-v2-m3"""
def __init__(self, model_path: str = "BAAI/bge-reranker-v2-m3", use_gpu: bool = True):
self.reranker = FlagReranker(
model_path,
use_fp16=use_gpu,
device="cuda:0" if use_gpu else "cpu"
)
def rerank(
self,
query: str,
candidates: List[Tuple[str, float, dict]], # (doc_id, rrf_score, payload)
top_k: int = 5,
score_threshold: float = 0.3, # 低于此分数的结果丢弃
) -> List[dict]:
"""
重排序:对 RRF 融合后的候选集做精细排序
输入:RRF 融合后的 Top-20 候选
输出:精排后的 Top-5 结果
"""
if not candidates:
return []
# 构建 query-passage 对
pairs = [
[query, candidate[2]["content"]]
for candidate in candidates
]
# 计算重排序分数
scores = self.reranker.compute_score(
pairs,
normalize=True, # 归一化到 [0, 1]
)
# 如果只有一个候选,scores 不是列表
if isinstance(scores, float):
scores = [scores]
# 组装结果
results = []
for i, (doc_id, rrf_score, payload) in enumerate(candidates):
rerank_score = scores[i]
if rerank_score >= score_threshold:
results.append({
"doc_id": doc_id,
"content": payload["content"],
"category": payload["category"],
"source": payload["source"],
"rrf_score": rrf_score,
"rerank_score": rerank_score,
"final_score": 0.3 * rrf_score + 0.7 * rerank_score, # 加权融合
})
# 按 final_score 降序排列,取 Top-K
results.sort(key=lambda x: x["final_score"], reverse=True)
return results[:top_k]
# ========== 使用示例 ==========
if __name__ == "__main__":
reranker = FunPlayReranker(use_gpu=True)
query = "活动取消了怎么退款"
candidates = [
("doc_001", 0.85, {
"content": "活动开始前24小时以上申请退款,全额退还票款。",
"category": "退款政策", "source": "退款规则.md"
}),
("doc_002", 0.72, {
"content": "活动取消时,系统自动发起全额退款,原路退回。",
"category": "退款政策", "source": "退款规则.md"
}),
("doc_003", 0.65, {
"content": "金牌会员享受活动票价9折优惠。",
"category": "会员权益", "source": "会员手册.md"
}),
]
results = reranker.rerank(query, candidates, top_k=2)
for r in results:
print(f"[{r['category']}] score={r['final_score']:.3f}: {r['content'][:50]}")
八、召回评估体系
8.1 评估指标定义
8.2 评估数据集构建
"""
趣玩搭 RAG 评估数据集
基于真实业务场景构建,覆盖7大查询类别
"""
FUNPLAY_EVAL_DATASET = [
# === 退款类(高频) ===
{
"query": "活动明天就开始了,现在还能退票吗",
"relevant_doc_ids": ["refund_001", "refund_002"],
"category": "退款政策",
"difficulty": "easy",
},
{
"query": "退款多久能到微信零钱",
"relevant_doc_ids": ["refund_005"],
"category": "退款政策",
"difficulty": "easy",
},
{
"query": "活动取消了钱自动退吗",
"relevant_doc_ids": ["refund_003", "refund_004"],
"category": "退款政策",
"difficulty": "medium",
},
# === 票务类(高频) ===
{
"query": "候补排队大概要等多久",
"relevant_doc_ids": ["ticket_003"],
"category": "票务操作",
"difficulty": "medium",
},
{
"query": "票卖完了还有办法参加吗",
"relevant_doc_ids": ["ticket_003", "ticket_004"],
"category": "票务操作",
"difficulty": "medium",
},
# === 会员类(中频) ===
{
"query": "怎么升到金牌会员",
"relevant_doc_ids": ["member_001", "member_002"],
"category": "会员权益",
"difficulty": "easy",
},
{
"query": "vip打几折",
"relevant_doc_ids": ["member_003"],
"category": "会员权益",
"difficulty": "easy",
},
# === 主理人类(低频) ===
{
"query": "怎么设置早鸟票和普通票两种价格",
"relevant_doc_ids": ["host_005", "host_006"],
"category": "主理人操作",
"difficulty": "hard",
},
# === 口语化/模糊查询(测试鲁棒性) ===
{
"query": "人不够咋整",
"relevant_doc_ids": ["rule_002"],
"category": "活动规则",
"difficulty": "hard",
},
{
"query": "钱还没退过来",
"relevant_doc_ids": ["refund_005", "refund_006"],
"category": "退款政策",
"difficulty": "hard",
},
]
8.3 自动化评估代码
import numpy as np
from typing import List, Dict
class FunPlayRAGEvaluator:
"""趣玩搭 RAG 检索质量评估器"""
def __init__(self, retrieval_pipeline):
"""retrieval_pipeline: 完整的检索链路(Query→Embedding→Search→Rerank)"""
self.pipeline = retrieval_pipeline
def evaluate(self, eval_dataset: List[Dict]) -> Dict:
"""运行完整评估"""
metrics = {
"recall_at_5": [], "precision_at_5": [],
"mrr": [], "ndcg_at_5": [], "hit_rate": [],
}
for sample in eval_dataset:
query = sample["query"]
relevant_ids = set(sample["relevant_doc_ids"])
# 执行检索
results = self.pipeline.search(query, top_k=5)
retrieved_ids = [r["doc_id"] for r in results]
# 计算各指标
metrics["recall_at_5"].append(
self._recall_at_k(relevant_ids, retrieved_ids, k=5)
)
metrics["precision_at_5"].append(
self._precision_at_k(relevant_ids, retrieved_ids, k=5)
)
metrics["mrr"].append(
self._mrr(relevant_ids, retrieved_ids)
)
metrics["ndcg_at_5"].append(
self._ndcg_at_k(relevant_ids, retrieved_ids, k=5)
)
metrics["hit_rate"].append(
1.0 if len(relevant_ids & set(retrieved_ids)) > 0 else 0.0
)
# 汇总
summary = {k: np.mean(v) for k, v in metrics.items()}
summary["total_samples"] = len(eval_dataset)
return summary
@staticmethod
def _recall_at_k(relevant: set, retrieved: list, k: int) -> float:
retrieved_set = set(retrieved[:k])
if not relevant:
return 0.0
return len(relevant & retrieved_set) / len(relevant)
@staticmethod
def _precision_at_k(relevant: set, retrieved: list, k: int) -> float:
retrieved_set = set(retrieved[:k])
if not retrieved_set:
return 0.0
return len(relevant & retrieved_set) / len(retrieved_set)
@staticmethod
def _mrr(relevant: set, retrieved: list) -> float:
for i, doc_id in enumerate(retrieved):
if doc_id in relevant:
return 1.0 / (i + 1)
return 0.0
@staticmethod
def _ndcg_at_k(relevant: set, retrieved: list, k: int) -> float:
dcg = 0.0
for i, doc_id in enumerate(retrieved[:k]):
if doc_id in relevant:
dcg += 1.0 / np.log2(i + 2)
ideal_dcg = sum(1.0 / np.log2(i + 2) for i in range(min(len(relevant), k)))
return dcg / ideal_dcg if ideal_dcg > 0 else 0.0
# ========== 执行评估 ==========
# evaluator = FunPlayRAGEvaluator(pipeline)
# results = evaluator.evaluate(FUNPLAY_EVAL_DATASET)
#
# 预期输出:
# {
# "recall_at_5": 0.92, ← 目标 ≥ 0.90 ✅
# "precision_at_5": 0.65, ← 目标 ≥ 0.60 ✅
# "mrr": 0.83, ← 目标 ≥ 0.80 ✅
# "ndcg_at_5": 0.87, ← 目标 ≥ 0.85 ✅
# "hit_rate": 0.96, ← 目标 ≥ 0.95 ✅
# "total_samples": 10
# }
九、完整检索 Pipeline 整合
9.1 端到端 Pipeline(Spring AI + Python 混合架构)
趣玩搭 RAG Pipeline 架构:
┌─────────────────────────────────────────────────────────┐
│ Spring Boot 应用层 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ Spring AI ChatClient │ │
│ │ ├── QuestionAnswerAdvisor (RAG Advisor) │ │
│ │ ├── MessageChatMemoryAdvisor (上下文记忆) │ │
│ │ └── Prompt Template (业务模板) │ │
│ └─────────────────────────────────────────────────────┘ │
│ │ 调用 │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ VectorStore (Spring AI Qdrant Store) │ │
│ │ └── QdrantVectorStore.similaritySearch() │ │
│ └─────────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────┘
│ gRPC
┌─────────────────────────────────────────────────────────┐
│ Python Sidecar 服务(FastAPI) │
│ ├── /embed → bge-m3 向量化 │
│ ├── /rerank → bge-reranker-v2-m3 重排序 │
│ └── /health → 健康检查 │
└─────────────────────────────────────────────────────────┘
9.2 Spring AI 端核心配置
/**
* 趣玩搭 RAG 配置类
* 基于 Spring AI 框架集成 Qdrant + DeepSeek/Qwen
*/
@Configuration
public class FunPlayRAGConfig {
/**
* Qdrant 向量存储 Bean
*/
@Bean
public QdrantVectorStore vectorStore(EmbeddingModel embeddingModel) {
return QdrantVectorStore.builder(
new QdrantClient(
QdrantGrpcClient.newBuilder("localhost", 6334, false).build()
),
embeddingModel
)
.collectionName("funplay_knowledge_base")
.build();
}
/**
* DeepSeek 主力模型
*/
@Bean
@Primary
public ChatModel deepSeekChatModel() {
return OpenAiChatModel.builder()
.openAiApi(OpenAiApi.builder()
.baseUrl("https://api.deepseek.com")
.apiKey("${DEEPSEEK_API_KEY}")
.build())
.defaultOptions(OpenAiChatOptions.builder()
.model("deepseek-chat")
.temperature(0.1) // 客服场景低温度,减少幻觉
.maxTokens(1024)
.build())
.build();
}
/**
* Qwen2-7B 本地兜底模型
*/
@Bean("localChatModel")
public ChatModel qwenLocalChatModel() {
return OllamaChatModel.builder()
.ollamaApi(new OllamaApi("http://localhost:11434"))
.defaultOptions(OllamaChatOptions.builder()
.model("qwen2:7b")
.temperature(0.1)
.build())
.build();
}
/**
* RAG Advisor(核心:检索增强生成顾问)
*/
@Bean
public Advisor ragAdvisor(VectorStore vectorStore) {
return QuestionAnswerAdvisor.builder(vectorStore)
.searchRequest(SearchRequest.builder()
.topK(5)
.similarityThreshold(0.7) // 相似度阈值
.build())
.build();
}
}
9.3 RAG 客服服务核心代码
/**
* 趣玩搭智能客服 RAG 服务
* 核心流程:Query → 缓存检查 → Embedding → 混合检索 → Rerank → LLM生成 → 缓存写入
*/
@Service
@Slf4j
public class FunPlayRAGService {
private final ChatModel chatModel;
private final ChatModel localChatModel;
private final VectorStore vectorStore;
private final Advisor ragAdvisor;
private final StringRedisTemplate redisTemplate;
// 趣玩搭业务 Prompt 模板
private static final String SYSTEM_PROMPT = """
你是趣玩搭平台的智能客服助手,名字叫"搭搭"。
## 你的角色
- 专业、友好、耐心的客服助手
- 只回答与趣玩搭平台相关的问题
- 如果知识库中没有相关信息,诚实告知用户并建议转人工
## 回答规则
1. 严格基于以下知识库内容回答,不要编造信息
2. 回答简洁明了,控制在200字以内
3. 涉及金额/时间/规则时必须准确引用知识库内容
4. 如果用户问题超出知识库范围,回复:"抱歉,这个问题我暂时无法回答,
建议您联系人工客服获取帮助哦~"
5. 语气亲和,适当使用 emoji
## 知识库内容
{context}
""";
private static final String USER_PROMPT = """
用户问题:{question}
请基于知识库内容回答。如果知识库中没有相关信息,请诚实告知。
""";
/**
* 处理用户查询(端到端)
*/
public String handleQuery(String userId, String question, String sessionId) {
// 1. 检查 Redis 缓存
String cacheKey = "rag:cache:" + DigestUtils.md5Hex(question);
String cached = redisTemplate.opsForValue().get(cacheKey);
if (cached != null) {
log.info("Cache hit for query: {}", question);
return cached;
}
// 2. 调用 RAG 链路
String answer;
try {
answer = callRAGPipeline(question, sessionId);
} catch (Exception e) {
log.error("Primary model failed, falling back to local model", e);
answer = callLocalModel(question);
}
// 3. 写入缓存(TTL 30分钟)
redisTemplate.opsForValue().set(cacheKey, answer, Duration.ofMinutes(30));
// 4. 异步记录日志(用于评估优化)
asyncLogService.logQuery(userId, question, answer, sessionId);
return answer;
}
private String callRAGPipeline(String question, String sessionId) {
ChatResponse response = ChatClient.create(chatModel)
.prompt()
.system(SYSTEM_PROMPT)
.user(USER_PROMPT.replace("{question}", question))
.advisors(
ragAdvisor,
MessageChatMemoryAdvisor.builder(
new RedisChatMemory(redisTemplate, sessionId)
)
.chatMemoryRetrieveSize(10) // 保留最近10轮对话
.build()
)
.call()
.chatResponse();
return response.getResult().getOutput().getText();
}
private String callLocalModel(String question) {
// 降级到本地 Qwen2-7B
return ChatClient.create(localChatModel)
.prompt()
.system(SYSTEM_PROMPT)
.user(USER_PROMPT.replace("{question}", question))
.advisors(ragAdvisor)
.call()
.chatResponse()
.getResult().getOutput().getText();
}
}
十、文档实时更新机制(5分钟内同步)
10.1 更新流程
sequenceDiagram
participant Admin as 管理员/运营
participant API as 文档管理API
participant MQ as RocketMQ
participant Worker as 向量化Worker
participant EMB as bge-m3
participant QD as Qdrant
participant RC as Redis Cache
Admin->>API: 上传/更新文档
API->>API: 文档预处理 + 切分
API->>MQ: 发送消息(doc_id, chunks, metadata)
MQ->>Worker: 消费消息
Worker->>EMB: 批量向量化
EMB-->>Worker: 向量结果
Worker->>QD: 按 doc_id 删除旧向量
Worker->>QD: 批量写入新向量
Worker->>RC: 清除相关缓存(模糊匹配)
Worker->>API: 回调通知:同步完成
Note over Admin, RC: 全流程 < 5分钟(实测通常 < 1分钟)
10.2 增量更新策略
class FunPlayKnowledgeUpdater:
"""知识库增量更新器"""
def update_document(self, source: str, new_content: str, metadata: dict):
"""
增量更新策略:
1. 按 source(文件名)删除旧的所有 Chunk
2. 重新切分并写入新 Chunk
3. 清除相关 Redis 缓存
"""
# Step 1: 删除旧向量(按 source 过滤)
self.qdrant_client.delete(
collection_name="funplay_knowledge_base",
points_selector=FilterSelector(
filter=Filter(
must=[FieldCondition(key="source", match=MatchValue(value=source))]
)
)
)
# Step 2: 重新处理并写入
processor = FunPlayDocumentProcessor()
docs = processor.process_document(new_content, **metadata)
embedding_service = FunPlayEmbeddingService()
vectors = embedding_service.encode_documents([d.page_content for d in docs])
vector_store = FunPlayVectorStore()
vector_store.upsert_documents(docs, vectors["dense"], vectors["sparse"])
# Step 3: 清除缓存(模糊匹配该类别的所有缓存)
category = metadata.get("category", "")
cache_pattern = f"rag:cache:*"
# 生产环境建议使用 SCAN 而非 KEYS
keys = self.redis_client.scan_iter(match=cache_pattern, count=100)
for key in keys:
self.redis_client.delete(key)
return {"status": "success", "chunks_updated": len(docs)}
十一、成本估算与扩容方案
11.1 当前阶段成本估算(DAU 1000+,知识库 <10K 文档)
11.2 扩容路线图
graph LR
subgraph Phase1["阶段1: 当前(DAU 1K)"]
A1[单机部署<br/>bge-m3 + Qdrant<br/>+ Qwen2-7B<br/>共用1台GPU服务器]
end
subgraph Phase2["阶段2: 增长期(DAU 5K-10K)"]
B1[Embedding服务独立<br/>2C4G专用容器]
B2[Qdrant 扩容<br/>4C8G + SSD]
B3[Redis缓存集群<br/>3节点哨兵]
end
subgraph Phase3["阶段3: 规模化(DAU 50K+)"]
C1[Embedding GPU集群<br/>2+ GPU卡]
C2[Qdrant分布式集群<br/>3节点分片]
C3[LLM推理独立集群<br/>vLLM部署]
end
Phase1 -->|DAU > 5K| Phase2
Phase2 -->|DAU > 50K| Phase3
十二、风险与优化建议
12.1 已识别风险
12.2 后续优化方向
阶段2技术方案总结
本方案为趣玩搭量身定制了一套高性价比、可落地的企业级 RAG 系统: