RAG与Agent层:召回不准比没有更危险

摘要
RAG上线后召回率虚高、引用张冠李戴、Agent循环调用空耗算力。本文从某企业知识库项目真实复盘切入,剖析召回、重排、规划、记忆、知识库五个痛点,给出语义切分、cross-encoder重排、迭代终止、双层记忆、增量流水的量化方案。
1. 召回不准的根因是chunk切分破坏语义完整性
痛点现场
某企业知识库RAG上线,用户问"信用卡年费减免条件",系统召回了一段"信用卡权益介绍"的chunk,答案张冠李戴——权益介绍里确实提到"年费"二字,但没说减免条件。排查发现chunk按固定500 token切分,把"年费介绍"和"年费减免条件"分到了两个chunk,召回命中了前者,后者没被召回。
固定长度切分破坏语义边界的杀伤力在于召回率高但精度低——总能召回到含关键词的chunk,但chunk内容可能不是用户要的。某项目离线评测召回率85%(看起来不错),用户实际满意度只有60%——召回的chunk沾边但不切题,答案似是而非。
更隐蔽的是chunk内语义断裂,固定500 token切到句子中间,chunk前半句是上一段结尾,后半句是下一段开头,两个不相关内容拼成一个chunk,embedding编码出来的向量谁都不像,召回相关性骤降。
根因剖析
固定切分失效的机理是语义单位与长度单位不对齐。文档的语义单位是段落/章节,长度可能200也可能800 token,固定500切分必然切到语义边界中间,要么截断段落要么拼接无关段落。
embedding对语义完整性敏感,同一概念的完整段落编码出的向量集中,被截断的半句编码出的向量漂移。检索时query"年费减免条件"的向量与完整段落的向量近,与被截断半句的向量远,召回不到正确chunk。
工程方案:语义边界切分+overlap冗余
方案分两步。切分:按句切分后计算句间语义相似度,相似度高于阈值的归入同一chunk,低于阈值开启新chunk,保证chunk内语义连贯。冗余:chunk之间加10-20%overlap,避免边界信息丢失——某chunk结尾与下一chunk开头重叠,边界概念两个chunk都能召回到。
// 来源:langchain 0.1.0 / langchain/text_splitter + 自研语义切分
import numpy as np
from sentence_transformers import SentenceTransformer
class SemanticBoundarySplitter:
def __init__(self, embed_model_name='m3e-base',
similarity_threshold=0.75,
max_chunk_tokens=400,
overlap_ratio=0.15):
"""
similarity_threshold: 句间相似度低于此值开启新chunk
max_chunk_tokens: chunk最大长度,防止过长稀释相关性
overlap_ratio: chunk间重叠比例,保边界信息
"""
self.embedder = SentenceTransformer(embed_model_name)
self.sim_threshold = similarity_threshold
self.max_chunk = max_chunk_tokens
self.overlap_ratio = overlap_ratio
def split(self, document):
"""语义边界切分主流程"""
# 1. 按句切分,中文按。!?,英文按. ! ?
sentences = self._split_sentences(document)
if not sentences:
return []
# 2. 编码所有句子的embedding
embeddings = self.embedder.encode(sentences,
normalize_embeddings=True)
# 3. 语义聚类形成chunk
chunks = []
current_chunk = [sentences[0]]
current_embedding = embeddings[0]
for i in range(1, len(sentences)):
# 计算当前句与chunk首句的相似度
sim = np.dot(embeddings[i], current_embedding)
# 检查chunk长度是否超限
current_tokens = sum(len(s) for s in current_chunk)
sentence_tokens = len(sentences[i])
if sim > self.sim_threshold and \
current_tokens + sentence_tokens <= self.max_chunk:
# 语义相近且未超长,归入当前chunk
current_chunk.append(sentences[i])
else:
# 语义断或超长,关闭当前chunk开启新chunk
chunks.append(' '.join(current_chunk))
# overlap冗余:新chunk带上一chunk末尾几句
overlap_count = max(1, int(len(current_chunk) *
self.overlap_ratio))
current_chunk = current_chunk[-overlap_count:] + \
[sentences[i]]
current_embedding = embeddings[i]
# 收尾
if current_chunk:
chunks.append(' '.join(current_chunk))
return chunks
def _split_sentences(self, text):
"""按标点切分句子,保留标点"""
import re
# 中英文标点切分
pattern = r'(?<=[。!?.!?])\s*'
sentences = re.split(pattern, text.strip())
return [s.strip() for s in sentences if s.strip()]
# 量化对比:固定500切分召回精度72%,语义切分召回精度89%
# overlap冗余让边界chunk召回率从65%提到83%
量化指标与边界
某项目落地语义切分后,召回精度从72%提到89%,用户满意度从60%提到82%。sim_threshold=0.75是经验值,过低chunk过长稀释相关性(一个chunk讲多件事),过高chunk过碎丢失上下文(一个概念被拆到多chunk)。overlap_ratio=0.15平衡边界信息与冗余,过长冗余导致重复召回,过短边界丢失。
边界与踩坑:语义切分依赖embedding模型质量,m3e-base对中文效果好但对专业术语仍弱,金融/医疗建议用领域微调的embedding。chunk长度max_chunk_tokens=400是经验,长文档场景可到600,对话场景保持300。语义切分的计算开销比固定切分高一个量级(每句要编码),离线建库时无所谓,实时切分需缓存embedding。overlap导致存储膨胀15%和召回重复,需配合去重逻辑。
2. 重排缺失让召回排序靠概率蒙
痛点现场
某项目用向量召回取top50,直接喂给LLM生成答案,用户反馈答案经常"沾边但不切题"。排查发现top50里正确chunk排在第30位,LLM上下文窗口只放得下top5,正确chunk被截掉了。向量召回的相关性排序不稳,正确chunk可能排第1也可能排第30,靠概率蒙。
向量召回排序不稳的机理是embedding只做query与chunk的独立编码,query和chunk分别编码成向量再算相似度,没有两者的交叉注意力。cross-encoder把query和chunk拼一起做联合编码,能捕捉细粒度的语义交互,排序精度高一个量级。
但团队没上重排,原因有二:一是不知道重排的价值,以为向量召回够用;二是担心重排的计算开销,每对query-chunk要联合前向,50候选约30ms,觉得贵。实际上30ms换召回精度从72%到91%是划算的trade-off。
工程方案:cross-encoder重排+粗精两段
方案是粗精两段。粗排用向量召回取top50,保证召回率(正确chunk在top50内)。精排用cross-encoder对50候选重排取top5,保证精度(正确chunk排到top5)。LLM只看top5,上下文窗口够用且答案切题。
// 来源:sentence-transformers 2.2.0 / CrossEncoder + 自研重排
import torch
import numpy as np
from sentence_transformers import CrossEncoder
class CrossEncoderReranker:
def __init__(self, model_name='bge-reranker-base',
top_k_after_rerank=5,
batch_size=32):
"""
cross-encoder重排器
top_k_after_rerank: 重排后取前k个喂LLM
batch_size: 联合前向的batch,平衡延迟与吞吐
"""
# cross-encoder加载,输入query-chunk对联合编码
self.reranker = CrossEncoder(model_name,
max_length=512)
self.top_k = top_k_after_rerank
self.batch_size = batch_size
def rerank(self, query, candidates):
"""
query: 用户原始查询
candidates: 粗排召回的候选chunk列表
返回: 重排后的top_k chunks
"""
if not candidates:
return []
# 1. 构造query-chunk对
pairs = [(query, chunk.text) for chunk in candidates]
# 2. cross-encoder批量打分
# 联合编码含交叉注意力,捕捉细粒度语义交互
scores = self.reranker.predict(
pairs, batch_size=self.batch_size
)
# 3. 按分数降序排序
ranked_indices = np.argsort(scores)[::-1]
# 4. 取top_k返回
reranked = []
for idx in ranked_indices[:self.top_k]:
candidates[idx].rerank_score = float(scores[idx])
reranked.append(candidates[idx])
return reranked
def rerank_with_metadata(self, query, candidates):
"""带元数据加权重排,业务定制排序"""
pairs = [(query, c.text) for c in candidates]
scores = self.reranker.predict(pairs,
batch_size=self.batch_size)
for i, candidate in enumerate(candidates):
# 基础相关性分数
relevance_score = scores[i]
# 元数据加权:时效性、权威性、来源权重
time_weight = self._time_decay(candidate.timestamp)
authority_weight = self._authority_weight(candidate.source)
# 综合分数 = 相关性 × 时效 × 权威
final_score = relevance_score * time_weight * \
authority_weight
candidate.rerank_score = float(final_score)
# 按综合分数排序
candidates.sort(key=lambda c: c.rerank_score, reverse=True)
return candidates[:self.top_k]
def _time_decay(self, timestamp, half_life_days=180):
"""时效性衰减,半年权重减半"""
import time
age_days = (time.time() - timestamp) / 86400
return 0.5 ** (age_days / half_life_days)
def _authority_weight(self, source):
"""来源权威性权重,官方文档>内部wiki>社区"""
weights = {
"official": 1.0,
"internal_wiki": 0.85,
"community": 0.6
}
return weights.get(source, 0.7)
# 量化对比:
# 向量召回top1精度:72%
# cross-encoder重排top1精度:91%
# 重排延迟:50候选约30ms,占总延迟8%可接受
量化指标与边界
某项目落地cross-encoder重排后,top1精度从72%提到91%,top5内命中率从85%提到97%(正确chunk几乎总在top5)。重排延迟30ms占端到端延迟8%,用户无感。元数据加权让时效性强的文档优先召回,过期文档自然下沉。
边界与踩坑:cross-encoder的max_length限制512 token,长chunk需截断可能丢信息,建议chunk控制在400 token内。50候选是经验,过多重排延迟线性增长(100候选约60ms),过少召回率不够。bge-reranker对中文效果好,英文场景用ms-marco-MiniLM。元数据加权需谨慎,权重失衡导致权威但过期的文档压过新鲜但非权威的,需按业务调half_life_days和authority_weight。
3. Agent规划失控是ReAct循环无终止条件
痛点现场
某Agent上线处理工单,用户问"帮我查下订单123的状态顺便看下物流",Agent规划成两步:查订单→查物流。但查物流工具返回异常,Agent进入"思考-重试-再思考"循环,10次迭代后token耗尽超时崩溃,用户等了2分钟收到错误。
ReAct循环失控的根因是模式本身没有终止条件——思考-行动-观察循环理论上可以无限进行,模型遇到错误就重试,重试又失败再重试,空耗token。团队没加迭代上限,模型自己不知道何时该放弃。
更隐蔽的是工具调用错误率,Agent调工具时参数解析错误(如把订单号"123"解析成整数123但工具要字符串"123"),工具返回错误,Agent误以为是业务错误重试,实际是参数错误越重试越错。这类错误占Agent失败的40%,是首要失败模式。
工程方案:迭代上限+收敛检测+强制兜底
方案分三层。迭代层:硬上限5次,达到即强制兜底。错误层:工具调用错误先分类,参数错误自动修正重试,超时换工具或降级。收敛层:检测输出是否含Final Answer,无收敛迹象提前终止避免无效循环。
// 来源:langchain 0.1.0 / agents/react/base + 自研安全Agent
import time
from enum import Enum
class ErrorType(Enum):
PARAM = "param_error" # 参数错误,可修正重试
TIMEOUT = "timeout" # 超时,换工具或降级
BUSINESS = "business_error" # 业务错误,非重试可解决
UNKNOWN = "unknown" # 未知错误
class SafeReActAgent:
def __init__(self, llm, tools, max_iterations=5,
token_budget=4000, tool_timeout=10):
self.llm = llm
self.tools = {t.name: t for t in tools}
self.max_iter = max_iterations
self.token_budget = token_budget
self.tool_timeout = tool_timeout
self.consecutive_errors = 0 # 连续错误计数
def run(self, query):
scratchpad = ""
total_tokens = 0
for iteration in range(self.max_iter):
# token预算检查,超预算提前终止
if total_tokens > self.token_budget:
return self._fallback(query, "token预算耗尽",
scratchpad)
# 构造ReAct prompt
prompt = self._build_prompt(query, scratchpad)
output = self.llm(prompt)
# 收敛检测:检查是否给出Final Answer
if "Final Answer:" in output:
return self._parse_final(output)
# 解析工具调用
try:
tool_name, tool_input = self._parse_action(output)
except ParseError:
# 解析失败计数,连续3次终止
self.consecutive_errors += 1
if self.consecutive_errors >= 3:
return self._fallback(query, "连续解析失败",
scratchpad)
scratchpad += f"\n{output}\nObservation: 解析失败\n"
continue
# 工具调用,带超时和错误分类
observation = self._safe_call_tool(
tool_name, tool_input, iteration
)
scratchpad += f"\n{output}\nObservation: {observation}\n"
# 成功调用重置错误计数
if "错误" not in observation:
self.consecutive_errors = 0
# 达到迭代上限强制兜底
return self._fallback(query, f"达到迭代上限{self.max_iter}",
scratchpad)
def _safe_call_tool(self, tool_name, tool_input, iteration):
"""安全工具调用,带超时和错误分类"""
if tool_name not in self.tools:
return f"工具{tool_name}不存在,可用: {list(self.tools.keys())}"
tool = self.tools[tool_name]
# 同一工具连续重试限制,避免死循环
if iteration > 0 and self._is_repeat_call(tool_name,
tool_input):
return "重复调用同一工具同一参数,停止重试"
try:
# 带超时调用
result = tool.run(tool_input, timeout=self.tool_timeout)
return result
except TimeoutError:
return "工具调用超时,请换方案或降级处理"
except ParamError as e:
# 参数错误自动修正重试一次
corrected = self._auto_fix_param(tool, tool_input, e)
if corrected and iteration < self.max_iter - 1:
try:
return tool.run(corrected,
timeout=self.tool_timeout)
except Exception:
return "参数修正后仍失败,转人工"
return f"参数错误: {e},自动修正失败"
except Exception as e:
return f"工具调用异常: {type(e).__name__}"
def _auto_fix_param(self, tool, tool_input, error):
"""参数自动修正,常见类型转换"""
# 整数→字符串:工具要字符串但传了整数
for key, val in tool_input.items():
if isinstance(val, int) and "str" in str(error):
tool_input[key] = str(val)
return tool_input
# 字符串→整数:工具要整数但传了字符串数字
if isinstance(val, str) and val.isdigit() and \
"int" in str(error):
tool_input[key] = int(val)
return tool_input
return None
def _is_repeat_call(self, tool_name, tool_input):
"""检测重复调用,避免死循环"""
# 简化实现:记录最近调用,对比是否重复
return hasattr(self, '_last_call') and \
self._last_call == (tool_name, str(tool_input))
def _fallback(self, query, reason, scratchpad):
"""强制兜底,绝不返回错误让用户裸奔"""
# 1. 尝试RAG检索已知答案
knowledge = self.rag.retrieve(query, top_k=3)
if knowledge and knowledge[0].score > 0.8:
return f"基于知识库回答: {knowledge[0].answer}\n" \
f"(Agent无法自主解决,{reason})"
# 2. 转人工
return f"无法在{self.max_iter}步内解决,{reason},转人工处理"
量化指标与边界
某项目落地安全Agent后,循环崩溃率从12%压到0%(硬上限兜底),工具调用成功率从60%提到85%(参数自动修正),平均迭代次数从7.3压到3.2(收敛检测提前终止)。max_iterations=5是经验,过少复杂任务做不完,过多空耗成本——某工单场景调到8支持多步任务,token预算配6000。
边界与踩坑:迭代上限是硬约束,复杂任务(如多步骤工单)可能5步不够,需按业务调。参数自动修正覆盖常见类型转换,复杂参数错误仍需人工兜底。收敛检测依赖LLM输出"Final Answer:"标记,模型有时不按格式输出,需配合输出解析兜底。连续错误计数3次终止,防止模型在错误模式里打转。兜底转人工会增加人工负载,需配合分流——高频兜底问题回流训练数据补强Agent能力。
4. 长对话上下文爆炸靠滑动窗口必然丢信息
痛点现场
某客服Agent对话到第8轮,用户问"刚才说的那个订单号是多少",Agent答不上——滑动窗口只保留最近5轮,第3轮说的订单号被截掉了。团队加长窗口到10轮,token爆炸每次对话消耗8000 token,成本激增且模型注意力被稀释,重要信息淹没在长上下文里反而答得更差。
滑动窗口的本质缺陷是机械截断,按轮数丢早期对话,不管丢的是不是关键信息。用户对话中的关键事实(订单号、人名、数字)经常在第1轮就确认,后续轮次依赖这些事实,但窗口截断把它们丢了,Agent失忆。
更隐蔽的是注意力稀释,即使上下文窗口够大(如32k token),模型对长上下文的注意力也是不均匀的——首尾token注意力强,中间token注意力弱(lost in the middle现象),关键信息放中间照样丢失。
工程方案:摘要压缩+近期原文双层记忆
方案是双层记忆。长期层:早期对话压缩成摘要,保留关键事实(订单号、人名、决策)。近期层:最近N轮保留原文,保留细节。两者拼接入Prompt,既不token爆炸又不失忆。补充关键事实提取表,把对话中的关键事实结构化存储,随Prompt注入确保不丢。
// 来源:langchain 0.1.0 / memory/summary + 自研双层记忆
import json
from collections import defaultdict
class HybridMemory:
def __init__(self, llm, max_recent_turns=5,
summary_max_tokens=500):
"""
max_recent_turns: 保留近期原文的轮数
summary_max_tokens: 历史摘要的token上限
"""
self.llm = llm
self.max_recent = max_recent_turns
self.summary_max = summary_max_tokens
self.summary = "" # 历史摘要
self.recent = [] # 近期对话原文
self.key_facts = defaultdict(str) # 关键事实表
def add(self, human_msg, ai_msg):
"""添加一轮对话到记忆"""
self.recent.append({
"human": human_msg,
"ai": ai_msg,
"turn": len(self.recent) + len(self._summarized_count)
})
# 提取本轮关键事实
self._extract_key_facts(human_msg, ai_msg)
# 超过近期窗口触发摘要压缩
if len(self.recent) > self.max_recent:
# 取最早的一轮滚动入摘要
oldest = self.recent.pop(0)
self._update_summary(oldest)
def _extract_key_facts(self, human_msg, ai_msg):
"""从对话中提取关键事实,结构化存储"""
# 关键事实类型:订单号、人名、数字、决策
extract_prompt = f"""
从以下对话提取关键事实,JSON格式返回:
对话:用户:{human_msg}\n助手:{ai_msg}
提取格式:{{"订单号":"", "人名":"", "金额":"",
"决策":"", "时间":""}}
无则对应字段留空。
"""
result = self.llm(extract_prompt)
try:
facts = json.loads(result)
# 累积更新事实表,新值覆盖旧值
for key, value in facts.items():
if value:
self.key_facts[key] = value
except json.JSONDecodeError:
pass # 提取失败容错,不影响主流程
def _update_summary(self, oldest_turn):
"""把最早一轮对话滚动入摘要"""
summary_prompt = f"""
已有摘要:{self.summary}
新增对话:用户:{oldest_turn['human']}
助手:{oldest_turn['ai']}
请更新摘要,保留关键信息,控制在{self.summary_max}token内。
重点保留:订单号、人名、数字、已做决策、未解决问题。
"""
self.summary = self.llm(summary_prompt)
def build_context(self, current_query):
"""构建入Prompt的上下文"""
context_parts = []
# 1. 历史摘要提供长期记忆
if self.summary:
context_parts.append(f"历史摘要:{self.summary}")
# 2. 关键事实表结构化注入,确保不丢
if self.key_facts:
facts_str = "\n".join(
f" {k}: {v}" for k, v in self.key_facts.items() if v
)
context_parts.append(f"已知关键事实:\n{facts_str}")
# 3. 近期对话原文提供细节
for turn in self.recent:
context_parts.append(
f"用户: {turn['human']}\n助手: {turn['ai']}"
)
# 4. 当前查询
context_parts.append(f"用户: {current_query}")
return "\n\n".join(context_parts)
def total_tokens(self):
"""估算当前上下文token数"""
# 简化估算:字符数/1.5近似token数
context = self.build_context("")
return len(context) // 2 # 中文约2字符1token
# 量化对比:
# 滑动窗口5轮:token 2000,关键事实丢失率40%
# 双层记忆:token 2500(摘要500+近期1500+事实500),关键事实丢失率5%
量化指标与边界
某项目落地双层记忆后,长对话(10轮+)关键事实丢失率从40%压到5%,token消耗从滑动10轮的8000压到2500(摘要500+近期5轮1500+事实表500),模型注意力稀释问题缓解。max_recent_turns=5是经验,客服场景5轮够用,复杂工单场景可到8。summary_max_tokens=500平衡摘要信息量与token成本。
边界与踩坑:摘要压缩有信息损失,重要事实(人名/数字)需在摘要prompt中明确要求保留,否则LLM可能省略。关键事实提取依赖LLM能力,提取失败需容错不影响主流程。事实表的覆盖策略是新值覆盖旧值,但同一字段多次变化时可能丢失历史值——如订单号换了,旧订单号被覆盖,需评估是否保留历史。双层记忆增加了每轮的处理延迟(摘要+事实提取各一次LLM调用约500ms),实时性要求高的场景需异步压缩。
5. 知识库构建慢是入库链路无流水化
痛点现场
某企业知识库1万份文档,首次入库花了6小时——解析PDF、清洗、切分、embedding、建索引全串行,单文档平均2秒,1万份就是5.5小时。期间知识库不可用,业务方等不了。后续每周增量更新100份,全量重跑要重算所有embedding(即使文档没变),又花2小时,团队干脆不更新,知识库内容 stale。
入库慢的根因是全量重算+串行处理。没有增量识别,每次更新都重算所有文档的embedding。没有流水化,解析、切分、embedding、建索引串行执行,GPU embedding阶段卡住CPU解析就闲着,资源利用率低。
更隐蔽的是失败重试机制缺失,100份文档有3份PDF解析失败,整个批次报错回滚,已处理的97份白算,需从头再来。失败处理粗暴导致入库可靠性差,团队怕出错干脆不更新。
工程方案:增量入库+任务流水化+失败重试
方案分三层。增量层:文档content_hash去重,已入库的hash跳过,只处理新增/变更文档。流水层:解析、切分、embedding、建索引并行流水化,CPU解析与GPU embedding并行。重试层:失败入队列异步重试,不影响整批,最终一致。
// 来源:haystack 2.0 / indexing + 自研增量流水
import hashlib
import time
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
from collections import deque
class IncrementalIndexer:
def __init__(self, vector_store, embedder, splitter,
num_workers=4):
self.vector_store = vector_store
self.embedder = embedder
self.splitter = splitter
self.num_workers = num_workers
self.failed_queue = deque() # 失败重试队列
self.max_retries = 3 # 最大重试次数
def ingest(self, documents):
"""增量流水化入库主流程"""
# 第一阶段:增量识别,过滤已入库文档
new_docs = []
for doc in documents:
content_hash = self._content_hash(doc)
if not self.vector_store.exists(content_hash):
doc.content_hash = content_hash
new_docs.append(doc)
# 已入库跳过,避免重复计算
print(f"增量识别:{len(documents)}份中{len(new_docs)}份需处理")
if not new_docs:
return {"processed": 0, "skipped": len(documents),
"failed": 0}
# 第二阶段:CPU解析清洗切分(ProcessPool并行)
with ProcessPoolExecutor(max_workers=self.num_workers) as pool:
chunk_results = list(pool.map(
self._parse_clean_split, new_docs
))
# 收集成功chunk和失败文档
all_chunks = []
failed_docs = []
for doc, result in zip(new_docs, chunk_results):
if result.get("error"):
failed_docs.append((doc, result["error"]))
else:
for chunk in result["chunks"]:
chunk.content_hash = doc.content_hash
all_chunks.append(chunk)
# 第三阶段:GPU embedding(批量并行,充分利用GPU)
embeddings = self._batch_embed(all_chunks)
# 第四阶段:建索引入库(原子批量写入)
self.vector_store.batch_upsert(
chunks=all_chunks,
embeddings=embeddings,
content_hashes=[c.content_hash for c in all_chunks]
)
# 失败文档入队列异步重试
for doc, error in failed_docs:
self.failed_queue.append({
"doc": doc, "error": error, "retries": 0
})
return {
"processed": len(new_docs) - len(failed_docs),
"skipped": len(documents) - len(new_docs),
"failed": len(failed_docs)
}
def _parse_clean_split(self, doc):
"""单文档解析清洗切分,独立进程隔离失败"""
try:
# 解析PDF/Word/HTML
text = self._parse(doc.path)
# 清洗去噪(去页眉页脚、OCR错误修正)
cleaned = self._clean(text)
# 语义切分
chunks = self.splitter.split(cleaned)
return {"chunks": chunks}
except Exception as e:
# 单文档失败不影响整批,返回错误标记
return {"error": str(e), "chunks": []}
def _batch_embed(self, chunks, batch_size=64):
"""GPU批量embedding,充分利用并行"""
all_embeddings = []
# 分批处理,每批送GPU批量编码
for i in range(0, len(chunks), batch_size):
batch = chunks[i:i + batch_size]
texts = [c.text for c in batch]
# GPU批量编码,比逐条快10倍
embeddings = self.embedder.encode(
texts, batch_size=len(batch),
show_progress_bar=False
)
all_embeddings.extend(embeddings)
return all_embeddings
def _content_hash(self, doc):
"""文档内容hash,内容变更即识别"""
with open(doc.path, 'rb') as f:
return hashlib.md5(f.read()).hexdigest()
def retry_failed(self):
"""异步重试失败文档,指数退避"""
retry_queue = list(self.failed_queue)
self.failed_queue.clear()
for item in retry_queue:
if item["retries"] >= self.max_retries:
# 超过重试上限,记录日志转人工
self._log_persistent_failure(item)
continue
# 指数退避等待
backoff = 2 ** item["retries"]
time.sleep(backoff)
# 重试入库
result = self.ingest([item["doc"]])
if result["failed"] > 0:
# 仍失败,重试计数+1回队列
item["retries"] += 1
self.failed_queue.append(item)
# 量化对比:
# 全量串行:1万份6小时
# 增量流水:首次1万份1.5小时(4worker并行+GPU批量),后续100份3分钟
# 失败重试:3份失败不影响97份,异步重试最终一致
量化指标与边界
某项目落地增量流水后,首次1万份入库从6小时压到1.5小时(4worker并行+GPU批量embedding),后续每周100份增量从2小时压到3分钟(只处理变更文档)。失败重试让97份成功不被3份失败拖累,入库可靠性从70%提到99%。content_hash去重让全量重算变成增量计算,embedding计算量减少95%。
边界与踩坑:content_hash只识别内容变更,不识别切分规则变更——切分规则变了需全量重入库,建议切分规则版本化,规则变更触发全量。ProcessPool隔离单文档失败,但进程创建开销大,小批量文档不如直接串行+try。GPU批量embedding的batch_size=64是经验,过大显存溢出,过小GPU利用率低。失败重试的指数退避避免雪崩,但持久失败的文档需人工介入,不能无限重试。增量入库依赖vector_store的exists查询性能,百万级hash查询需建索引,否则去重阶段慢。
总结
RAG与Agent层的本质是检索精度与决策可控的博弈。语义边界切分把召回精度从72%提到89%,cross-encoder重排把top1精度从72%提到91%,迭代上限+收敛检测把Agent循环崩溃率从12%压到0,双层记忆把长对话事实丢失率从40%压到5%,增量流水把入库从6小时压到1.5小时且增量3分钟。五个支点都有量化指标与边界,落地顺序建议:召回切分与重排先行(决定答案质量上限),Agent安全与记忆紧随(决定可用性),知识库流水作为运营终局支撑持续更新。
更多推荐



所有评论(0)