Day 40 / 共 60 天 · 阶段6 记忆与知识

embedding 与存储后端:向量从哪来、存哪去

阶段6 收官。记忆和知识库都绕不开两件事:把文字变成向量(embedding)把向量存下来能搜(storage)。今天读这两个"可插拔点"。嵌入侧:rag/embeddings/ 用一个 provider 注册表 + 工厂 统一了 OpenAI/Cohere/Ollama 等十几种嵌入模型。存储侧:LanceDBStorage(669 行)的生产级细节——维度自动检测、维度不匹配主动报错、提交冲突指数退避重试、自动 compaction、scope 标量索引。这些是"能上生产"和"玩具"的分水岭。

📍 你在 60 天里的位置(阶段6 记忆与知识 · 收官)
D33 记忆总览 D34 unified_memory D35 recall/encoding D36 短/长/entity D37 scope 作用域 D38 RAG D39 knowledge D40 embedding/存储
💡 先用一个类比兜住今天 嵌入模型像不同厂牌的"翻译机"——都把文字翻成向量,但用词习惯(维度、语义空间)不同。CrewAI 用一张"翻译机名录"(provider 注册表),你报个牌子("openai")它就给你配好那台机器。存储像一个高并发的档案库:多人同时往里塞文件难免"撞车"(提交冲突),好的档案库会自动重试;文件多了会碎片化,它会定期整理(compaction);还会给常查的字段建索引加速。今天看这两样怎么做到"稳"。
L01

痛点:换模型、换库、还要扛并发

🤔 痛点三个真实需求:① 我不想用 OpenAI 嵌入(要花钱/数据出境),想换成本地 Ollama 或别的——框架得让我一行配置就换。② 记忆库跑久了、并发写多了,会不会写坏、会不会越来越慢?③ 我升级了 CrewAI,默认嵌入模型变了,旧库还能用吗?这三件事分别对应嵌入可插拔、存储的健壮性、维度兼容
💡 一句话本质 嵌入侧用"字符串 provider 名 → 注册表 → 动态 import → 实例化"把十几种嵌入模型统一成一个 build_embedder(spec)。存储侧用"维度自检 + 不匹配抛专用错 + 提交冲突指数退避重试 + 后台 compaction + 标量索引"把嵌入式向量库做到可上生产。可插拔靠工厂,健壮靠这一堆细节。
L02

嵌入 provider 工厂:一张注册表统一十几种

核心是一张 PROVIDER_PATHS 表 + build_embedderrag/embeddings/factory.py:90):

# rag/embeddings/factory.py:90(节选)
PROVIDER_PATHS = {
    "openai":  "crewai.rag.embeddings.providers.openai.openai_provider.OpenAIProvider",
    "cohere":  "crewai.rag.embeddings.providers.cohere.cohere_provider.CohereProvider",
    "ollama":  "crewai.rag.embeddings.providers.ollama.ollama_provider.OllamaProvider",
    "google":  "crewai.rag.embeddings.providers.google.generative_ai.GenerativeAiProvider",
    "huggingface": "...huggingface_provider.HuggingFaceProvider",
    "voyageai": "...voyageai_provider.VoyageAIProvider", ...    # 共 ~18 种
}

# rag/embeddings/factory.py:223(节选)
def build_embedder_from_dict(spec):
    provider_name = spec["provider"]
    if provider_name not in PROVIDER_PATHS:
        raise ValueError(f"Unknown provider: {provider_name}. Available: {list(PROVIDER_PATHS.keys())}")
    provider_path = PROVIDER_PATHS[provider_name]
    provider_class = import_and_validate_definition(provider_path)   # ★动态 import
    provider_config = spec.get("config", {})
    provider = provider_class(**provider_config)          # 实例化 provider
    return build_embedder_from_provider(provider)

# rag/embeddings/factory.py:113
def build_embedder_from_provider(provider):
    return provider.embedding_callable(                   # 造出真正的嵌入函数
        **provider.model_dump(exclude={"embedding_callable"}))
PROVIDER_PATHS 注册表字符串名 → provider 类的"点路径"。加新嵌入模型只要注册一行——又是 D38/D39 见过的"标签→实现"套路。
import_and_validate_definition动态 import:用到哪个 provider 才导入哪个模块。因为每个嵌入库(cohere/voyageai…)是可选依赖,不装也不影响其他——不能在顶部全 import。
provider_class(**config)provider 是个 BaseSettingsbase_embeddings_provider.py:14),能从 config dict 和环境变量读参数。
build_embedder_from_providerprovider 里存了 embedding_callable(真正的嵌入函数类),用 provider 的字段把它实例化出来——最终得到一个"输入 list[str]、输出向量"的可调用对象。
数据结构:嵌入 spec → provider → 嵌入函数 spec 字典 {provider,config} PROVIDER_PATHS 名→类路径 + 动态import Provider 实例 BaseSettings 嵌入函数 list[str]→向量 build_embedder("openai") → 一台"翻译机",memory/knowledge 都用它
图注:字符串 provider 名经注册表动态导入、实例化,最终产出统一的"文字→向量"嵌入函数。
L03

