在这里插入图片描述

摘要

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加overlap冗余

语义完整chunk

方案分两步。切分:按句切分后计算句间语义相似度,相似度高于阈值的归入同一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重排+粗精两段

LLM生成 Cross-encoder精排 向量库 向量粗排 用户Query LLM生成 Cross-encoder精排 向量库 向量粗排 用户Query 粗排保召回率,精排提精度 query embedding 检索top50 候选chunks query-chunk对精排 top5重排结果 精排top5喂LLM 答案

方案是粗精两段。粗排用向量召回取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%,是首要失败模式。

工程方案:迭代上限+收敛检测+强制兜底

参数错误

超时

有Final Answer

小于上限

大于上限

用户问题

迭代1: 思考+行动

工具调用成功

观察结果

错误分类

修正参数重试

换工具或降级

收敛检测

输出答案

迭代次数

强制兜底

转人工或返回已知

方案分三层。迭代层:硬上限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现象),关键信息放中间照样丢失。

工程方案:摘要压缩+近期原文双层记忆

历史对话10轮

是否超近期窗口

全量原文入Prompt

最早轮次摘要压缩

摘要入长期记忆

近期轮次保留原文

摘要+近期原文入Prompt

关键事实提取表

随Prompt注入

方案是双层记忆。长期层:早期对话压缩成摘要,保留关键事实(订单号、人名、决策)。近期层:最近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

建索引

入库成功

失败队列

异步重试

hash登记

方案分三层。增量层:文档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安全与记忆紧随(决定可用性),知识库流水作为运营终局支撑持续更新。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