跳到主要内容

数据截至 (上游 commit 8c51b8dc5408)

两套多向量存储:同一个接口,两条完全不同的路

本章讲什么: MaxSim 精确算太贵,存储层怎么变着法子省。项目给了两个实现,取舍完全相反。


1. 它要解决的小问题

第 02 章留下的账:一页有上千个 patch 向量,一次查询要对每一页算一个 n_query_token × n_patch 的相似度矩阵。一万页文档就是一万次矩阵乘。没法直接上生产。

两个实现分别下注在不同的地方:

慢路 MultiVectorStore快路 FastMultiVectorStore
省法把向量压成 1 bit/维,用位运算算先用一个定长向量粗筛,只对候选做精确算
精度近似(二值化有损),一步到位粗筛近似 + 重排精确
依赖只要 PostgresTurbopuffer(外部托管服务)+ 对象存储
配置[multivector_store] provider = "postgres"(默认)provider = "morphik" + TURBOPUFFER_API_KEY
代码core/vector_store/multi_vector_store.py:39core/vector_store/fast_multivector_store.py:305

两者实现同一个抽象基类 BaseVectorStore,core/services_init.py:141-190 按配置二选一,上层完全无感。


2. 慢路:二值量化 + 自建 SQL 函数

思路

128 维 float32 是 512 字节。但如果只关心"每一维是正是负",128 维就只要 16 字节——压缩 32 倍。两个二值向量的相似度可以用异或加popcount(数 1 的个数)算出来,这在 CPU 上极快。

项目自己在注释里给了出处(core/utils/fast_ops.py:7-10):Qdrant 的二值量化文章、SimSIMD 库。

表结构

CREATE TABLE IF NOT EXISTS multi_vector_embeddings (
id BIGSERIAL PRIMARY KEY,
document_id TEXT NOT NULL,
chunk_number INTEGER NOT NULL,
content TEXT NOT NULL, -- 这里存的是对象存储的 key,不是内容本身
chunk_metadata TEXT,
embeddings BIT(128)[] -- 一页 = 一个位串数组
)

core/vector_store/multi_vector_store.py:242-250BIT(128)[] 是整个设计的核心:一行存一页,数组里每个元素是一个 patch 的 128 位二值向量。

那个 SQL 函数

initialize() 会往数据库里装一个自定义函数(core/vector_store/multi_vector_store.py:287-311):

