回顾 -> 问题引出
- DeepRearchSystem 0x00:初识
- DeepRearchSystem 0x01:Agent 基础
- DeepRearchSystem 0x02:Graph 构建
- DeepRearchSystem 0x03:HITL
- DeepResearchSystem 0x04:MAS 进阶
- DeepResearchSystem 0x06:LLM as Judge
- DeepResearchSystem 0x07:查询缓存
上一章] 实现了 DeepRearchSystem 的搜索缓存。解决了短时间内同样主题重复请求的问题。现在又遇到了新的挑战:
search_cache本身解决的是短时间内的重复请求问题,我们设计的TTL也是小时级别。那么如果是第二天出现了同样的query该怎么办呢?是否只能按正常流程重新跑一遍。search_cache的 key设计是:(prompt, count),本质是query-prompt的精确匹配。比如 result1 = query(“AI芯片销量”), 下次 query(“GPU出货量”):因为 prompt 不能匹配,就必须重新跑一遍。但是呢,我们都知道 “AI芯片销量” 和 “GPU出货量”是息息相关的,上次的检索结果对这次查询完全有可复用的资料。
3.search_cache 解决的是短期内并发请求的重复调用问题。一旦跨会话/跨项目,DeepRearchSystem 每次只能重来。Agent 是有能力但只会蛮干,就像 “鱼的记忆只有7s”,如果能让 DeepRearchSystem 像人一样有记忆就好了。
怎么办?(思考方向)
其实,到这里要做的事情就有点清晰了:要让 Agent 具备长期记忆,当 query 来临时先查自己的记忆,记忆里没有再查短期缓存,短期缓存没有再走搜索流程重新沉淀。沉淀后的知识再反哺记忆。概括起来需要具备:
- 具备长期记忆,可以优先从记忆中检索
- 不局限于精确匹配,而是支持语义匹配。
咦,听起来是不是像 RAG 知识库。确实有点像,但还是有点差别的。至于哪些地方有差别,我们先去设计和实现之后再回来看待这个问题。
干他 (架构设计)
事实提取器
extractor: 事实提取器。从搜索摘要中精准剥离出离散、可验证的事实,并自动打标(置信度、类别等)
- 置信度:对事实进行验证,记忆只存储经验证的有置信度打分的数据
- 类别:用于
lifecycle对数据的 TTL 管理,不同类型数据的过期时间不同
事实记忆器
fact_store: 事实存储器。基于 Milvus 的向量数据库存储。
vector选型:支持高效的语义检索
TTL 架构
lifecycle: TTL 规则器。决定哪些类型的数据有效期短,哪些类型有效期长。用来保证知识新鲜度。
年龄阈值
FRESHNESS_MAX_AGE: 年龄阈值(天)
high: 7。7 天有效期medium: 30low: 180
生命周期
KBLifecycleMode: KB模式
- off/关闭:不进行过滤、不衰减、不添加年龄标记。
- inform/通知:添加年龄标记并更改措辞,但不进行过滤或置信度衰减
- 开关打开,为搜索数据增加
age_days时效时效时效 - 根据
age_days参数修改 Prompt:
* “知识库中已有的相关事实(请勿重复搜索这些内容)”
* ”标记较早的事实可能已过时,请优先搜索获取最新信息“
- 开关打开,为搜索数据增加
- freshness/新鲜度:Plan LLM判定 → 统一的年龄过滤和置信度衰减。
- 置信度衰减:随着时间推移记忆数据的置信度理应衰减
* 曲线衰减:decay = e^(-lambda * age)
* 类别敏感度:对应market_data到historical分布不同
- 置信度衰减:随着时间推移记忆数据的置信度理应衰减
- lifecycle/生命周期:事实类别 → 按类别进行生命周期时间 (TTL) 过滤和衰减。
market_data7, 市场数据——7天过期product_info: 30, 产品信息——30天strategy: 90, 战略方向——90天technology: 180, 技术定义一般有效期更长——180天historical: None, 历史事件——永不过期
至此,整个 KB 框架就构思好了,接下来就是搬砖🧱实现了。
Coding
extractor
核心代码
注意最终返回的结构:
fact:事实数据source_url:数据来源 (必须是真实 URL, 而不是前面流程构造的短url)confidence:置信度,由 LLM 判断的 事实可信度category:事实类型,是技术类还是产品类
1class FactExtractor: 2 """ 3 从搜索结果摘要中提取结构化事实。 4 """ 5 6 def __init__(self, model_id="deepseek-v4-pro"): 7 self.model_id = model_id 8 9 def extract(self, summary: str, research_topic: str = "") -> list[dict]: 10 """ 11 从单个搜索结果摘要中提取事实。 12 13 Args: 14 summary: 网页搜索结果文本(已由 LLM 进行摘要) 15 research_topic: 用于提供上下文的整体研究主题 16 17 Returns: 18 List of {fact, source_url, confidence} dicts. 19 """ 20 if not summary or len(summary) < 50: 21 logger.debug("[KB-extractor] 摘要过短,省略提取") 22 return [] 23 24 agent = Agent(model_id=self.model_id) if self.model_id else Agent() 25 agent.set_step_prompt(EXTRACTION_INSTRUCTIONS) 26 27 raw = agent.step( 28 research_topic=research_topic or "(extracting facts)", 29 summary=summary[:8000], # truncate for safety 30 ) 31 32 try: 33 json_str = JsonUtils.extract_pattern(raw, pattern="json") 34 facts = json.loads(json_str) 35 if isinstance(facts, list): 36 facts = self._validate(facts) 37 logger.info( 38 f"[KB-extractor] extracted {len(facts)} facts " 39 f"from summary ({len(summary)} chars)" 40 ) 41 return facts 42 except json.JSONDecodeError as exc: 43 logger.warning( 44 f"[KB-extractor] JSON parse failed — LLM returned malformed JSON: {exc}" 45 ) 46 except (ValueError, TypeError) as exc: 47 logger.warning( 48 f"[KB-extractor] data validation failed ({type(exc).__name__}): {exc}" 49 ) 50 except Exception as exc: 51 logger.warning( 52 f"[KB-extractor] failed to parse facts ({type(exc).__name__}): {exc}" 53 ) 54 55 return [] 56 57 @staticmethod 58 def _validate(facts: list) -> list[dict]: 59 """ 60 过滤和规整提取的facts. 61 """ 62 valid = [] 63 for f in facts: 64 if not isinstance(f, dict): 65 continue 66 fact_text = f.get("fact", "").strip() 67 if not fact_text or len(fact_text) < 10: 68 continue 69 valid.append({ 70 "fact": fact_text, 71 "source_url": f.get("source_url", ""), 72 "confidence": min(1.0, max(0.0, float(f.get("confidence", 0.7)))), 73 "fact_category": f.get("fact_category", "strategy"), 74 }) 75 if len(valid) >= 10: 76 break 77 return valid 78
Prompt
那么 LLM 是怎么提取出以上事实的?就依赖于我们写的 Prompt。
1EXTRACTION_INSTRUCTIONS = """# 任务说明 2你是一个知识提取专家。从给定的搜索摘要中提取出离散的、可验证的事实陈述。 3 4# Instruction 51. 每条 fact 必须是一个独立的、可验证的陈述(一句话),避免复合句 62. 每条 fact 标注来源 URL(从摘要中的 [来源名](url) 格式提取) 73. 每条 fact 标注 confidence(0.0-1.0): 8 - 0.9-1.0:摘要明确引用具体数据/事件 9 - 0.7-0.8:摘要中有较明确的表述 10 - 0.5-0.6:摘要中隐含但未直接说明 114. 标注每条 fact 的 category(事实时效性类别): 12 - market_data: 市场份额、价格、排名、增长率、营收数据 13 - product_info: 产品功能、规格、版本、发布信息 14 - strategy: 公司战略、投资方向、合作、收购 15 - technology: 技术原理、架构定义、标准、协议 16 - historical: 历史事件、里程碑、已发生的事实 175. 不要提取过于泛泛的陈述(如"这是一个重要市场") 186. 每条 fact 控制在 80 字以内 197. 最多提取 10 条 fact 20 21# 输出格式 22```json 23[ 24 { 25 "fact": "2025年Q1全球AI编程助手市场规模达到15亿美元", 26 "source_url": "https://xxxx.com/xxxx", 27 "confidence": 0.9, 28 "category": "market_data" 29 }, 30 { 31 "fact": "GitHub Copilot 占据AI编程助手市场约35%的份额", 32 "source_url": "https://xxxx.com/xxxx", 33 "confidence": 0.8, 34 "category": "market_data" 35 } 36]
研究主题(用于上下文理解)
{research_topic}
搜索摘要
{summary}
输出"""
### fact\_store
基于 Milvus 的研究知识库事实存储库。
#### Milvus
和之前我们使用 `ChromaDB` 一样,`Milvus` 使用基本流程相同。
```Python
@property
def client(self) -> MilvusClient:
if self._client is None:
self._client = MilvusClient(uri=self.uri)
return self._client
def _ensure_collection(self) -> None:
"""创建集合(如果不存在)。"""
# 检查集合是否已存在
try:
if self.client.has_collection(self.collection):
logger.info(f"[KB] collection '{self.collection}' already exists")
return
except Exception as e:
raise KBConnectionError(
f"Milvus 连接失败(检查集合是否存在时): {e}"
) from e
# 创建集合
try:
self.client.create_collection(
collection_name=self.collection,
dimension=self.embedding_dim,
metric_type="COSINE",
auto_id=True,
enable_dynamic_field=True,
)
except Exception as e:
msg = str(e).lower()
if "dimension" in msg or "param" in msg or "schema" in msg:
raise KBConfigError(
f"Milvus 集合创建配置错误: {e}"
) from e
raise KBConnectionError(
f"Milvus 集合创建失败: {e}"
) from e
# 创建 IVF_FLAT 索引以提高搜索效率。
# Note: 某些 Milvus 版本会自动创建默认索引。
try:
index_params = IndexParams()
index_params.add_index(
field_name="vector",
index_type="IVF_FLAT",
metric_type="COSINE",
params={"nlist": 128},
)
self.client.create_index(
collection_name=self.collection,
index_params=index_params,
)
self.client.load_collection(self.collection)
except Exception as exc:
# Milvus may raise if index already exists; this is harmless
error_msg = str(exc).lower()
if "already exist" in error_msg or "duplicate" in error_msg:
logger.debug(f"[KB] index already exists, skipping creation")
else:
logger.warning(
f"[KB] index creation skipped ({type(exc).__name__}): {exc}"
)
logger.info(
f"[KB] created collection '{self.collection}' "
f"(dim={self.embedding_dim}, metric=COSINE)"
)
Embedding
这块就很熟悉了,和之前 RAG 项目中的向量化流程一毛一样。这里要注意几点:
base_url拼接要注意检查给定厂商是否带了v1- 做好异常捕获和处理
- 实际dim与配置dim的差异:即使输出警告排查
- 重试策略:按错误码分情况处理,429/5xx 可重试,重试采取指数退避算法
- Json 解析的错误处理
1def _embed(self, texts: list[str]) -> list[list[float]]: 2 """ 3 从兼容 OpenAI 的端点获取嵌入向量。 4 5 最多重试 3 次,失败后采用退避策略。 6 7 配置优先级: 8 - EMBEDDING_BASE_URL → LLM_BASE_URL(fallback) 9 - EMBEDDING_API_KEY → APP_TOKEN(fallback) 10 """ 11 base_url = os.getenv("EMBEDDING_BASE_URL") or os.getenv("LLM_BASE_URL") 12 api_key = os.getenv("EMBEDDING_API_KEY") or os.getenv("APP_TOKEN") 13 14 if not base_url: 15 raise KBConfigError( 16 "Embedding URL 未配置:请设置 EMBEDDING_BASE_URL 或 LLM_BASE_URL 环境变量" 17 ) 18 if not api_key: 19 raise KBConfigError( 20 "Embedding API Key 未配置:请设置 EMBEDDING_API_KEY 或 APP_TOKEN 环境变量" 21 ) 22 23 # 大多数兼容 OpenAI 的端点都支持 /v1/embeddings 24 # 如果base URL 已经以 /v1 结尾,需要相应调整。 25 if base_url.endswith("/v1"): 26 url = f"{base_url}/embeddings" 27 else: 28 url = f"{base_url}/v1/embeddings" 29 30 headers = { 31 "Authorization": f"Bearer {api_key}", 32 "Content-Type": "application/json", 33 } 34 payload = { 35 "model": self.embedding_model, 36 "input": texts, 37 } 38 39 last_exc: Exception | None = None 40 for attempt in range(3): 41 try: 42 logger.debug(f"[KB] embedding {len(texts)} texts via {url} (model={self.embedding_model}, attempt={attempt + 1})") 43 resp = requests.post(url, json=payload, headers=headers, timeout=30) 44 resp.raise_for_status() 45 46 data = resp.json() 47 embeddings = [item["embedding"] for item in data["data"]] 48 # 按索引排序以保持顺序 49 embeddings.sort(key=lambda x: x.get("index", 0) if isinstance(x, dict) else 0) 50 result = [item["embedding"] if isinstance(item, dict) else item for item in embeddings] 51 52 # 自动检测实际dim与配置dim的差异 53 if result and len(result[0]) != self.embedding_dim: 54 logger.warning( 55 f"[KB] embedding dim mismatch: configured={self.embedding_dim}, " 56 f"actual={len(result[0])}. Update EMBEDDING_DIM env var." 57 ) 58 59 return result 60 61 except requests.HTTPError as e: 62 last_exc = e 63 status = e.response.status_code if e.response is not None else 0 64 65 # 不可恢复:认证/权限/参数错误 — 不重试 66 if status in (401, 403): 67 logger.error( 68 f"[KB] embedding auth error HTTP {status}, not retrying: {e}" 69 ) 70 raise KBEmbeddingFatalError( 71 f"Embedding API 认证/权限失败 (HTTP {status})" 72 ) from e 73 if status == 400: 74 logger.error( 75 f"[KB] embedding bad request HTTP 400, not retrying: {e}" 76 ) 77 raise KBEmbeddingFatalError( 78 f"Embedding API 参数错误 (HTTP 400)" 79 ) from e 80 81 # 可恢复:429 / 5xx — 重试 82 if status == 429 and attempt < 2: 83 wait = 3 * (attempt + 1) 84 logger.warning(f"[KB] embedding 429 rate-limited, retrying in {wait}s (attempt {attempt + 1}/3)") 85 time.sleep(wait) 86 elif status >= 500 and attempt < 2: 87 wait = 2 * (attempt + 1) 88 logger.warning(f"[KB] embedding server error {status}, retrying in {wait}s (attempt {attempt + 1}/3)") 89 time.sleep(wait) 90 else: 91 logger.error(f"[KB] embedding HTTP {status}, no more retries: {e}") 92 raise KBEmbeddingError( 93 f"Embedding API HTTP {status} 重试耗尽" 94 ) from e 95 96 except (requests.ConnectionError, requests.Timeout) as e: 97 last_exc = e 98 if attempt < 2: 99 wait = 1.5 * (attempt + 1) 100 logger.warning(f"[KB] embedding network error, retrying in {wait:.1f}s (attempt {attempt + 1}/3): {e}") 101 time.sleep(wait) 102 else: 103 logger.error(f"[KB] embedding network error, no more retries: {e}") 104 raise KBEmbeddingError( 105 f"Embedding 网络错误重试耗尽: {e}" 106 ) from e 107 108 except (ValueError, KeyError, TypeError) as e: 109 # JSON 解析 / 数据结构错误 — 永久错误,不重试 110 logger.error( 111 f"[KB] embedding response parse error ({type(e).__name__}), " 112 f"not retrying: {e}" 113 ) 114 raise KBEmbeddingFatalError( 115 f"Embedding 响应解析失败: {e}" 116 ) from e 117 118 except Exception as e: 119 # 未知异常 — 保守重试 120 last_exc = e 121 if attempt < 2: 122 wait = 1.5 * (attempt + 1) 123 logger.warning( 124 f"[KB] embedding unexpected error ({type(e).__name__}), " 125 f"retrying in {wait:.1f}s (attempt {attempt + 1}/3): {e}" 126 ) 127 time.sleep(wait) 128 else: 129 logger.error( 130 f"[KB] embedding unexpected error ({type(e).__name__}), " 131 f"no more retries: {e}" 132 ) 133 raise KBEmbeddingError( 134 f"Embedding 未知错误重试耗尽 ({type(e).__name__}): {e}" 135 ) from e 136 137 raise KBEmbeddingError( 138 f"[KB] embedding failed after 3 attempts: {last_exc}" 139 ) from last_exc 140 141async def _aembed(self, texts: list[str]) -> list[list[float]]: 142 """ 143 异步 embedding——用 asyncio.to_thread 包裹同步方法,不阻塞事件循环. 144 """ 145 return await asyncio.to_thread(self._embed, texts) 146
主类
1class FactStore: 2 """ 3 Milvus-backed storage and retrieval of research facts. 4 """ 5 def __init__( 6 self, 7 uri: str | None = None, 8 collection: str = COLLECTION_NAME, 9 embedding_dim: int | None = None, 10 embedding_model: str | None = None, 11 ): 12 self.uri = uri or os.getenv("MILVUS_URI", "http://localhost:19530") 13 self.collection = collection 14 self.embedding_dim = embedding_dim or int( 15 os.getenv("EMBEDDING_DIM", str(DEFAULT_EMBEDDING_DIM)) 16 ) 17 self.embedding_model = embedding_model or os.getenv( 18 "EMBEDDING_MODEL", DEFAULT_EMBEDDING_MODEL 19 ) 20 21 self._client: MilvusClient | None = None 22 self._ensure_collection() 23
set
好像没什么可说的,就是 数据的分块向量化存储那一块。
1def add_facts(self, facts: list[dict]) -> int: 2 """ 3 将fact嵌入并插入到 Milvus 中。返回已插入的事实数量。 4 5 每个fact字典必须包含: 6 - fact (str): 事实陈述 7 - source_url (str): 事实来源 8 Optional: 9 - research_topic (str): 触发此搜索的研究主题 10 - confidence (float): 0.0-1.0 11 """ 12 if not facts: 13 return 0 14 15 texts = [f["fact"] for f in facts] 16 embeddings = self._embed(texts) 17 18 data = [] 19 now = int(time.time()) 20 for i, fact in enumerate(facts): 21 data.append({ 22 "vector": embeddings[i], 23 "fact_text": fact["fact"], 24 "source_url": fact.get("source_url", ""), 25 "research_topic": fact.get("research_topic", ""), 26 "confidence": float(fact.get("confidence", 1.0)), 27 "category": fact.get("category", "strategy"), 28 "created_at": now, 29 }) 30 31 result = self.client.insert(collection_name=self.collection, data=data) 32 inserted = result.get("insert_count", len(data)) 33 logger.info( 34 f"[KB] stored {inserted} facts → Milvus/{self.collection} " 35 f"(topics: {set(f.get('research_topic', '') for f in facts)})" 36 ) 37 return inserted 38
get
查询的时候要注意根据 lifecycle 策略做好年龄时效性和置信度衰减的过滤。同时置信度衰减要按不同类型的数据设置不同的衰减系数。
1def query( 2 self, 3 topic: str, 4 top_k: int = 10, 5 min_confidence: float = 0.0, 6 max_age_days: int | None = None, 7 decay: bool = False, 8 lifecycle_mode: bool = False, 9) -> list[dict]: 10 """ 11 对与研究主题相关的fact进行语义搜索。 12 13 Args: 14 topic: 用于语义搜索的研究主题文本。 15 top_k: 期望返回的事实数量。若启用 reranker,Milvus 会召回 16 max(top_k * 3, 20) 条候选,再由 reranker 精排到 top_k。 17 min_confidence: 过滤前的最小置信度(0.0-1.0)。 18 max_age_days: 排除超过指定天数的事实。None 表示无限制。 19 在 lifecycle_mode 模式下,此参数将被忽略,而是使用按类别划分的 TTL(生存时间)。 20 decay: 如果为 True,则应用基于时间的置信度衰减。 21 lifecycle_mode: 如果为 True,则使用按类别划分的 TTL(CATEGORY_TTL), 22 而不是单一的 max_age_days 阈值。 23 24 Returns list of {fact, source_url, confidence, research_topic, relevance, 25 created_at, age_days, category}. 26 """ 27 28 embedding = self._embed([topic]) 29 milvus_limit = max(top_k * 3, 20) 30 results = self.client.search( 31 collection_name=self.collection, 32 data=[embedding[0]], 33 limit=milvus_limit, 34 output_fields=[ 35 "fact_text", "source_url", "research_topic", 36 "confidence", "created_at", "category", 37 ], 38 ) 39 40 if not results or not results[0]: 41 logger.info(f"[KB] query '{topic[:60]}...' → 0 results") 42 return [] 43 44 # ── 收集 Milvus 原始候选 ────────────────────────────────── 45 now = time.time() 46 raw_candidates = [] 47 for hit in results[0]: 48 entity = hit.get("entity", {}) 49 distance_score = 1.0 - hit.get("distance", 0) 50 raw_candidates.append({ 51 "entity": entity, 52 "distance_score": distance_score, 53 }) 54 55 candidates = raw_candidates 56 57 # ── 置信度 & 时效过滤 ───────────────────────────────────── 58 hits = [] 59 for cand in candidates: 60 entity = cand["entity"] 61 62 raw_confidence = entity.get("confidence", 1.0) 63 if raw_confidence < min_confidence: 64 continue 65 66 created = entity.get("created_at", 0) 67 age_days = (now - created) / 86400 if created else 36500 68 69 # ── 年龄时效性过滤 ── 70 if lifecycle_mode: 71 category = entity.get("category", "strategy") 72 category_max_age = CATEGORY_TTL.get(category) 73 if category_max_age is not None and age_days > category_max_age: 74 continue 75 elif max_age_days and age_days > max_age_days: 76 continue 77 78 # ── 置信度衰减 ── 79 confidence = raw_confidence 80 if decay: 81 category = entity.get("category", "strategy") 82 ttl = CATEGORY_TTL.get(category, 180) 83 84 # 1. 如果超过TTL,直接归零(除非是historical) 85 if ttl is not None and age_days > ttl: 86 confidence = 0.0 87 else: 88 # 2. 指数衰减公式: decay = e^(-lambda * age) 89 # lambda 是衰减系数,不同类别不同 90 decay_lambda = { 91 "market_data": 0.1, # 衰减极快 92 "product_info": 0.05, # 衰减快 93 "strategy": 0.02, # 衰减中等 94 "technology": 0.005, # 衰减慢 95 "historical": 0.0 # 不衰减 96 }.get(category, 0.01) 97 98 decay_factor = math.exp(-decay_lambda * age_days) 99 # 设置一个硬下限,防止完全消失(除非超TTL) 100 decay_factor = max(0.1, decay_factor) if category != "historical" else 1.0 101 confidence = raw_confidence * confidence * decay_factor 102 103 # relevance 优先用 rerank_score,否则用 Milvus 距离得分 104 orig_idx = cand.get("_orig_idx") 105 106 hits.append({ 107 "fact": entity.get("fact_text", ""), 108 "source_url": entity.get("source_url", ""), 109 "research_topic": entity.get("research_topic", ""), 110 "confidence": round(confidence, 2), 111 "created_at": created, 112 "age_days": round(age_days, 1), 113 "category": entity.get("category", "strategy"), 114 }) 115 116 logger.info( 117 f"[KB] query '{topic[:60]}...' → {len(hits)} hits " 118 f"(top score={hits[0]['relevance']:.3f})" 119 if hits else f"[KB] query '{topic[:60]}...' → 0 hits" 120 ) 121 return hits 122
lifecycle
基本定义
1class KBLifecycleMode(str, Enum): 2 OFF = "off" 3 INFORM = "inform" 4 FRESHNESS = "freshness" 5 LIFECYCLE = "lifecycle" 6 7 8# ── 统一年龄阈值 (freshness mode) ─────────────────────────── 9FRESHNESS_MAX_AGE: dict[str, int] = { 10 "high": 7, 11 "medium": 30, 12 "low": 180, 13} 14 15# ── 按类别划分的 TTL (lifecycle mode) ───────────────────────────────── 16CATEGORY_TTL: dict[str, int | None] = { 17 "market_data": 7, # 市场数据——7天过期 18 "product_info": 30, # 产品信息——30天 19 "strategy": 90, # 战略方向——90天 20 "technology": 180, # 技术定义——180天 21 "historical": None, # 历史事件——永不过期 22} 23
模式获取
1def get_mode() -> KBLifecycleMode: 2 """ 3 从环境变量中读取当前生命周期模式。 4 5 如果值无效或未设置,则回退到 FRESHNESS。 6 """ 7 raw = os.getenv("KB_LIFECYCLE_MODE", "freshness") 8 try: 9 return KBLifecycleMode(raw) 10 except ValueError: 11 return KBLifecycleMode.FRESHNESS 12
关键状态
should_filter:当且仅当返回True,fact_store-query才需要对记忆做时效性和置信度过滤should_decay:当且仅当返回True,fact_store-query才需要对记忆的置信度做时间衰减处理should_tag:当且仅当返回True,research节点才需要给事实添加人类可读的时效性标签should_warn:当且仅当返回True,research节点才会鼓励提示词(prompt)措辞重新验证
1def should_filter(mode: KBLifecycleMode | None = None) -> bool: 2 """ 3 当应该根据事实的时效性来排除事实时,返回 True 4 """ 5 if mode is None: 6 mode = get_mode() 7 return mode in (KBLifecycleMode.FRESHNESS, KBLifecycleMode.LIFECYCLE) 8 9 10def should_decay(mode: KBLifecycleMode | None = None) -> bool: 11 """当需要对置信度分数进行基于时间的衰减时,返回 True 12 """ 13 if mode is None: 14 mode = get_mode() 15 return mode in (KBLifecycleMode.FRESHNESS, KBLifecycleMode.LIFECYCLE) 16 17 18def should_tag(mode: KBLifecycleMode | None = None) -> bool: 19 """ 20 当需要给事实添加人类可读的时效性标签时,返回 True 21 """ 22 if mode is None: 23 mode = get_mode() 24 return mode != KBLifecycleMode.OFF 25 26 27def should_warn(mode: KBLifecycleMode | None = None) -> bool: 28 """ 29 当提示词(prompt)的措辞应该鼓励重新验证,而不是禁止重新搜索时,返回 True""" 30 if mode is None: 31 mode = get_mode() 32 return mode != KBLifecycleMode.OFF 33
调用
初始化
KB 知识库应用主要在子图 reasearch_graph,我们在对应文件下实现初始化。
1# ── KB 单例(延迟初始化,在代理运行之间共享) ────────────── 2_kb_store: FactStore | None = None 3_kb_extractor: FactExtractor | None = None 4 5_GENERATE_QUERIES = "generate_queries" 6_WEB_SEARCH = "web_search" 7_CRITIQUE = "critique" 8 9def _get_kb_extractor() -> FactExtractor: 10 """ 11 获取 KB 事实生成器 12 :return: 13 """ 14 global _kb_extractor 15 if _kb_extractor is None: 16 _kb_extractor = FactExtractor() 17 return _kb_extractor 18 19def _get_kb_store() -> FactStore | None: 20 """ 21 获取 KB 事实存储器 22 :return: 23 """ 24 global _kb_store 25 if _kb_store is None: 26 try: 27 _kb_store = FactStore() 28 logger.info("[KB] FactStor 连接到 Milvus") 29 except (KBConnectionError, KBConfigError) as exc: 30 logger.warning(f"[KB] FactStore 初始化失败(KB 不可用,将静默降级): {exc}") 31 _kb_store = False # type: ignore — sentinel 32 except Exception as exc: 33 logger.warning(f"[KB] FactStore 初始化失败(未知错误,将静默降级): {exc}") 34 _kb_store = False # type: ignore — sentinel 35 return _kb_store if _kb_store is not False else None 36
set
存储这里需要在 _web_search 节点,因为只有在这里才会返回搜索数据。这里要注意构造 KB 事实的数据源必须是原始可用 URL。
1# ── KB/知识库 存储 ──────────────────────────────────────────────── 2try: 3 store = _get_kb_store() 4 if store: 5 extractor = _get_kb_extractor() 6 topic = get_research_topic(state.get("messages", [])) 7 facts = extractor.extract(summary, research_topic=topic) 8 if facts: 9 # 将短链接还原为真实 URL,避免 KB 中存储不可解析的过期引用 10 short2long = {v: k for k, v in long2shorts.items()} 11 for f in facts: 12 f["source_url"] = short2long.get(f["source_url"], f["source_url"]) 13 f["research_topic"] = topic 14 store.add_facts(facts) 15except (KBConnectionError, KBEmbeddingError) as exc: 16 logger.warning(f"[KB] 跳过存储(瞬时错误): {exc}") 17except (KBConfigError, KBEmbeddingFatalError) as exc: 18 logger.error(f"[KB] 跳过存储(永久错误,需人工修复): {exc}") 19except Exception as exc: 20 logger.warning(f"[KB] 跳过存储(未知错误): {exc}") 21
get
Code
对应地,在 _generate_search 节点优先读取 KB 记忆,搜索生成器只 针对 KB 记忆 known_facts 以外的知识生成相关搜索主题。 也就是基于当前的 KB 记忆去补充搜索。
1# ── KB/知识库检索 ────────────────────────────────────────────── 2known_facts_text = "" 3try: 4 store = _get_kb_store() 5 if store: 6 mode = get_mode() 7 topic = get_research_topic(state["messages"]) 8 freshness = state.get("fresh_level", "medium") 9 max_age = FRESHNESS_MAX_AGE.get(freshness, 30) if should_filter(mode) else None 10 decay = should_decay(mode) 11 use_lifecycle = mode == KBLifecycleMode.LIFECYCLE 12 13 # 命中缓存 14 hits = store.query( 15 topic, top_k=20, min_confidence=0.6, 16 max_age_days=max_age, 17 decay=decay, 18 lifecycle_mode=use_lifecycle, 19 ) 20 if hits: 21 facts_lines = [] 22 for h in hits: 23 line = f"- [{h['confidence']:.0%}] {h['fact']}" 24 # 添加时效性标签 25 if should_tag(mode): 26 age_days = h.get("age_days", (time.time() - h["created_at"]) / 86400) 27 age_tag = ( 28 "🕐 刚刚" if age_days < 1 else 29 f"{age_days:.0f}天前" if age_days < 30 else 30 f"{age_days / 30:.0f}个月前" 31 ) 32 line += f" ({age_tag}, 来源: {h['source_url'][:60]})" 33 else: 34 line += f" (来源: {h['source_url'][:60]})" 35 facts_lines.append(line) 36 # 提示词加强处理:是否需要鼓励重新验证 37 if should_warn(mode): 38 header = "\n## 📚 知识库中已有的相关事实\n" 39 footer = "\n\n⚠️ 标记较早的事实可能已过时,请优先搜索获取最新信息。" 40 else: 41 header = "\n## 知识库中已有的相关事实(请勿重复搜索这些内容)\n" 42 footer = "" 43 known_facts_text = header + "\n".join(facts_lines) + footer 44 logger.info(f"[KB] 检索到 {len(hits)} 个facts 用作查询生成上下文") 45except (KBConnectionError, KBEmbeddingError) as exc: 46 logger.warning(f"[KB] retrieval skipped(瞬时错误,下次可能恢复): {exc}") 47except (KBConfigError, KBEmbeddingFatalError) as exc: 48 logger.error(f"[KB] retrieval skipped(永久错误,需人工修复): {exc}") 49except Exception as exc: 50 logger.warning(f"[KB] retrieval skipped(未知错误): {exc}") 51logger.info(f"[ResearchAgent] _generate_queries使用模型: {configuration.query_generator_model}") 52 53agent = JsonAgent(model_id=configuration.query_generator_model, keys=SearchQueryList) 54agent.set_step_prompt(query_writer_instructions) 55result = agent.step( 56 current_date=get_current_date(), 57 research_topic=get_research_topic(state["messages"]), 58 number_queries=state["initial_search_query_count"], 59 research_proposal=state.get("plan", ""), # 这里加了人类确定的研究计划 60 known_facts = known_facts_text 61) 62
Prompt改造
上面不难发现原来的提示词模板里多了一个 known_facts 参数,这里需要对原来的 query_writer_instructions 针对性补充。这里需要注意的是:
- 新增
known_facts参数,同时必须和## KB知识库已知事实(已验证的结构化数据,带时效标签)做显示绑定,让 Prompt 知道它是 KB 知识 - 新增基于
known_facts的拆解约束,拆出的搜索主题必须满足:- 禁止重复搜索 KB 已验证且时效充足的事实
- 对 KB 过时事实生成「更新验证」类查询
- 基于 KB 已知事实填补维度缺口
- 如果 KB 已知事实之间存在矛盾,或你认为已知事实可能不准确,生成用于交叉验证的搜索主题,并在
rationale中说明
1uery_writer_instructions = """# 角色定义 2你是一个主题研究拆解大师,你擅长将给定的研究主题拆解为不同的子主题,并给出这么拆解的理由,最后按照给定的格式输出结果 3# 任务说明 4你的任务是根据当前的研究主题决定多个用于网络搜索的标题,这些标题会被用于从网页搜集信息,并整合成一份专业的研究报告 5 6# Instruction 7- 针对当前的研究主题,你可以将其拆解成若干个搜索主题,每个搜索主题都应该是针对当前研究主题不同维度的切分 8- 针对当前研究主题,最多不产生{number_queries}条搜索主题 9- 你的搜索主题应该尽可能的广泛,如果研究主题本身就非常宽泛,则产出1条以上的搜索主题 10- 每个搜索主题应该具备独立性,即不要同时产出多个相似或者耦合的搜索主题 11- 搜索主题应该考虑时间,即除非研究主题要求,不然尽可能搜集近期的资料,当前时间是{current_date} 12 13- KB知识库已知事实 使用规则(重要) 14**本规则专用于下方 `# 背景信息 → ## KB知识库已知事实` 区块中的内容。** 15**注意:请勿将`研究计划`中的假设、推论或目标,误认为是`KB知识库已知事实`。** 16 17下述「已知事实」是从KB知识库中检索到的、与研究主题相关且已经过验证的事实。 18你必须按以下规则使用它们: 191. **禁止重复搜索已覆盖且时效充足的事实** 20 对于已知事实中标注为「刚刚/X天前」(且属于 strategy / technology / historical 类别)的内容,不得生成与之相同或高度相似的搜索主题。 21 222. **对过时事实生成「更新验证」类查询** 23 对于已知事实中标注为「X个月前」或属于 market_data / product_info 类别的内容,你必须生成用于**获取最新数据**的搜索主题(例如在原主题前加「2025年最新」、「近期更新」等时间限定),而不是忽略它。 24 253. **基于已知事实填补维度缺口** 26 已知事实可能只覆盖了研究主题的部分维度。你必须识别尚未被已知事实覆盖的维度(如:市场规模、主要厂商、政策环境、技术瓶颈、未来趋势等),并针对这些**缺口维度**生成搜索主题。 27 284. **冲突处理** 29 如果已知事实之间存在矛盾,或你认为已知事实可能不准确,生成用于**交叉验证**的搜索主题,并在 rationale 中说明。 30 31# Output Format 32你生成的内容应该是一个标准的json格式的内容,并包含两个字端 33<param> 34 <attribute>rationale</attribute> 35 <type>string</type> 36 <description>你的思考,即为什么要产出如下的几个搜索主题</description> 37</param> 38 39<param> 40 <attribute>query</attribute> 41 <type>List</type> 42 <description>用于做网络搜索的搜索主题</description> 43</param> 44下面是一个输出样例 45```json 46{ 47 "rationale": "xxxx", 48 "query": ["搜索主题1", "搜索主题2", ...] 49}
背景信息
研究课题
{research_topic}
研究计划(未经证实的假设与目标)
{research_proposal}
KB知识库已知事实(已验证的结构化数据,带时效标签)
{known_facts}
Output"""
## KB 知识 🆚 search\_cache 搜索缓存
### 对比
我们在前面实现了 [search\_cache 搜索缓存](https://link.juejin.cn?target=), 本章开头也有过简单的讨论。我们在完全实现玩 KB 知识记忆之后再系统地看下两个缓存方式。
首先,来看下加入二者后 DeepResearchSystem 的整体执行流程:
```sql
[User Query]
│
▼
[1. KB Pre-retrieval] ──▶ 语义查询KB
│ ├─ Hit (High Confidence, Fresh): 直接作为已知事实传入Prompt
│ └─ Miss / Low Confidence: 继续下一步
│
▼
[2. Search Cache Check] ──▶ 精确匹配(prompt, count)
│ ├─ Hit: 直接返回Raw Snippets -> 跳至 [4. Extractor]
│ └─ Miss: 继续下一步
│
▼
[3. Web Search MCP] ──▶ 调用昂贵API
│
▼
[4. Extractor] ──▶ 从Raw Snippets中提取结构化Facts
│
▼
[5. KB Update] ──▶ 将新Fact存入KB (Add/Update)
│
▼
[6. Search Cache Set] ──▶ 将Raw Snippets存入Search Cache
│
▼
[7. Agent Reasoning] ──▶ 基于KB中的结构化事实生成报告
我们不难发现:
- Search Cache 保护的是 Web Search。它挡住了99%的重复搜索请求。
- KB 保护的是 Agent 的智力水平。它挡住了幻觉和低质量信息。
- Extractor 是连接二者的桥梁。它将 Search Cache 里的“原材料”加工成 KB 里的“成品”。
| 维度 | Search Cache (search_cache.py) | Knowledge Base (KB System) |
|---|---|---|
| 心智模型 | 速度至上,成本优先,最终一致性。 核心思想是“空间换时间” | 质量至上,可信优先,事实一致性。 核心思想是“验证即真理”,强调知识的准确性、可追溯性和长效性。 |
| 解决的痛点 | 1. API成本爆炸:防止同一Prompt在毫秒级被并发请求多次调用昂贵的Web Search API。 2. 网络抖动/Redis故障:通过本地内存缓存降级,保证单实例在Redis宕机时仍能工作。 3. 数据分裂:多实例部署时,通过Redis共享缓存,并通过回填机制解决Redis重启后各实例数据不一致的问题。 | 1. 大模型幻觉:防止Agent编造事实,所有结论必须有据可依。 2. 知识断层:避免每次启动Agent都从零开始,实现跨会话、跨项目的知识积累。 3. 信息过时:通过TTL和置信度衰减,确保决策基于最新信息,而非陈旧的训练数据。 |
| 数据本质 | Raw Data(原材料) : 未经清洗的网页摘要(Title, Snippet, URL)。包含噪音、广告、过时信息。 | Refined Knowledge(黄金) : 经过Extractor提取的结构化原子事实(Fact Text, Confidence, Category)。 |
| 匹配模式 | Key精确匹配 : MD5(prompt + count)。 | Embedding 语义匹配 |
| TTL设计 | 极短 (Short-lived) : 默认TTL=7200秒(2小时)。 网页内容瞬息万变,长时间缓存等于主动使用脏数据。 | 分级长期 (Long-lived) : 基于CATEGORY_TTL。 历史事件永不过期,市场数据7天过期。符合知识本身的客观寿命。 |
| 一致性模型 | 最终一致性 : 通过get_cached中的回填机制(Backfill) 实现。本地缓存可能比Redis新,此时会触发回写,保证Redis最终是最新的。 | 强一致性 : 写入即真理。冲突时通过置信度或来源权重仲裁。 |
| 失效策略 | 被动过期 : 设定失效时间 | 主动治理 : 置信度衰减、新旧事实冲突消解 |
| 故障处理 | 优雅降级 : Redis挂了 -> 降级为本地内存缓存 -> 继续服务。 | **熔断与告警 ** : KB挂了 -> Agent主流程不影响(长步骤高成本) -> 触发告警,人工介入。 |
🤔思考🤔🤔
- KB 与 Search Cache 的联动失效怎么处理 ?
- 当 KB 更新了一条高置信度的 Fact 时(例如修正了一个错误数据),理论上应该清除 Search Cache 中所有与该 Fact 相关的Key。否则走到 Web_search 时拿缓存,又拿到了低置信度数据。
- Cache 匹配模式是 key,这显然很难做到...
- 目前我能想到的就是:
* 直接全放弃 Cache ,重头搞起。(好像略有点浪费。。。)
* Prompt 加强: ”请忽略缓存中的旧数据,使用KB中的最新事实“
更多 AI 技术干货 请订阅 AI技术手札专栏。
《DeepResearchSystem 0x08:KB 知识记忆》 是转载文章,点击查看原文。
