跳到主要内容

数据截至 (上游 commit 8c51b8dc5408)

存储分层、多租户与运维面

本章讲什么: chunk 的内容为什么不进数据库、租户之间怎么隔开、跑起来之后怎么看它在干什么。这一章偏工程,不偏算法。


1. 三层存储,各存各的

Postgres对象存储(S3 / 本地盘)Turbopuffer(可选)
documents 表:元数据、system_metadatachunk_ids原始上传文件(app 专属 bucket)定长编码向量 + 属性
multi_vector_embeddings:位串向量页面图片 / 文本内容 ({app_id}/{doc}/{n}.png)
完整多向量 multivector/{doc}/{n}.npy

为什么内容不进数据库

一页 150 DPI 的 PNG 有几百 KB。一万页就是几个 GB。塞进 Postgres 的 TEXT 列会:

  • 让每一次 SELECT * 都拖出巨量数据;
  • 让备份和复制变得昂贵;
  • 让 TOAST 表膨胀。

所以 multi_vector_embeddings.content 存的是一个 key,形如 {app_id}/{document_id}/{chunk_number}.png,生成逻辑在 _generate_storage_key(core/vector_store/multi_vector_store.py:400)。

检索时再按 key 下载。这也是为什么 output_format="url" 这个选项有价值——让前端自己去取,服务端连下载都省了。

怎么判断一个字符串是 key 还是内容

is_storage_key(core/vector_store/utils.py:30)。规则大意是:短、带斜杠、不像 base64。这个判据在存和取两端都要一致,所以抽成了共享函数。

取内容时的候选键阶梯

_retrieve_content_from_storage(core/vector_store/multi_vector_store.py:510)里有一段长长的候选键构造。它按顺序试这些形式:

顺序候选形式为什么会有这种形式
1derive_repaired_image_key(...)修正过扩展名的键
2原样 key正常情况
3{bucket}/{key}历史上 bucket 名被拼进 key 的写法
4{key}.txt遗留文本 chunk
5{bucket}/{key}.txt上面两种叠加
6换掉扩展名的 .txt / .txt.txt更早的遗留形式

最后用 dict.fromkeys 去重保序。这段代码是历史包袱的化石层——每一条候选都对应线上某个时期写进去的数据。想理解一个项目的年龄,看它的兼容分支比看 CHANGELOG 快。

拿到字节之后还要判类型:是图片就还原成 data URI(先试解码看是不是已经是 data URI,再按魔数嗅探 MIME,支持 PNG/JPEG/GIF/BMP/TIFF/WEBP,第 588-604 行);不是就当 UTF-8 文本,解不了再退回 base64。

元数据的双格式兼容

@staticmethod
def _parse_metadata(meta: Optional[str]) -> Dict[str, Any]:
"""Robustly parse metadata stored as JSON or Python dict string.

Some historical rows stored `str(dict)` rather than JSON; handle both.
"""

core/vector_store/multi_vector_store.py:97-115。先试 json.loads,失败再试 ast.literal_eval。又一处化石。


2. 多租户:一列贯穿到底

app_id 这一个字段串起了整条链:

JWT token
| verify_token 解出 app_id
v
AuthContext(user_id, app_id)
|
+--> SQL: WHERE app_id = :app_id (文档表)
+--> 存储: {app_id}/{doc_id}/{chunk}.png (对象存储路径前缀)
+--> 向量: self.ns(app_id) (Turbopuffer 命名空间)

三层用的是同一个值,所以租户边界在数据库、文件系统、向量库上是一致的。快路存储那行 self.ns = lambda app_id: self.tpuf.namespace(app_id)(core/vector_store/fast_multivector_store.py:321)把每个租户放进独立命名空间,隔离度最高。

没有 app_id 时(自托管、开发模式)回落到 owner_id,存储路径用常量 DEFAULT_APP_ID = "default"(multi_vector_store.py:36)。

认证与吊销

core/auth_utils.py 除了验签,还做应用有效性检查(ensure_app_is_active:88)。因为 JWT 一旦签发就无法撤回,所以额外维护了两个 Redis 键:

键前缀含义命中后行为
auth:app_revoked:{app_id}已吊销直接 401
auth:app_active:{app_id}有效,值是 token 版本号版本号一致就放行,省一次数据库查询

Redis 挂了会打警告并回落数据库(第 113-116 行),不阻断请求。bypass_auth_mode = true 时整个检查跳过(自托管默认开)。

用量与配额

云模式下有两处闸门(core/limits_utils.py:92check_and_increment_limits):

  1. 接口层预检:_verify_ingest_and_storage_limits(ingestion_service.py:463),用 verify_only=True 做干跑,超了直接拒绝上传。
  2. worker 层复检:解析出真实字符数后再干跑一次(ingestion_worker.py:712-725),因为估算的页数这时才准。