CREATE OR REPLACE FUNCTION public.max_sim(document bit[], query bit[])
RETURNS double precision
LANGUAGE SQL IMMUTABLE PARALLEL SAFE
AS $$
WITH queries AS (SELECT row_number() OVER () AS query_number, *
FROM (SELECT unnest(query) AS query) AS foo),
documents AS (SELECT unnest(document) AS document),
similarities AS (
SELECT query_number,
1.0 - (bit_count(document # query)::float /
greatest(bit_length(query), 1)::float) AS similarity
FROM queries CROSS JOIN documents),
max_similarities AS (
SELECT MAX(similarity) AS max_similarity FROM similarities GROUP BY query_number)
SELECT COALESCE(SUM(max_similarity), 0.0) FROM max_similarities
$$

逐行对应第 02 章那个 Python 演示:

SQL 片段干什么
queries CROSS JOIN documents查询 token × 页面 patch 的全组合(就是那个矩阵)
bit_count(document # query)# 是按位异或,bit_count 数 1 的个数 = 汉明距离
1.0 - 距离 / 位长汉明距离归一化成 [0,1] 的相似度
MAX(...) GROUP BY query_number每个查询 token 取最大值
SUM(max_similarity)求和得到页面得分

标了 IMMUTABLE PARALLEL SAFE,允许 Postgres 并行执行。

检索长什么样

query = (
"SELECT id, document_id, chunk_number, content, chunk_metadata, "
f"max_sim(embeddings, {array_literal}) AS similarity "
"FROM multi_vector_embeddings"
)
# ... WHERE document_id IN (...) ORDER BY similarity DESC LIMIT k

core/vector_store/multi_vector_store.py:746-760

注意这里没有索引可用。 max_sim 是自定义函数,ORDER BY max_sim(...) 只能全表扫描——对 WHERE document_id IN (...) 圈定的范围内每一行都算一次。这就是它叫"慢路"的原因:精度可以接受,但延迟随文档数线性增长。

还有一个细节:查询向量是拼进 SQL 字符串的(第 743 行 array_literal),不是参数化的:

array_literal = "ARRAY[" + ",".join(f"B'{s}'" for s in bit_strings) + "]::bit(128)[]"

代码注释标注 "internal usage only"。内容来自模型输出的位串(只含 0/1),不是用户输入,所以注入面有限,但这是个需要留意的写法。

量化本身:Rust 加速,有 Python 兜底

_binary_quantize(core/vector_store/multi_vector_store.py:329)有一个漂亮的分支:

# Bit(bytes) infers the bit length as 8 * len(bytes), so the packed path
# is only valid when the dimension is byte-aligned (always true for ColPali's 128).
if np.shape(embeddings)[-1] % 8 == 0:
return [Bit(packed) for packed in binary_quantize_packed(embeddings)]
binary_lists = binary_quantize(embeddings)
return [Bit(bits) for bits in binary_lists]

维度能被 8 整除时走打包成字节的路径(128 维 → 16 字节),否则退回布尔列表。128 永远满足,所以实际总是走快路径。

底层实现在 core/utils/fast_ops.py:有 morphik_rust 扩展就用 Rust(带 rayon 并行、手工展开成 SIMD 友好的形式,见 morphik_rust/src/binary_ops.rs:81),没有就用纯 Python 兜底。每个 Rust 函数都配了一份行为等价的 Python 实现,连正则匹配的字符集都特意对齐(fast_ops.py:22-26)。


3. 快路:先粗后精

思路

慢路的问题是"每一页都要精算"。快路的办法是加一层粗筛:

第一层(粗):把一页的上千个向量压成 ONE 个定长向量
→ 可以建普通 ANN 索引 → 毫秒级取回 top-75 候选

第二层(精):只对这 75 页下载完整多向量,做真正的 MaxSim
→ 排出最终 top-k

粗筛用的是一个叫 fixed-dimensional-encoding(定长编码) 的库。它的作用是:把一个不定长的向量集合编码成一个定长向量,使得两个定长向量的内积近似原来两个集合之间的 MaxSim/Chamfer 相似度。这样多向量检索就退化成了普通的单向量 ANN 检索。

诚实说明: 这个库以本地依赖形式引入(pyproject.toml[tool.uv.sources] 指向 fde),但克隆里 fde/ 目录只有 pyproject.toml,没有实现源码。所以本文只能描述它的调用方式和配置,算法内部无法从这个克隆核实。

编码参数

core/vector_store/fast_multivector_store.py:335-341:

self.fde_config = fde.FixedDimensionalEncodingConfig(
dimension=128,
num_repetitions=20,
num_simhash_projections=5,
projection_dimension=16,
projection_type="AMS_SKETCH",
)
参数含义(据参数名推断,inferred)
dimension128输入向量维度,和 ColPali 输出对齐
num_simhash_projections5SimHash 分桶位数,即 2^5 = 32 个空间分区
num_repetitions20重复 20 组独立分区以降方差
projection_dimension16每个桶内降维到 16 维
projection_typeAMS_SKETCH用的投影草图类型

写入路径

store_embeddings(core/vector_store/fast_multivector_store.py:439)一次写三份东西:

一页的数据
|
+-- 定长编码向量 --> Turbopuffer(可 ANN 检索)
+-- 完整多向量 .npy --> 对象存储 (multivector/{doc_id}/{chunk_no}.npy)
+-- 页面图片/文本 --> 对象存储 ({app_id}/{doc_id}/{chunk_no}.png)

Turbopuffer 那行同时带上 content(图片的存储 key)和 multivector(.npy 的 bucket + key)两个属性,检索时一次取回不用再查库。

.npy 存成 float32 而不是 float64,注释解释了原因(fast_multivector_store.py:723-724):ColPali 输出 bfloat16,和 float32 共享指数范围,存 float32 既不损动态范围又省一半。

有一个刻意的"不做":上传失败必须炸,不能留空洞(fast_multivector_store.py:472-480):

if self.chunk_storage is not None and any(key is None for key in storage_keys):
raise RuntimeError(
f"{missing}/{len(storage_keys)} chunk-content uploads returned no storage key; "
"aborting batch to avoid storing null content"
)

宁可让整批失败重试,也不能往 Turbopuffer 里写一条内容指向 None 的记录。

检索路径

query_similar(core/vector_store/fast_multivector_store.py:523)六步,每步都打了埋点计时:

[1] 查询多向量 --> 定长编码 fde.generate_query_encoding(放到线程池,CPU 密集)
|
[2] Turbopuffer ANN 检索 top_k = min(10*k, 75),按 document_id 过滤
|
[3] 并发下载候选页的 .npy load_multivector_from_storage(带磁盘 LRU 缓存)
|
[4] 精确 MaxSim 重排 processor.score_multi_vector(放到线程池)
| torch.topk 取前 k
[5] 并发下载 top-k 的页面内容 _retrieve_content_from_storage
|
[6] 组装 DocumentChunk 返回

第 2 步那个 top_k=min(10 * k, 75) 是整条路的调参旋钮:候选取多了重排慢,取少了召回掉。上限 75 是硬编码。

第 1、4 步都用 asyncio.to_thread 挪出事件循环。注释直说是因为 "CPU-bound FDE" 和 "CPU-bound torch"——这两处不挪走会阻塞整个 FastAPI 进程。

一个坏页不能拖垮整个查询

第 3 步之后有一段防御,注释写得非常具体(fast_multivector_store.py:570-583):

# A single missing/corrupt/lagging .npy (delete-reingest race, S3
# eventual consistency) must degrade one result, not 500 the whole
# query for every user of the app. Drop failed candidates, score the rest.

下载失败的候选直接丢掉,剩下的照样打分;全部失败才返回空。

同一处还留了个踩坑记录:

# NB: turbopuffer Row is a pydantic model — it supports row["id"]
# (subscript) but NOT row.get(...). Using .get() here would raise
# AttributeError and 500 the whole query on exactly this degraded
# path, defeating the purpose of the fix.

磁盘缓存:只在读路径回填

.npy 文件反复被下载,所以有一个 LRU 磁盘缓存 FileCacheManager(fast_multivector_store.py:79),按字节数限额淘汰。

有意思的是摄取时故意不写缓存(fast_multivector_store.py:750-755):

刚摄取的页绝大多数在被 LRU 淘汰之前根本不会被查询,这次写入只会挤掉热数据,还给每一页的摄取加一次同步 I/O。读路径首次未命中时自然会回填,检索延迟不受影响。

这是一条很典型的"看起来该做但实测不该做"的优化。


4. 迁移期的第三个实现

DualMultiVectorStore(core/vector_store/dual_multivector_store.py:24)是个包装器,行为很直白:

操作行为
store_embeddings并发写两个库
query_similar只读慢路
get_chunks_by_id只读慢路
delete_chunks_by_document_id两个库都删

失败处理不对称(第 89-100 行):快路失败只打错误日志,慢路失败直接抛。因为检索走慢路,慢路才是真相来源。

用法是"先双写一段时间攒够数据,再把读切过去"。由 ENABLE_DUAL_MULTIVECTOR_INGESTION 控制。


5. 还有一条完全不同的路:传统 pgvector

不开 ColPali 时用 PGVectorStore(core/vector_store/pgvector_store.py)。它就是标准做法:

  • 一块文本一个向量,<=> 余弦距离排序。
  • 分数换算:score = 1.0 - distance / 2.0(第 499 行),把 [0,2] 的余弦距离映射到 [0,1]。
  • ivfflat.probes 调召回(第 459-461 行,用 set_config 而不是 SET,因为 SET 不支持参数化)。
  • 一个小优化:查询时 defer(VectorEmbedding.embedding) 不取向量列。注释算过账——每行约 30KB 传输加每维一次 float(),而调用方转手就丢弃(第 465-467 行)。

此外还有一个更新的 ChunkV2Store(core/vector_store/chunk_v2_store.py:56),表 chunk_v2app_idfolder_pathpage_numberfilename 全部平铺成列并建 ivfflat 索引,而不是塞在 JSON 里。这是新一代 v2 通道用的存储,和 v1 并存。


6. 怎么选

场景选哪个为什么
自托管、文档量千级以内慢路(默认)零外部依赖,延迟可接受
文档量上万、要求亚秒快路全表扫描撑不住
已有慢路数据要迁移Dual,双写一段时间再切读不停机
纯文本、无视觉需求关掉 ColPali,用 pgvector省一个 3B 模型的资源

7. 代码地图

主题文件路径符号名
慢路存储core/vector_store/multi_vector_store.pyMultiVectorStore
建表与 max_sim 函数安装core/vector_store/multi_vector_store.pyinitialize
二值量化(含打包分支)core/vector_store/multi_vector_store.py_binary_quantize
慢路检索 SQL 拼装core/vector_store/multi_vector_store.pyquery_similar
批量插入core/vector_store/multi_vector_store.py_bulk_insert_rows
快路存储core/vector_store/fast_multivector_store.pyFastMultiVectorStore
定长编码配置core/vector_store/fast_multivector_store.pyfde_config
快路两阶段检索core/vector_store/fast_multivector_store.pyquery_similar
多向量落盘 / 载入core/vector_store/fast_multivector_store.py_save_multivector_to_storage_with_cache_timeload_multivector_from_storage
磁盘 LRU 缓存core/vector_store/fast_multivector_store.pyFileCacheManager
双写迁移包装core/vector_store/dual_multivector_store.pyDualMultiVectorStore
传统单向量存储core/vector_store/pgvector_store.pyPGVectorStore
v2 平铺列存储core/vector_store/chunk_v2_store.pyChunkV2Store
Rust / Python 双实现工具core/utils/fast_ops.pymorphik_rust/src/binary_ops.rsbinary_quantize_packedhamming_distance