embedding 与存储后端:向量从哪来、存哪去
阶段6 收官。记忆和知识库都绕不开两件事:把文字变成向量(embedding) 和 把向量存下来能搜(storage)。今天读这两个"可插拔点"。嵌入侧:rag/embeddings/ 用一个 provider 注册表 + 工厂 统一了 OpenAI/Cohere/Ollama 等十几种嵌入模型。存储侧:LanceDBStorage(669 行)的生产级细节——维度自动检测、维度不匹配主动报错、提交冲突指数退避重试、自动 compaction、scope 标量索引。这些是"能上生产"和"玩具"的分水岭。
"openai")它就给你配好那台机器。存储像一个高并发的档案库:多人同时往里塞文件难免"撞车"(提交冲突),好的档案库会自动重试;文件多了会碎片化,它会定期整理(compaction);还会给常查的字段建索引加速。今天看这两样怎么做到"稳"。痛点:换模型、换库、还要扛并发
build_embedder(spec)。存储侧用"维度自检 + 不匹配抛专用错 + 提交冲突指数退避重试 + 后台 compaction + 标量索引"把嵌入式向量库做到可上生产。可插拔靠工厂,健壮靠这一堆细节。嵌入 provider 工厂:一张注册表统一十几种
核心是一张 PROVIDER_PATHS 表 + build_embedder(rag/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 是个 BaseSettings(base_embeddings_provider.py:14),能从 config dict 和环境变量读参数。build_embedder_from_providerprovider 里存了 embedding_callable(真正的嵌入函数类),用 provider 的字段把它实例化出来——最终得到一个"输入 list[str]、输出向量"的可调用对象。OpenAIProvider:一个具体 provider 长什么样
看最常用的 OpenAIProvider(rag/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 的主角。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 细讲)。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 的单线程池,双重保证不并发写坏。EmbeddingDimensionMismatchError(memory/storage/backend.py:11)的注释道破关键:它故意不继承 RuntimeError。因为 D34 我们看到,后台保存的容错逻辑会把 RuntimeError 当作"进程正在关闭"而静默丢弃这次保存。如果维度不匹配错误也是 RuntimeError,就会被这层容错悄悄吞掉——用户永远看不到"你该重建库了"这个可操作的提示,只会觉得"记忆莫名其妙没存进去"。所以它继承 ValueError:一个不会被静默吞、带完整迁移指引(reset 或钉旧模型)的错。异常的继承层次不是随便选的——它决定了这个错会被哪层 except 捕获、是被吞还是被抛。提交冲突:指数退避重试
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)。给"版本快速推进"的高负载留出追赶时间,避免所有写挤在一起反复撞。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 后台 compact★compact_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。 技巧
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 存储的前缀扫描时,这是个值得记住的手法。阶段6 收官:记忆与知识全景
👶 小白:我想换成本地嵌入、不用 OpenAI,怎么做?
👨🏫 老师:一行配置。Memory(embedder={"provider":"ollama","config":{"model":"nomic-embed-text"}}) 或给 crew 传 embedder=...。工厂(L02)会按 "ollama" 从注册表找到 OllamaProvider 动态导入、实例化。但切记(L05 的坑):不同嵌入模型维度不同,换了模型 = 旧向量库作废,得先 crewai reset-memories 重建,否则触发 EmbeddingDimensionMismatchError。知识库同理。
🧠 今天你应该能回答
- 嵌入 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
"
Flow(@start/@listen/@router)来写记忆流水线。从明天起进入阶段7 Flow 事件驱动,正式拆 flow/:Flow 是什么、和 Crew 什么关系、装饰器怎么把方法串成有向图。