Day 38 / 共 60 天 · 阶段6 记忆与知识
RAG 子系统:向量库客户端的统一抽象层
记忆(memory/)用的是自带的 LanceDB 存储;而 CrewAI 还有一套独立的 RAG 子系统(rag/),专门给"知识库"(knowledge/,D39)当底座。它把 ChromaDB、Qdrant 等各种向量库统一抽象成一个 BaseClient 协议:建集合、加文档、语义搜、删集合。今天读这套抽象:BaseClient 接口、文档怎么生成稳定 doc_id、add_documents 的 upsert + 批处理、以及最关键的 distance(距离)怎么换算成 score(相似分)。
📍 你在 60 天里的位置(阶段6 记忆与知识 · 共 8 天)
D33 记忆总览→
D34 unified_memory→
D35 recall/encoding→
D36 短/长/entity→
D37 scope 作用域→
D38 RAG→
D39 knowledge→
D40 embedding/存储
💡 先用一个类比兜住今天
rag/ 就像一个"数据库驱动"标准(类比 JDBC/ODBC):不管你后面接的是 ChromaDB 还是 Qdrant,上层代码只对着一个统一接口 BaseClient 编程——"给我建个集合""把这些文档塞进去""按这句话搜最像的 5 条"。换库只要换驱动实现,上层一行不改。而"搜出来"的向量库返回的是距离(越小越像),RAG 层帮你翻译成相似分(越大越像,0~1),因为人更习惯"分越高越好"。L01
痛点:记忆有存储了,为什么还要一套 RAG?
🤔 痛点D34 我们看到记忆自己就有一套 LanceDB 存储和向量检索,那为什么 CrewAI 还要额外搞一个
rag/ 子系统?两套不是重复吗?——不重复。它们服务两种不同的东西:记忆是 Agent 运行时动态产生的经验(带 scope/importance/新鲜度这些"记忆语义");而知识库(knowledge)是你预先灌进去的静态资料(PDF、CSV、文档),只需要"存 chunk + 按查询搜 chunk"这种纯粹的 RAG。用途不同,抽象也就不同。💡 一句话本质
rag/ = 面向"知识库"的通用向量库抽象:一个 BaseClient 协议统一 ChromaDB/Qdrant,数据单元是简单的 BaseRecord(content + 可选 metadata + doc_id),没有 scope/importance/新鲜度这些记忆专属概念。它更薄、更通用,是 D39 knowledge 的地基。| memory/(记忆) | rag/(知识库底座) | |
|---|---|---|
| 数据来源 | Agent 运行时动态产生 | 预先灌入的静态资料 |
| 数据单元 | MemoryRecord(scope/importance/时间…) | BaseRecord(content/metadata/doc_id) |
| 默认后端 | LanceDB | ChromaDB |
| 检索 | 复合评分 + Flow 自适应深度 | 纯向量相似搜 + 阈值过滤 |
L02
rag/ 目录分层:core / 后端 / embeddings
目录结构(lib/crewai/src/crewai/rag/)分得很清楚:
rag/
├── core/
│ ├── base_client.py # ★BaseClient 协议(所有后端的契约)
│ ├── base_embeddings_callable.py # EmbeddingFunction 协议
│ └── base_embeddings_provider.py # 嵌入 provider 基类(D40)
├── types.py # BaseRecord / SearchResult / EmbeddingFunction
├── chromadb/ # ★默认后端 ChromaDB 实现
│ ├── client.py # ChromaDBClient(BaseClient)
│ ├── config.py / factory.py # 配置与工厂
│ └── utils.py # 文档准备 / distance→score / 名称清洗
├── qdrant/ # 备选后端 Qdrant 实现
├── embeddings/ # 嵌入 provider 工厂 + 十几个 provider(D40)
│ ├── factory.py
│ └── providers/{openai,cohere,ollama,...}
├── config/ # 全局 RAG 配置(选后端、可选依赖)
└── factory.py # create_client:按 config 造出对应后端 client
core/ = 契约只放协议/基类:BaseClient(向量库操作)、EmbeddingFunction(嵌入)。上层只依赖这里,不碰具体后端。chromadb/ · qdrant/ = 实现各后端一个目录,实现 BaseClient。默认用 chromadb。加新后端只要新增一个目录 + 实现协议。embeddings/ = 嵌入插件十几个嵌入 provider(OpenAI/Cohere/Ollama…)+ 工厂(D40 主角)。向量库和"怎么把文字变向量"解耦。factory.py = 组装create_client(config) 按配置把"后端 + 嵌入函数"组装成一个可用 client。图注:上层只依赖 core 契约;具体后端和嵌入 provider 都是可替换实现,由 factory 组装。
L03
BaseClient:所有向量库的统一契约
核心协议 BaseClient(rag/core/base_client.py:66),方法都用 Unpack[TypedDict] 收关键字参数:
# rag/core/base_client.py:66(节选)
@runtime_checkable
class BaseClient(Protocol):
client: Any # 底层真实的向量库 client
embedding_function: EmbeddingFunction # 文字→向量
@classmethod
def __get_pydantic_core_schema__(cls, _src, _handler) -> CoreSchema:
return core_schema.any_schema() # ★让 Protocol 能塞进 pydantic 模型字段
@abstractmethod
def create_collection(self, **kwargs: Unpack[BaseCollectionParams]) -> None: ...
@abstractmethod
def get_or_create_collection(self, **kwargs: Unpack[BaseCollectionParams]) -> Any: ...
@abstractmethod
def add_documents(self, **kwargs: Unpack[BaseCollectionAddParams]) -> None: ...
@abstractmethod
def search(self, **kwargs: Unpack[BaseCollectionSearchParams]) -> list[SearchResult]: ...
@abstractmethod
def delete_collection(self, **kwargs: Unpack[BaseCollectionParams]) -> None: ...
@abstractmethod
def reset(self) -> None: ...
# 每个都有 a* 异步版:acreate_collection / aadd_documents / asearch ...
Protocol + @abstractmethod又是协议(D33 见过):定契约不定实现。ChromaDBClient(BaseClient) 这里显式继承,但本质仍是"实现这些方法即可"。Unpack[TypedDict] 参数★方法参数不是散的 name, docs, batch_size,而是用 TypedDict 描述"允许哪些关键字"(BaseCollectionAddParams 等),编辑器能补全、类型能检查,又能各后端加自己的扩展字段。__get_pydantic_core_schema__★让这个 Protocol 能当 pydantic 字段类型用(返回 any_schema),不必开 arbitrary_types_allowed——D39 的 KnowledgeStorage._client: BaseClient 就靠它。同步 + 异步双份每个操作都有 a* 异步版——知识入库/检索可同步可异步,看调用场景。大白话这就是"接口编程"的教科书写法:先把"一个向量库该会干哪些活"写成一张清单(协议),谁想当后端就把这张清单上的活都干了。上层(knowledge)只对着清单调用,从来不知道背后是 Chroma 还是 Qdrant——想换库,换个实现清单的类就行。
L04
BaseRecord 与稳定 doc_id:内容哈希
入库前,文档要被"准备"成向量库能吃的三元组,关键是生成稳定的 doc_id(rag/chromadb/utils.py:56):
# rag/chromadb/utils.py:72(节选)
for doc in documents:
if "doc_id" in doc:
doc_id = str(doc["doc_id"]) # ① 显式给了就用
else:
metadata = doc.get("metadata")
if metadata and isinstance(metadata, dict) and "doc_id" in metadata:
doc_id = str(metadata["doc_id"]) # ② metadata 里有也用
else:
content_for_hash = doc["content"]
if metadata:
metadata_str = json.dumps(metadata, sort_keys=True) # 稳定序列化
content_for_hash = f"{content_for_hash}|{metadata_str}"
doc_id = hashlib.sha256(content_for_hash.encode()).hexdigest() # ③ 内容哈希
三级 doc_id 来源优先用显式 doc_id,其次 metadata 里的,最后拿内容算 SHA-256。BaseRecord 就是 {content, doc_id?, metadata?} 这么个简单字典。sort_keys=True把 metadata 序列化进哈希时排序键——保证同样的 metadata 无论字典顺序如何,算出的哈希一致。内容哈希做 id★妙处:同样的内容 → 同样的 id。配合下一讲的 upsert,同一段资料重复灌入不会产生重复记录(同 id 覆盖)——天然幂等。💡 设计取舍①:为什么用"内容哈希"当默认 ID,而不是自增/随机 UUID?
随机 UUID 做法:每条文档一个新 id。问题是——你把同一个 PDF 重新灌一遍,就会产生一整份重复的 chunk,检索时同样内容出现两次。内容哈希做法:id 由内容决定,重复灌入 = 相同 id =
upsert 覆盖,天然去重、天然幂等。这对"知识库经常整体重建"的场景太重要了:不用先清空再灌,直接重灌就行,一样的内容不会翻倍。代价是内容变一个字 id 就变(算新记录)——但这正符合"内容变了就是新知识"的直觉。L05
add_documents:upsert + 分批入库
ChromaDB 的入库实现(rag/chromadb/client.py:309):
# rag/chromadb/client.py:309(节选)
def add_documents(self, **kwargs: Unpack[BaseCollectionAddParams]) -> None:
if not _is_sync_client(self.client): # ★类型守卫(L07)
raise TypeError("Synchronous method add_documents() requires a ClientAPI. "
"Use aadd_documents() for AsyncClientAPI.")
collection_name = kwargs["collection_name"]
documents = kwargs["documents"]
batch_size = kwargs.get("batch_size", self.default_batch_size) # 默认 100
if not documents:
raise ValueError("Documents list cannot be empty")
with self._locked(): # 跨进程锁
collection = self.client.get_or_create_collection(
name=_sanitize_collection_name(collection_name),
embedding_function=self.embedding_function)
prepared = _prepare_documents_for_chromadb(documents) # 生成 ids/texts/metadatas
for i in range(0, len(prepared.ids), batch_size): # ★分批
batch_ids, batch_texts, batch_metadatas = _create_batch_slice(
prepared=prepared, start_index=i, batch_size=batch_size)
collection.upsert(ids=batch_ids, documents=batch_texts, # ★upsert 而非 add
metadatas=batch_metadatas)
get_or_create_collection集合(collection,类比一张表)不存在就建、存在就取——传入 embedding_function,之后 chromadb 会自动用它把 documents 转向量。_sanitize_collection_name集合名也要清洗(chromadb 对名字有长度/字符/大小写约束,utils.py:296)——和 memory scope 名清洗同理。分批 batch_size=100★一次灌太多会撞 embedding API 的 token 上限。按 100 条一批切片循环——大知识库也能稳稳灌进去。upsert 不是 add★用 upsert:同 id 存在就更新、不存在就插入。配合 L04 的内容哈希 id,实现"重灌不重复"的幂等入库。嵌入是"延迟到 chromadb 内部做"的:这里只传
documents=文本 和 embedding_function,chromadb 在 upsert 时自己调 embedding_function 把文本转向量。所以 RAG 层不用手动嵌入——和 memory 手动 embed_texts 再存不同(这是两套设计的又一处差异)。L06
search:距离 → 相似分的翻译
搜索先查 collection,拿到的是"距离",再翻译成"分数"(rag/chromadb/client.py:410):
# rag/chromadb/client.py:410(节选)
def search(self, **kwargs) -> list[SearchResult]:
if "limit" not in kwargs: kwargs["limit"] = self.default_limit # 默认 5
if "score_threshold" not in kwargs: kwargs["score_threshold"] = self.default_score_threshold # 0.6
params = _extract_search_params(kwargs)
collection = self.client.get_or_create_collection(
name=_sanitize_collection_name(params.collection_name),
embedding_function=self.embedding_function)
where = params.where if params.where is not None else params.metadata_filter
results = collection.query(query_texts=[params.query], n_results=params.limit,
where=where, ...) # chromadb 返回 distances
return _process_query_results(collection=collection, ...)
翻译公式在 _convert_distance_to_score(rag/chromadb/utils.py:167):
# rag/chromadb/utils.py:167
def _convert_distance_to_score(distance, distance_metric) -> float:
"""Convert ChromaDB distance to similarity score. [0,1], 1=most similar."""
if distance_metric == "cosine":
score = 1.0 - 0.5 * distance # 余弦距离 [0,2] → 分 [0,1]
return max(0.0, min(1.0, score))
if distance_metric == "l2":
score = 1.0 / (1.0 + distance) # L2 距离 [0,∞) → 分 (0,1]
return max(0.0, min(1.0, score))
raise ValueError(f"Unsupported distance metric: {distance_metric}")
为什么要翻译向量库返回距离(越小越像),但人和上层习惯分数(越大越像、0~1)。这个函数统一成"分越高越好",score_threshold 才好用。cosine: 1-0.5·dist余弦距离范围 [0,2](0=完全同向、2=完全反向)。1-0.5·dist 把它线性映射到 [1,0],即"距离 0 → 分 1"。l2: 1/(1+dist)欧氏距离 [0,∞)。1/(1+dist) 把它压到 (0,1],距离 0 → 分 1,距离越大分越趋近 0。和 memory LanceDB 的换算(1/(1+distance),D40)同一思路。clamp + 阈值过滤max(0,min(1,·)) 夹到 [0,1];再由 score_threshold(默认 0.6)把太不像的滤掉(utils.py:239)。图注:查询→嵌入→向量库返回距离→翻译成 0~1 相似分→阈值过滤→SearchResult。
L07
同步/异步双客户端:类型守卫防误用
ChromaDB 有同步 ClientAPI 和异步 AsyncClientAPI 两种,ChromaDBClient 都支持,但每个方法开头都守卫(rag/chromadb/client.py:328):
# rag/chromadb/client.py:328(同步方法开头)
if not _is_sync_client(self.client):
raise TypeError("Synchronous method add_documents() requires a ClientAPI. "
"Use aadd_documents() for AsyncClientAPI.")
# rag/chromadb/client.py:379(异步方法开头)
if not _is_async_client(self.client):
raise TypeError("Asynchronous method aadd_documents() requires an AsyncClientAPI. "
"Use add_documents() for ClientAPI.")
并发保护用跨进程锁(rag/chromadb/client.py:78):
# rag/chromadb/client.py:78
def _locked(self) -> AbstractContextManager[None]:
return store_lock(self._lock_name) if self._lock_name else nullcontext()
@asynccontextmanager
async def _alocked(self): # 异步版:在 executor 里拿/放同步锁
if not self._lock_name: yield; return
lock_cm = store_lock(self._lock_name)
loop = asyncio.get_event_loop()
await loop.run_in_executor(None, lock_cm.__enter__)
try: yield
finally: await loop.run_in_executor(None, lock_cm.__exit__, None, None, None)
_is_sync_client 守卫拿一个异步 client 调同步方法(或反过来)会得到协程没 await 之类的诡异 bug。守卫提前抛清晰 TypeError,告诉你"该用 aadd_documents"。_locked / nullcontext有锁名就用跨进程锁(多进程共享同一个 chromadb 目录时防写冲突),没锁名就 nullcontext(什么都不做的空上下文)——优雅地"可选加锁"。_alocked 在 executor 里拿同步锁★异步方法要用一个同步的文件锁:直接在事件循环里阻塞拿锁会卡住整个 loop,所以丢到 run_in_executor(线程池)里拿,不阻塞协程调度。L08
取舍 + 今日小结
💡 设计取舍②:为什么记忆和知识库不共用一套存储抽象?
看起来它俩都是"向量存储",为什么不合并成一套?因为关注点不同。记忆需要 scope 隔离、importance/新鲜度评分、LLM 合并、后台异步写——这些是"记忆语义",绑死在
StorageBackend(返回 MemoryRecord)上。知识库只要"存 chunk、按查询搜 chunk、带 metadata 过滤"——更通用,用 BaseClient(返回 SearchResult)。硬合并成一套,要么记忆抽象被塞进无关的 scope/importance 概念(污染通用性),要么知识库被迫背上记忆的复杂度(过度设计)。两套各自贴合自己的领域模型,是"不为复用而牺牲清晰"的务实取舍——代价是两套相似代码,但边界清楚、各自能独立演进。⚠️ 边界:distance metric 必须和建集合时一致
_convert_distance_to_score 要传 distance_metric(cosine/l2)——它从 collection 的配置里读(utils.py:271)。坑在于:如果你建集合用的是 cosine,检索时却按 l2 公式换算,分数就全错(尺度和方向都不对),阈值过滤会莫名其妙全滤掉或全放过。源码通过"从 collection 元数据读取实际 metric"来避免写死,但如果你手动建/迁移集合,务必保证 metric 前后一致。遇到 "Unsupported distance metric" 直接抛错也是提醒:ip(内积)等其他 metric 这里没实现,别乱配。🧠 今天你应该能回答
- rag/ 和 memory/ 分别服务什么?为什么两套?
- rag/ 目录三层(core/后端/embeddings)各是什么?
- BaseClient 有哪些核心方法?Unpack[TypedDict] 参数的好处?
- doc_id 为什么默认用内容哈希?带来什么幂等性?
- add_documents 为什么用 upsert + 分批?
- distance 怎么翻译成 score?cosine 和 l2 公式各是什么?
✋ 10 分钟动手
P=lib/crewai/src/crewai/rag
sed -n '66,113p' $P/core/base_client.py # BaseClient 协议
sed -n '56,95p' $P/chromadb/utils.py # doc_id 内容哈希
sed -n '167,189p' $P/chromadb/utils.py # distance→score
sed -n '309,358p' $P/chromadb/client.py # add_documents upsert + 分批
明日预告 · Day 39:有了 RAG 底座,上面就能搭"知识库"了。明天读
knowledge/:Knowledge 聚合多个知识源、BaseKnowledgeSource 怎么把 PDF/CSV/字符串切成 chunk、KnowledgeStorage 怎么把它们塞进今天的 ChromaDBClient。