Agent引用数据库知识过时的增量同步方案

发布时间:2026/7/27 2:45:56
Agent引用数据库知识过时的增量同步方案 数据库里的知识持续新增、修改、删除而 Agent 仍引用旧内容本质是向量索引与源头数据不同步。解决思路不是更频繁地全量重建而是搭建一条从源头变更到向量库、再到 Agent 检索生成的端到端增量同步链路并在 Agent 侧加一层自我反思 人工闭环来兜底。下面给出完整设计、步骤和代码。一、整体架构六层同步链路源数据库 ──CDC── 消息队列(Kafka) ── 嵌入服务 ── 向量库(带元数据版本) │ ▼ LangGraph Agent: [检索] ─ [相关性评分] ─ [生成] ─ [事实校验] ─ [人工审核闭环] │ ▼ 监控 / 评测 / 回滚核心原则增量更新负责日常变化定期全量重建负责长期健康索引时用的 Embedding 模型必须与查询时用的模型一致。二、步骤 1源头变更捕获CDC传统定时拉取存在轮询间隔与实时性成反比、空轮询浪费资源的问题。推荐 CDCChange Data CaptureMySQL/PostgreSQLDebezium 监听 binlog/WALMongoDBChange Streams带 resume token支持断点续传API 服务通过 Kafka 发布/订阅传递变更事件每条变更事件必须携带operation_typeinsert/update/delete、doc_id、以及整行数据。 关键每一次 insert/update/delete 都必须触发向量库的对应动作否则已下架商品/已删除条款还会被召回。三、步骤 2知识抽取与版本化Chunk 级别这是防引用过时最核心的一步。不要做文档级粗粒度更新要做Chunk 级增量3.1 文档切分与指纹每个 chunk 计算 SHA-256 作为内容寻址身份importhashlibdefnormalize(text:str)-str:归一化去空白、转小写保证 hash 确定性return .join(text.strip().lower().split())defchunk_id(doc_id:str,position:int,content:str)-str:hhashlib.sha256(normalize(content).encode()).hexdigest()returnf{doc_id}:{position}:{h[:16]}defsplit_with_metadata(row:dict,doc_id:str):把数据库行转为带元数据的 chunk 列表# 只把语义有意义的列拼进嵌入文本text_to_embedf{row[name]}.{row[description]}# 数值/状态/时间戳作为元数据不参与嵌入metadata{doc_id:doc_id,price:row.get(price),category:row.get(category),status:row.get(status),# active / discontinuedupdated_at:row.get(updated_at),version:row.get(version,1),}return[{content:text_to_embed,metadata:metadata}]3.2 变更检测逻辑在 Hash Store 中维护doc_id - [chunk_hash1, chunk_hash2, ...]的映射比较新旧版本检测结果含义动作新 hash该位置新增内容计算嵌入 插入向量库hash 变化内容被修改删除旧向量 插入新向量旧 hash 消失内容被删除从向量库删除hash 不变内容未变跳过省下嵌入成本这样可以把嵌入计算量从 O©全量重嵌降到 O(ΔC)仅变部分。四、步骤 3嵌入与向量库 Upsert / Delete4.1 CDC 消费者代码fromkafkaimportKafkaConsumerimportjson consumerKafkaConsumer(db_changes,bootstrap_servers[localhost:9092],value_deserializerlambdav:json.loads(v.decode()),group_idembedding_service,auto_offset_resetlatest,)formsginconsumer:eventmsg.value opevent[op]# c(insert) / u(update) / d(delete)doc_idevent[doc_id]rowevent.get(after,{})ifopd:# 硬删除直接移除该 doc 所有 chunk 向量vector_store.delete(filter{doc_id:doc_id})hash_store.pop(doc_id,None)continue# insert / update切分 指纹比对new_chunkssplit_with_metadata(row,doc_id)old_hasheshash_store.get(doc_id,[])fori,chunkinenumerate(new_chunks):chunk_hashchunk_id(doc_id,i,chunk[content])ifchunk_hashinold_hashes:continue# 未变化跳过# 先删旧若存在vector_store.delete(filter{doc_id:doc_id,position:i})# 嵌入 写入embeddingembed_model.embed(chunk[content])vector_store.upsert(idchunk_hash,vectorembedding,payload{**chunk[metadata],position:i,valid_from:row[updated_at],valid_to:None,content:chunk[content],})# 更新 hash 映射hash_store[doc_id][chunk_id(doc_id,i,c[content])fori,cinenumerate(new_chunks)]4.2 向量库元数据 Schema防过时关键每个向量必须带以下元数据字段doc_id源文档 IDversion单调自增版本号updated_at最后更新时间戳valid_from/valid_to版本生效时间窗TTL 淘汰依据status业务状态如商品已下架source_id溯源用⚠️ 最常见的坑只插入新向量不删除旧版本。这会导致新旧政策同时被召回用户拿到过期答案——比查不到还糟。删除文档时必须同步清理其所有 chunk 向量否则会残留幽灵文档。4.3 嵌入模型一致性硬规则索引时用的 Embedding 模型必须与查询时用的模型完全一致。换模型时必须全量重建。五、步骤 4Agent 检索层——过滤 重排5.1 元数据过滤排除过期内容defretrieve(query:str,top_k:int5):embeddingembed_model.embed(query)resultsvector_store.query(vectorembedding,top_ktop_k*3,# 多召回后续重排filter{status:{$eq:active},# 排除已下架valid_to:{$isnull:True},# 排除已被取代的旧版本updated_at:{$gte:cutoff_time},# 可选时间窗过滤})returnrerank(query,results)[:top_k]5.2 重排序用 Cross-Encoder 重排提升准确率fromsentence_transformersimportCrossEncoder rerankerCrossEncoder(BAAI/bge-reranker-v2-m3)defrerank(query:str,candidates:list):pairs[[query,c[content]]forcincandidates]scoresreranker.predict(pairs)forc,sinzip(candidates,scores):c[rerank_score]sreturnsorted(candidates,keylambdax:x[rerank_score],reverseTrue)六、步骤 5LangGraph Agent 自我反思Self-RAG光有同步还不够Agent 生成时仍需校验。用 LangGraph 构建带反思节点的图fromlanggraph.graphimportStateGraph,START,ENDfromtypingimportTypedDict,Listfromlangchain_core.messagesimportHumanMessageclassAgentState(TypedDict):query:strdocuments:List[dict]generation:strrelevance_grade:str# relevant / irrelevantsupport_grade:str# fully_supported / partially / no_support# 节点 1检索defretrieve_node(state:AgentState):docsretrieve(state[query])return{documents:docs}# 节点 2文档相关性评分defgrade_documents(state:AgentState):llmget_llm()prompt判断以下文档是否与问题相关。 问题{q} 文档{d} 只回答 relevant 或 irrelevant。.format(qstate[query],dstate[documents])gradellm.invoke([HumanMessage(contentprompt)]).contentreturn{relevance_grade:grade}# 节点 3生成defgenerate(state:AgentState):ifstate[relevance_grade]irrelevant:# 改写查询重新检索return{generation:,relevance_grade:rewrite_needed}llmget_llm()ctx\n.join([d[content]fordinstate[documents]])promptf基于以下上下文回答问题。若上下文矛盾优先使用 updated_at 更晚的内容。\n\n上下文{ctx}\n\n问题{state[query]}genllm.invoke([HumanMessage(contentprompt)]).contentreturn{generation:gen}# 节点 4事实校验防止幻觉引用过时内容defgrade_generation_vs_documents(state:AgentState):llmget_llm()prompt校验生成内容是否被检索文档充分支持。 文档{d} 生成{g} 只回答 fully_supported / partially / no_support。.format(dstate[documents],gstate[generation])gradellm.invoke([HumanMessage(contentprompt)]).contentreturn{support_grade:grade}# 构建图workflowStateGraph(AgentState)workflow.add_node(retrieve,retrieve_node)workflow.add_node(grade_docs,grade_documents)workflow.add_node(generate,generate)workflow.add_node(grade_generation,grade_generation_vs_documents)workflow.add_edge(START,retrieve)workflow.add_edge(retrieve,grade_docs)workflow.add_conditional_edges(grade_docs,lambdas:generateifs[relevance_grade]relevantelseretrieve,{generate:generate,retrieve:retrieve}# 改写查询后重检)workflow.add_edge(generate,grade_generation)workflow.add_conditional_edges(grade_generation,lambdas:ENDifs[support_grade]fully_supportedelsegenerate,{generate:generate,END:END})appworkflow.compile()这套 Self-RAG 流程能在检索质量差时自动改写查询重检在生成与文档不符时重新生成大幅降低引用过时/错误内容的概率。七、步骤 6人工闭环Human-in-the-Loop对于高风险的生成结果用 LangGraph 的interrupt()暂停等待人工审核fromlanggraph.typesimportinterrupt,Commanddefhuman_review_node(state:AgentState)-Command:# 中断等待人工审批decisioninterrupt({action:approve_or_edit,draft:state[generation],sources:[d[doc_id]fordinstate[documents]],})ifdecision[type]accept:returnCommand(gotoEND)elifdecision[type]edit:# 用户编辑后的内容可作为新知识写回向量库updated_contentdecision[edited_response]# 触发增量索引更新trigger_reindex(state[query],updated_content)returnCommand(update{generation:updated_content},gotoEND)else:# rejectreturnCommand(gotoEND)人工编辑后的内容自动存入向量库形成越用越准的闭环。八、步骤 7监控、评测与回滚8.1 关键监控指标同步延迟CDC 事件产生到向量库生效的时间差孤儿向量数向量库中存在但源库已删除的 chunk 数检索命中率含valid_to ! null旧版本的比例应趋近 0生成支持率Self-RAG 中fully_supported占比8.2 蓝绿部署 自动回滚defblue_green_switch():新旧索引并行原子切换# 新索引构建完成后通过配置中心一键切换流量# 旧索引保留 24-48 小时异常时秒级回滚pass8.3 定期全量重建增量更新为主但以下场景必须全量重建Embedding 模型升级每周/每月的健康检查处理索引碎片故障恢复# 每周日凌晨全量重建app.task(schedule0 2 * * 0)deffull_rebuild():vector_store.rebuild_from_source(source_dbdb,embed_modelembed_model,cleanupfull)九、最终总结保证 LangChain/LangGraph Agent 引用知识不过时本质是把同步做成一等公民而非事后补救。整套方案的关键设计决策1. 同步策略日常用CDC Chunk 级增量hash 指纹检测变更只处理 ΔC周期性全量重建维护索引健康Embedding 模型必须前后一致换模型即全量重建2. 防过时三道防线写入防线delete 事件必须传播到向量库旧版本向量带valid_to标记检索防线元数据过滤排除status ! active且valid_to ! null的向量生成防线Self-RAG 相关性评分 事实校验矛盾时优先采用updated_at更晚的内容3. 工程化保障向量库元数据 Schema 必须含doc_id / version / updated_at / valid_from / valid_to蓝绿部署实现零停机更新异常秒级回滚人工审核闭环让系统越用越准监控同步延迟、孤儿向量、检索命中率4. 成本与实时性权衡业务场景推荐策略商品库存/价格CDC 实时同步秒级客服知识库CDC 每分钟批量内部制度文档每小时/每天增量合规条款CDC 实时 人工审核闭环 一句话记忆法文档变了 → chunk 必须变 → 向量必须变 → 元数据版本必须对齐。三者任一脱节Agent 就会引用过时知识。按这套架构落地你的 Agent 就能在数据库持续新增/修改/删除的情况下始终引用最新、最准确的知识——而不是停留在第一天导入时的那个版本。