计量口径:estimate_pages_by_chars,4 字符/token、630 token/页,至少 1 页。

自托管模式(settings.MODE != "cloud")下这些检查全部短路,不产生开销。


3. 存储后端抽象

BaseStorage (core/storage/base_storage.py)
|
+-- LocalStorage 本地目录,bucket 参数被忽略
+-- S3Storage 真 S3,带并发上传信号量

选择在 core/services_init.py:64-79,由 [storage] provider 决定。

并发上传由 S3_UPLOAD_CONCURRENCY(默认 16)控制,并且被复用到多个地方:慢路的内容上传信号量(multi_vector_store.py:78-81)、快路的多向量上传信号量(fast_multivector_store.py:351-354)。一个配置管住所有出站并发。

快路还会尝试复用同一个存储实例给 chunk 内容和多向量用——只要 provider 和路径/bucket 相同(_init_vector_storage:392),省掉一份连接池。


4. 运维面:它跑起来之后能看到什么

进度

六步进度写进 system_metadata.progress:

{"current_step": 4, "total_steps": 6, "step_name": "Generating embeddings", "percentage": 67}

完成时被清成 None。前端轮询文档状态就能画条。

结构化进度日志

worker 里有一个独立的 progress_logger,打的是可 grep 的定长格式:

ingest start doc_id=... file=... colpali=True
ingest download doc_id=... size_mb=12.30 time_s=1.42
ingest chunks doc_id=... pages=87 text_chunks=0 image_chunks=87
ingest batching doc_id=... total_chunks=87 batch_size=16 batches=6
ingest batch doc_id=... 3/6 size=16 embed_s=8.21 store_s=2.10
ingest done doc_id=... status=completed total_s=124.55

分项计时

phase_times 字典贯穿整个 worker,ColPali 部分细到:排序、预处理、模型前向、张量转换,并且图片和文本分开统计(ingestion_worker.py:1227-1240)。想知道"到底是模型慢还是 I/O 慢",这份数据直接给答案。

检索侧对应的是 PerformanceTracker,从 API 层注入,支持父子阶段嵌套(add_suboperation)。

存储指标

build_store_metrics(core/vector_store/utils.py:82)统一了三个存储实现的指标口径:

指标含义
chunk_payload_upload_s / _objects / _bytes内容上传耗时、对象数、字节数
multivector_upload_s / _objects / _bytes多向量上传(仅快路)
vector_store_write_s / _rows向量库写入
chunk_payload_backend各层实际用的后端名

这些数字会回写到文档行上做用量统计(record_document_storage_deltas)。

遥测

TelemetryService(core/services/telemetry.py)基于 OpenTelemetry,配合 LogUploaderHeartbeat(在 core/app_factory.py:120-152 启动)。心跳打到 https://logs.morphik.ai/api/heartbeat,间隔 4 小时。可以通过 [telemetry] 段关掉。


5. 启动顺序

core/app_factory.pylifespan 按固定顺序初始化,失败策略不一样:

步骤失败了怎么办
数据库 initialize抛异常,启动失败
向量库 initialize记错误,继续
v2 chunk store initialize记错误,继续
ColPali 向量库 initialize记错误,继续
Redis 连接池抛异常,启动失败

只有数据库和 Redis 是硬依赖——没有它们连排队都做不了。向量库初始化失败会推迟到第一次真正使用时报错。


6. 代码地图

主题文件路径符号名
存储抽象与两个实现core/storage/base_storage.pylocal_storage.pys3_storage.pyBaseStorageLocalStorageS3Storage
存储键生成与判定core/vector_store/multi_vector_store.pycore/vector_store/utils.py_generate_storage_keyis_storage_keynormalize_storage_key
内容外置与取回(含候选阶梯)core/vector_store/multi_vector_store.py_store_content_externally_retrieve_content_from_storage
历史元数据兼容core/vector_store/multi_vector_store.py_parse_metadata
存储指标统一口径core/vector_store/utils.pybuild_store_metricsextract_storage_bytes
认证与应用吊销core/auth_utils.pyverify_tokenensure_app_is_activemark_app_revoked
配额检查core/limits_utils.pycheck_and_increment_limitsestimate_pages_by_chars
应用启动顺序core/app_factory.pylifespan
配置装载core/config.pymorphik.tomlSettingsget_settings
遥测与心跳core/services/telemetry.pycore/services/heartbeat.pyTelemetryServiceHeartbeat
进度上报core/workers/ingestion_worker.pyupdate_document_progressprogress_logger