OpenAIProvider:一个具体 provider 长什么样

看最常用的 OpenAIProviderrag/embeddings/providers/openai/openai_provider.py:13):

# rag/embeddings/providers/openai/openai_provider.py:13(节选)
class OpenAIProvider(BaseEmbeddingsProvider[OpenAIEmbeddingFunction]):
    @model_validator(mode="before")
    @classmethod
    def _normalize_model_alias(cls, data):
        if isinstance(data, dict) and "model" in data and "model_name" not in data:
            data = data.copy(); data["model_name"] = data["model"]   # model→model_name 兼容
        return data

    embedding_callable: type[OpenAIEmbeddingFunction] = Field(default=OpenAIEmbeddingFunction)
    api_key: str | None = Field(default=None,
        validation_alias=AliasChoices("EMBEDDINGS_OPENAI_API_KEY", "OPENAI_API_KEY"))  # ★从环境变量读
    model_name: str = Field(default="text-embedding-3-large",                          # ★默认大模型
        validation_alias=AliasChoices("EMBEDDINGS_OPENAI_MODEL_NAME", "model_name"))
    dimensions: int | None = Field(default=None, ...)
泛型 [OpenAIEmbeddingFunction]provider 类型参数就是它要造的嵌入函数类;embedding_callable 字段默认指向它。
_normalize_model_alias容错:用户写 model 还是 model_name 都认——把常见的写法差异抹平,少一个"为什么没生效"的坑。
validation_alias 读环境变量api_key 能从 OPENAI_API_KEY 环境变量自动读(BaseSettings 的能力)。所以 Memory() 不传 key 也能跑——它去环境里找。
默认 text-embedding-3-large★默认嵌入是 3072 维的 large 模型。记住这个数——它是 D33/L05 那个 EmbeddingDimensionMismatchError 的主角。
大白话每个 provider 就是"某厂牌翻译机的说明书":告诉工厂"我要用哪个嵌入函数类、需要哪些参数(key/模型名/维度)、这些参数能从哪些环境变量读"。工厂照着说明书就能把机器造出来。你换厂牌 = 换 provider 名,别的代码全不动。
L04

LanceDBStorage:维度自动检测 + 懒建表

转到存储侧。LanceDBStorage 初始化时不知道向量多少维,靠自检(memory/storage/lancedb_storage.py:97):

# memory/storage/lancedb_storage.py:97(节选)
try:
    self._table = self._db.open_table(self._table_name)         # 表已存在
    self._vector_dim = self._infer_dim_from_table(self._table)  # ★从表结构读维度
    with store_lock(self._lock_name):
        self._ensure_scope_index()                              # 确保 scope 索引在
    self._compact_if_needed()                                   # 顺手整理碎片
except Exception:
    self._table = None
    self._vector_dim = vector_dim or 0     # 0 = 还不知道,等第一次 save 再定

# memory/storage/lancedb_storage.py:116
@staticmethod
def _infer_dim_from_table(table) -> int:
    schema = table.schema
    for field in schema:
        if field.name == "vector":
            return int(field.type.list_size)   # 向量列的固定长度就是维度
    return DEFAULT_VECTOR_DIM   # 3072
open_table 成功→读维度已有表:从它的 schema 里 vector 列的 list_size 读出维度——旧库多少维就跟着多少维,自动适配。
open 失败→延后建表没有表就把 _table=None_vector_dim=0等第一次 save 时用真实嵌入的长度来建表(_ensure_table)——让维度由嵌入模型的实际输出决定,而不是瞎猜。
DEFAULT_VECTOR_DIM=3072兜底默认(:27),对齐 OpenAI text-embedding-3-large。和 L03 的默认嵌入呼应。
启动即建索引 + compact开表时顺手确保 scope 索引存在、后台整理上次残留的碎片(L07 细讲)。
L05

save:维度校验,宁可报错也不悄悄污染

save 存前严格校验维度(memory/storage/lancedb_storage.py:289):

