架构与边界
一个 Java 控制面、一个 Python 模型服务、一套 PostgreSQL。这一篇讲清楚职责怎么划、数据怎么存、版本与并发怎么处理、依赖挂掉时会发生什么。
一、先说结论
- Java 拥有身份、授权、入库与检索;Python 只做模型计算,不持有身份、不碰数据库。
- 两者只通过带版本的 OpenAPI 契约通信,双侧都有漂移检测。
- 单体优先:控制面是一个 Maven 模块,包边界由架构测试强制。
- 版本切换是单事务内的指针翻转,读者只看到完整的旧版本或新版本。
- 并发由数据库裁决:同一文档的并发入库在写事务内加行锁并重新判断。
二、信任边界
| 边界 | 通过的内容 | 规则 |
|---|---|---|
| 客户端 → 控制面 | 查询、Bearer token | 身份属性只来自验签后的 token,绝不来自请求体 |
| 控制面 → PostgreSQL | 带授权谓词的 SQL | 唯一筛选候选行的地方 |
| 控制面 → 模型服务 | 查询文本、已授权的 chunk 文本 | 模型服务收不到租户 id、身份属性或未授权文本 |
| 控制面 → 生成模型(M3) | 由已授权 chunk 构建的 prompt | 文档内容视为数据;模型没有工具调用能力 |
模型服务的请求体设置了 extra="forbid":误传的身份字段(例如 tenant_id)会直接 422,而不是悄悄越过边界。
三、控制面的包与约束
io.groundedaccess
├── identity 验签 JWT → Principal
├── authorization PolicyCompiler → 单一参数化 SQL 谓词
├── corpus 规范化 · 切分 · 内容哈希版本管理 · embedding 模型护栏
├── retrieval AuthorizedChunkQuery:chunk 表唯一的读取方
├── modelclient 模型服务客户端:分批、超时、固定模型
└── api REST controller、作用域、问题响应由测试强制的依赖规则:
- 只有
AuthorizedChunkQuery的源码可以出现from/join chunk; authorization不依赖 web、retrieval 与 JDBC;- controller 不接触数据库;
- 只有
modelclient知道模型服务的线上格式——这条规则曾拦下一次「护栏组件直接读模型服务配置」的提交,改为由EmbeddingClient暴露模型名。
四、数据模型
tenant(id, name)
document(id, tenant_id, external_key, status, active_version_id, …) -- external_key 租户内唯一
document_version(id, document_id, tenant_id, version_no, content_sha256, title, chunker_version, embedding_model, …)
chunk(id, tenant_id, document_id, version_id, ordinal, section_path, char_start, char_end, content, content_tsv, embedding, token_count)几个设计点:
document.active_version_id外键可延迟校验,允许同一事务内先插入新版本的所有 chunk,最后翻转指针;chunk.tenant_id冗余存储,租户条件无需 join 即可先行过滤,将来也可据此分区;content_tsv是english配置的生成列 + GIN 索引,embedding是vector(384);- 访问标签列(密级、部门、项目、region、有效期)在 M2 随决策表一起加入——不提前放置没有实现的列。
五、入库流程与并发
ingest(document)
规范化 + SHA-256
预检查:内容未变 → 跳过(仅作为优化)
切分 + 调模型服务算向量 ← 网络调用在事务之外
事务:
upsert 文档行,SELECT ... FOR NO KEY UPDATE 加锁
单独一条语句读活跃版本 ← 不能与加锁写在同一个 JOIN 里
哈希仍然相同 → UNCHANGED
写入新版本与 chunk → 翻转 active_version_id倒数第二步有个容易踩的坑:最初的实现用一条 LEFT JOIN 同时加锁并读取活跃版本。PostgreSQL 在等待锁之后会重新求值被锁的那一行,但与它 join 的行仍来自旧快照,于是等锁的事务读到过期的版本号,写出重复的 version_no。加锁与读取必须是两条语句。
测试用一个带屏障的 embedding 客户端制造真实竞态:相同内容只写一个版本,不同内容按顺序写两个,已存在文档的并发更新按序追加。去掉行锁后这些测试会失败。
异步任务
上面的 ingest 由任务 worker 逐篇调用。POST /ingestion-jobs 只把文档存进 ingestion_job_document 并返回 202,不解析、不调模型:
claim
租约在最后一次尝试中过期的任务 → failed(WORKER_LOST)
UPDATE … WHERE id = (SELECT 最早的可运行任务 … FOR UPDATE SKIP LOCKED)
可运行 = queued 且 run_after 已到
| running 但租约过期、仍有剩余次数
run
从上次记录的位置开始逐篇 ingest(每篇各自一个事务)
每篇之后:processed + 1、累加计数、续租
失败
模型服务不可达 / 5xx / 数据库暂时性错误,且有剩余次数 → queued,退避 5s × 2^(n-1)
其他错误 → failed(error_code)
worker 的每次写入:WHERE status = 'running' AND attempts = <自己认领的编号>几个取舍:
- 断点续跑而不是从头重来:重试不会重新 embedding 已写入的文档。投递语义是至少一次——worker 恰好在写完文档、记录进度之前死掉,这篇会再入库一次并计为
unchanged,由内容哈希幂等兜底。 - 租约而不是长事务:embedding 期间不持有数据库事务;丢了租约的 worker 因 attempt 编号不匹配而无法再改任务。
- 不对提交去重:按 manifest 哈希去重会挡住合理的重跑;重复提交只会得到一个全部
unchanged的任务。
停用、删除与清理
PATCH /api/v1/documents/{key} {"status": "active" | "disabled"}
DELETE /api/v1/documents/{key} → 204;未知、已删除或其他租户的 key 一律 404对下一次查询生效:状态是文档行上的一次 update,检索本来就只 join「active 文档的活跃版本」,查询语句一行不改。状态变更与入库用同一把行锁,同一 key 的并发操作串行,后提交者生效。
停用不被新内容撤销:例行重新同步写入新版本,但文档保持隐藏,直到管理员重新启用。
删除是墓碑:文档行保留 key,清理任务(每分钟)删掉被替换的旧版本和已删除文档的全部版本,chunk 随外键级联删除。清理与入库用同一把锁并
SKIP LOCKED跳过正被入库的文档;检索正确性从不依赖清理是否已运行。版本号永不复用:文档行上维护单调递增的
last_version_no。最初用max(version_no)计算下一个版本号,清理删掉旧版本后,被删除又重新入库的 key 会从 1 重新计数,key + version就指向了另一份内容。这个问题是在本地手工删除、重新入库后跑评测时暴露的:评测标注按「文档 + 版本」定位证据,召回率随之下降。SKIP LOCKED本身有专门的测试:一个事务锁住最早的任务时,认领应立即拿到下一个而不是等待——去掉SKIP LOCKED该测试失败。正确性(每个任务只被认领一次)靠行锁本身即可保证,SKIP LOCKED买的是多个 worker 互不阻塞。
六、模型服务
| 端点 | 说明 |
|---|---|
POST /v1/embed | input_type 区分 query 与 passage,可选 model 字段防漂移 |
GET /v1/models | 模型名、Hugging Face 快照 revision、维度、query 前缀、许可证 |
GET /healthz | 存活检查 |
实现上的几个选择:用 fastembed(ONNX Runtime)而不是 sentence-transformers,避免引入 PyTorch;权重在构建阶段下载进镜像,运行时 HF_HUB_OFFLINE=1;每个响应都带 revision,供评测与文档版本记录。
客户端固定使用 HTTP/1.1——JDK HttpClient 对明文 HTTP 默认发送 h2c 升级请求,而 uvicorn 会丢弃升级请求的 body,表现为「请求体为空」的诡异 422。这个坑由一个基于真实 HTTP server 的测试守住。
七、依赖故障时的行为
| 故障 | 行为 | 标记 |
|---|---|---|
| 查询时模型服务不可用 | dense 策略失败(503);hybrid 降级为仅 sparse | degraded: dense_unavailable |
| 重排超时(M3) | 回退到融合后的顺序 | degraded: rerank_unavailable |
| 生成模型不可用(M3) | 只返回证据 | degraded: generation_unavailable |
| 审计写入失败(M2) | 查询失败关闭(503) | 错误指标 |
| 入库时模型服务不可用 | 任务重试,耗尽后置为 failed | 任务 error_code |
评测会把任何降级的运行判为无效,避免降级结果混进报告。
八、可观测性(计划)
每次查询一条 trace:auth.verify → policy.compile → retrieval.sparse / retrieval.dense → fusion → rerank → context.build → generation → citation.validate。
允许记录的 span 属性:pipeline 配置哈希、policy 版本、各通道授权后的候选数、模型名与 revision、token 数、降级原因、错误码。默认不记录:查询文本、chunk 内容、prompt、模型输出、embedding、文档标题。审计与遥测分开:遥测可采样,审计不可。
小结
架构上没有炫技的部分:单体、一套数据库、一个模型服务。真正花力气的是边界——谁能读 chunk、谁知道模型格式、谁能拿到身份,这些都写成了会让构建失败的测试。最后看路线图与非目标。