# memory/storage/lancedb_storage.py:289(节选)
def save(self, records: list[MemoryRecord]) -> None:
    if not records: return
    dim = None
    for r in records:
        if r.embedding and len(r.embedding) > 0:
            if dim is None: dim = len(r.embedding)
            elif len(r.embedding) != dim:
                raise EmbeddingDimensionMismatchError(dim, len(r.embedding))   # 批内不一致
    is_new_table = self._table is None
    if not is_new_table and dim and self._vector_dim and dim != self._vector_dim:
        raise EmbeddingDimensionMismatchError(self._vector_dim, dim)           # 与既有表不一致
    with store_lock(self._lock_name):
        self._ensure_table(vector_dim=dim)      # 首存时用真实维度建表
        rows = [self._record_to_row(rec) for rec in records]
        for row in rows:
            if row["vector"] is None or len(row["vector"]) != self._vector_dim:
                row["vector"] = [0.0] * self._vector_dim   # 无向量→零向量占位
        self._do_write("add", rows)
批内维度一致检查一批记录里如果向量维度不一(不可能来自同一嵌入模型),立刻抛 EmbeddingDimensionMismatchError
与既有表维度检查★如果新向量维度和表里已有的不一致(典型:换了嵌入模型),主动抛错——而不是把不同维度的向量混进去,那会让整个检索算错。
零向量占位没算出嵌入的记录(如嵌入 API 临时失败)用全 0 向量占位存下——保住内容,只是这条搜不出来。宁可存个搜不到的,也别丢内容。
store_lock 跨进程锁写操作全程持锁(crewai_core.lock_store),配合 D34 的单线程池,双重保证不并发写坏。
💡 设计取舍①:维度不匹配为什么用专用异常、且故意不继承 RuntimeError? EmbeddingDimensionMismatchErrormemory/storage/backend.py:11)的注释道破关键:它故意不继承 RuntimeError。因为 D34 我们看到,后台保存的容错逻辑会把 RuntimeError 当作"进程正在关闭"而静默丢弃这次保存。如果维度不匹配错误也是 RuntimeError,就会被这层容错悄悄吞掉——用户永远看不到"你该重建库了"这个可操作的提示,只会觉得"记忆莫名其妙没存进去"。所以它继承 ValueError:一个不会被静默吞、带完整迁移指引(reset 或钉旧模型)的错。异常的继承层次不是随便选的——它决定了这个错会被哪层 except 捕获、是被吞还是被抛。
L06

提交冲突:指数退避重试

LanceDB 用乐观并发,高并发写会撞"提交冲突",_do_write 自动重试(memory/storage/lancedb_storage.py:128):

# memory/storage/lancedb_storage.py:128(节选)
_MAX_RETRIES = 5
_RETRY_BASE_DELAY = 0.2   # 秒,每次翻倍:0.2+0.4+0.8+1.6+3.2 ≈ 6.2s

def _do_write(self, op, *args, **kwargs):
    delay = _RETRY_BASE_DELAY
    for attempt in range(_MAX_RETRIES + 1):
        try:
            return getattr(self._table, op)(*args, **kwargs)   # add/delete/update
        except OSError as e:
            if "Commit conflict" not in str(e) or attempt >= _MAX_RETRIES:
                raise                                          # 不是冲突/次数用尽→抛
            try:
                self._table = self._db.open_table(self._table_name)  # ★重新打开(拿最新版本)
            except Exception:
                pass
            time.sleep(delay)                                  # 退避等待
            delay *= 2                                         # ★指数退避
乐观并发LanceDB 不是写时上锁,而是"先写、提交时检查版本"。两个写基于同一版本、几乎同时提交,后提交的会"Commit conflict"。
只重试真冲突只有错误信息含 "Commit conflict" 才重试;别的 OSError(磁盘满等)直接抛——不盲目重试无意义的错(D08 的容错哲学)。
重试前 open_table★重试前重新打开表,拿到最新版本再写——否则基于旧版本重试还会冲突。
指数退避 delay*=2每次等待翻倍(0.2→0.4→0.8→1.6→3.2s)。给"版本快速推进"的高负载留出追赶时间,避免所有写挤在一起反复撞。
控制流:_do_write 的冲突重试 table.add/update Commit conflict? 否 → 成功返回 ✅ 是 → open_table + sleep delay*=2(0.2→3.2s),最多重试 5 次;超限则抛出
图注:写冲突时重开表拿最新版本、指数退避后重试;5 次仍冲突或非冲突错误则抛出。
L07

compaction、scope 索引与相似分换算

每次 save 产生一个碎片文件,攒多了拖慢查询,所以定期后台整理(memory/storage/lancedb_storage.py:314):

# memory/storage/lancedb_storage.py:314(save 尾部)
self._save_count += 1
if self._compact_every > 0 and self._save_count % self._compact_every == 0:  # 每 100 次
    self._compact_async()               # ★后台守护线程 table.optimize()

# lancedb_storage.py:183 —— scope 标量索引(加速前缀过滤)
def _ensure_scope_index(self):
    self._table.create_scalar_index("scope", index_type="BTREE", replace=False)

# lancedb_storage.py:403 —— 距离→相似分(和 D38 chromadb 同思路)
distance = row.get("_distance", 0.0)
score = 1.0 / (1.0 + float(distance)) if distance is not None else 1.0
每 100 次 save 后台 compactcompact_every=100:攒够 100 次写就在守护线程optimize() 合并碎片文件——不阻塞前台,查询性能保持稳定。启动时也顺手 compact 一次上次的残留。
BTREE scope 索引★给 scope 列建标量索引。D37 的隔离靠 scope LIKE '前缀%'——有索引就不用全表扫,前缀过滤(list_records/get_scope_info)快得多。
1/(1+distance)LanceDB 返回 _distance,用 1/(1+d) 转成 (0,1] 相似分(距离 0→分 1)。和 D38 chromadb 的 l2 换算完全一致——两套存储对上层给出统一语义的"分"。
_scan_rows 只取需要的列做 scope 统计时只 select 需要的列(:466),不读又大又重的 vector 列——省内存省 IO。
⚠️ 边界:scope 前缀删除的 ￿ 技巧 reset(scope_prefix) 删除某 scope 时用了 scope >= 'prefix' AND scope < 'prefix/￿'memory/storage/lancedb_storage.py:615)。￿ 是 Unicode 最大码点——这个范围技巧确保删掉 prefix/ 下的所有子路径(因为任何 prefix/xxx 都排在 prefix/￿ 之前),又不会误删名字以 prefix 开头的兄弟 scope。这是 D37 L08 提过的"前缀包含"陷阱的一个正解——用范围查询而非纯 LIKE 精确框定"某路径及其子树"。写向量库/KV 存储的前缀扫描时,这是个值得记住的手法。
L08

阶段6 收官:记忆与知识全景

👶 小白:我想换成本地嵌入、不用 OpenAI,怎么做?

👨‍🏫 老师:一行配置。Memory(embedder={"provider":"ollama","config":{"model":"nomic-embed-text"}}) 或给 crew 传 embedder=...。工厂(L02)会按 "ollama" 从注册表找到 OllamaProvider 动态导入、实例化。但切记(L05 的坑):不同嵌入模型维度不同,换了模型 = 旧向量库作废,得先 crewai reset-memories 重建,否则触发 EmbeddingDimensionMismatchError。知识库同理。

💡 阶段6 八天串起来 D33 总览(MemoryRecord/复合评分/StorageBackend 协议)→ D34 Memory 方法(异步写 + 读屏障 + shallow/deep)→ D35 两条 Flow(RecallFlow 置信度路由 + EncodingFlow 五步四组)→ D36 短/长/实体如何被一个 Memory 统一(新鲜度/重要性+合并/抽取元数据)→ D37 scope 作用域(路径前缀隔离 + slice 合并)→ D38 RAG 底座(BaseClient/doc_id 哈希/distance→score)→ D39 knowledge(源切块 + KnowledgeStorage 对接 RAG)→ D40 嵌入 provider 工厂 + LanceDB 生产级细节。一句话:记忆是"会评分、会合并、按作用域组织的动态经验",知识是"切块灌入的静态资料",二者共享嵌入与向量存储这两个可插拔底座。

🧠 今天你应该能回答

  • 嵌入 provider 工厂怎么统一十几种模型?为什么动态 import?
  • provider 怎么从环境变量读 key?默认嵌入模型和维度是什么?
  • LanceDB 怎么自动检测/确定向量维度?
  • 维度不匹配为什么用专用异常、且不继承 RuntimeError?
  • 提交冲突怎么重试?为什么重试前要 open_table?
  • compaction/scope 索引/范围删除各解决什么问题?

✋ 10 分钟动手

R=lib/crewai/src/crewai/rag; M=lib/crewai/src/crewai/memory
sed -n '90,124p'  $R/embeddings/factory.py            # provider 注册表 + 工厂
sed -n '13,40p'   $R/embeddings/providers/openai/openai_provider.py  # 具体 provider
sed -n '289,318p' $M/storage/lancedb_storage.py       # save + 维度校验
sed -n '128,153p' $M/storage/lancedb_storage.py       # 提交冲突重试
sed -n '11,42p'   $M/storage/backend.py               # 维度不匹配异常
python -c "
from crewai.rag.embeddings.factory import build_embedder
ef = build_embedder({'provider':'openai','config':{}})
print('默认嵌入函数:', type(ef).__name__)  # 需要 OPENAI_API_KEY
"
明日预告 · Day 41(阶段7 Flow):阶段6 里我们反复用到 CrewAI 自家的 Flow@start/@listen/@router)来写记忆流水线。从明天起进入阶段7 Flow 事件驱动,正式拆 flow/:Flow 是什么、和 Crew 什么关系、装饰器怎么把方法串成有向图。
← Day 39 knowledge Day 41 · Flow 总览 →