① 只读对账 tools/reconcile_knowledge_vectors.py
按集合列出:孤儿向量 / 死向量 / 缺向量 / 重复正文 / 低信息量碎片 / 纯标题。
关键口径:非数字 id(FAQ-0013 这类语义 id)是灌库脚本有意写进 Milvus 的,
单独归类、不建议删;向量数取自 query 实际行数,不用 get_collection_stats
(后者含已软删未 compaction 的行)。
② 导入侧幂等:同 source_file + 集合重传 = 覆盖上一版
app/service/knowledge_ingest_service.py 新增 _supersede_previous_version:
把上一版 active 行置为 expired,并逐行投 knowledge.vector_delete_requested
(与本次入库同事务)。写入侧只认 active 而检索侧不看 status,旧向量不清掉
会继续参与排序、和同题活块抢答。
顺带修掉一个真 bug:改为先判 chunks 非空再下线 —— 否则传一份解析出 0 块的
文档会把上一版下架、新版一行没写,这份文档在检索侧凭空消失。
③ 清理入口:POST /api/v1/knowledge/{knowledge_id}/vector-cleanups
给历史上"被别的途径置为 expired、从未投过删除事件"的行补投向量清理。
DELETE 对已过期行返回 404 的口径保持不变(重复删除静默成功会让调用方
分不清"这次真下线了"和"早就过期了"),因此新开一个语义明确的端点:
不存在 404 / 仍是 active 422(请改用 DELETE)/ 已 expired 200 并回传事件名。
配套 tools/purge_expired_knowledge_vectors.py(默认 dry-run)批量驱动该端点。
文档:docs/演示用/知识库向量对账与清理-2026-09-15.md(含真机验证输出),
并对 docs/演示用/知识库问答诊断-2026-09-14.md 做两处更正 —— 实测孤儿向量 0 条、
那 175 行历史副本从来没有向量(不参与排序),当时的差额来自 get_collection_stats
把已软删行算进去。
新发现(未修,需业务拍板):661 条向量里 451 条正文不到 40 字,是灌库时把
markdown 表格/标题切碎产生的碎片。「风险评估问卷怎么评分」实测前 4 名是 4 条
一模一样的 19 字碎片(gap 0.0024),真正 2828 字的答案排第 5 → 客服必然转人工。
属灌库切分缺陷,补内容救不了,也不应靠放宽 MIN_GAP 解决。
验证:pytest tests/unit tests/contract → 1500 passed, 2 skipped, 0 failed;
mypy app → 3 个错全在组员文件中(与本次改动无关);ruff 本次改动文件 0 错。
真机端到端:重传 → 旧行 expired + 删除事件 published + 旧向量已从 Milvus 删除;
两个问句回归仍正常回答(r1到r5 gap 0.0766;申购确认 0.8453)。
348 lines
18 KiB
Python
348 lines
18 KiB
Python
"""知识入库服务:文件落存储 → 元数据落库(切块逐行)→ 投递向量同步事件。
|
||
|
||
## 分层边界
|
||
|
||
Service 层(MVC+S):使用**传入的** `AsyncSession`,不自建 Session、不直连 MySQL/Redis/Milvus,
|
||
也不调用 embedding —— Milvus 的写入与向量化全部由 `knowledge.vector_sync_requested` 事件的
|
||
Outbox handler(`app/worker/knowledge_vector_worker.py`)完成。这里只把"这件事需要做"写进
|
||
`domain_event_outbox`,与知识行处于**同一个事务**(调用方提交才一起生效),因此不存在
|
||
"知识入库了但向量同步事件丢了"的中间态。
|
||
|
||
构造器收 `embedder` / `endpoint_resolver` 不是摆设:它们是**后续对账自愈**要用的依赖
|
||
(裁定 3 —— 自愈补投的事件最终要靠它们落地向量),由调用方在组合根一次性注入,
|
||
避免 Service 以后为了自愈去偷偷 new 一个模型网关。
|
||
|
||
## 为什么 `ingest` 返回 `list[int]` 而不是单个 `int`
|
||
|
||
一份文档会被切分成**多个 chunk**(`DocumentParser`:512 字符/块、overlap 64,且 overlap 会在
|
||
相邻块间重复出现,见该模块 docstring),每个 chunk 是 `fin_knowledge_meta` 里**独立的一行**
|
||
(独立 `knowledge_id`、独立向量)。返回单个 int 无法表达"这次入库产生了哪些知识",
|
||
调用方也没法为多行知识建同义问法或回滚。所以返回**按 chunk 顺序排列的 `knowledge_id` 列表**;
|
||
空文档返回空列表。
|
||
|
||
## 尚未接入生产 worker(**已知缺口,不在 Task 5 范围**)
|
||
|
||
本服务投出的事件(`knowledge.vector_sync_requested`)**目前没有生产消费者**:
|
||
`WorkerRuntime.dispatch_one`(`app/worker/runtime.py`)的 handlers 是硬编码白名单,里面没有这两个
|
||
知识事件,需由组装层(`app/service/agent/bootstrap.py`,Task 7/9 范围)调用
|
||
`app/worker/knowledge_vector_worker.py` 的 `build_knowledge_handlers(...)` 并**注册进
|
||
`WorkerRuntime.dispatch_one` 的 handlers 字典**才能被领取。在那之前事件只会堆在
|
||
`domain_event_outbox`(现库已有 408 条),"入库成功"只代表**知识行与事件同一事务落库**,
|
||
**不代表向量已同步**。
|
||
|
||
## 校验口径
|
||
|
||
- `knowledge_type` 必须属于白名单;
|
||
- 集合名**不得由调用方随意指定**:只允许取 `TYPE_TO_COLLECTION[knowledge_type]`,
|
||
显式传入就必须**等于**它,否则 `ValidationAgentError`(全库只有三个知识集合,
|
||
写错集合等于让检索永远查不到这条知识);
|
||
- `created_by` 解析顺序:`ingest(created_by=...)` > 构造器 `created_by`;两者都没有则
|
||
`ValidationAgentError`(`fin_knowledge_meta.reviewer_id` 是"谁导入的"唯一线索,不能静默写空)。
|
||
"""
|
||
|
||
import hashlib
|
||
import json
|
||
import logging
|
||
from datetime import UTC, datetime
|
||
from pathlib import Path
|
||
from typing import Any
|
||
from uuid import uuid4
|
||
|
||
from sqlalchemy import select, text, update
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
from app.core.errors import ValidationAgentError
|
||
from app.core.knowledge_contracts import ALLOWED_COLLECTIONS
|
||
from app.model.knowledge import KnowledgeMeta
|
||
from app.service.document_parser import DocumentParser
|
||
|
||
# 裁定 5:事件常量只有一份真相,在 worker 模块里;这里 import 而不是另写字面量。
|
||
from app.worker.knowledge_vector_worker import (
|
||
KNOWLEDGE_ACTIVE_STATUS,
|
||
KNOWLEDGE_AGGREGATE_TYPE,
|
||
VECTOR_DELETE_EVENT,
|
||
VECTOR_SYNC_EVENT,
|
||
)
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
#: 下线状态取值。与 `knowledge_management_service.EXPIRED_STATUS` 同值;两个服务
|
||
#: 互相 import 会成环,故各自声明(DB 侧该列只有 'active' / 'expired' 两种取值)。
|
||
EXPIRED_STATUS = "expired"
|
||
ACTIVE_STATUS = KNOWLEDGE_ACTIVE_STATUS
|
||
|
||
#: 允许入库的知识类型与它们**唯一**对应的 Milvus 集合(`ALLOWED_COLLECTIONS` 是契约层白名单)。
|
||
TYPE_TO_COLLECTION: dict[str, str] = {
|
||
"faq": "fin_faq_collection",
|
||
"product": "fin_product_collection",
|
||
"policy": "fin_policy_collection",
|
||
}
|
||
ALLOWED_KNOWLEDGE_TYPES = frozenset(TYPE_TO_COLLECTION)
|
||
|
||
#: 入库默认版本号(同一份文档重导不会改变它;版本治理由后续任务的审核流程负责)。
|
||
DEFAULT_VERSION = "v1"
|
||
|
||
#: 标题列是 `String(256)`,chunk 正文可能远超;标题取正文前若干字符即可。
|
||
TITLE_MAX_LENGTH = 120
|
||
|
||
INSERT_KNOWLEDGE_META = text(
|
||
"INSERT INTO fin_knowledge_meta"
|
||
" (knowledge_type, title, source_file, minio_path, milvus_collection, version,"
|
||
" content_text, tags, reviewer_id, review_status, status, created_at, updated_at)"
|
||
" VALUES (:knowledge_type, :title, :source_file, :minio_path, :milvus_collection, :version,"
|
||
" :content_text, :tags, :reviewer_id, 'published', 'active', :now, :now)"
|
||
)
|
||
|
||
SELECT_LAST_INSERT_ID = text("SELECT LAST_INSERT_ID()")
|
||
|
||
INSERT_OUTBOX_EVENT = text(
|
||
"INSERT INTO domain_event_outbox"
|
||
" (event_id, event_type, aggregate_type, aggregate_id, trace_id, payload, status,"
|
||
" retry_count, occurred_at, created_at, updated_at)"
|
||
" VALUES (:event_id, :event_type, :aggregate_type, :aggregate_id, :trace_id, :payload,"
|
||
" 'pending', 0, :now, :now, :now)"
|
||
)
|
||
|
||
|
||
class KnowledgeIngestService:
|
||
def __init__(
|
||
self,
|
||
*,
|
||
session: AsyncSession,
|
||
parser: DocumentParser,
|
||
storage: Any,
|
||
embedder: Any,
|
||
endpoint_resolver: Any,
|
||
created_by: int | None = None,
|
||
) -> None:
|
||
self._session = session
|
||
self._parser = parser
|
||
self._storage = storage
|
||
self._embedder = embedder
|
||
self._endpoint_resolver = endpoint_resolver
|
||
self._created_by = created_by
|
||
|
||
async def ingest(
|
||
self,
|
||
*,
|
||
filename: str,
|
||
content: bytes,
|
||
knowledge_type: str,
|
||
collection: str | None = None,
|
||
created_by: int | None = None,
|
||
) -> list[int]:
|
||
"""把一份文档切成多块写入知识表,并为每一块投一条向量同步事件。
|
||
|
||
返回**该文档产生的 `knowledge_id` 列表**(按 chunk 顺序;空文档返回 `[]`)。
|
||
不 commit:事务归调用方(通常是与审核/审计记录一起提交,或整体回滚)。
|
||
"""
|
||
expected_collection = self._collection_for(knowledge_type, collection)
|
||
reviewer_id = self._resolve_created_by(created_by)
|
||
key = self._storage_key(knowledge_type, filename, content)
|
||
await self._storage.save(
|
||
key=key, content=content, content_type="application/octet-stream"
|
||
)
|
||
# 切分在落存储之后、写库之前:解析失败(UnsupportedDocumentError)时库里不留半截知识,
|
||
# 只在存储里留一个未引用的对象,比留一堆无事件的孤儿知识行安全。
|
||
chunks = self._parser.parse(filename=filename, content=content)
|
||
now = datetime.now(UTC).replace(tzinfo=None)
|
||
# 同一份文档**重新入库 = 覆盖上一版**:先把上一版下线(并清掉它的向量),再写新版。
|
||
# 见 `_supersede_previous_version` 里记录的后果。
|
||
# **必须先判 chunks 非空**:解析出 0 块(空文档、或全部是不支持的扩展名)时若照旧
|
||
# 下线上一版,结果是"旧版被下架、新版一行没写" —— 这份文档在检索侧就凭空消失了。
|
||
if chunks:
|
||
await self._supersede_previous_version(
|
||
filename=filename, collection=expected_collection, now=now
|
||
)
|
||
knowledge_ids: list[int] = []
|
||
for chunk in chunks:
|
||
knowledge_id = await self._insert_knowledge_row(
|
||
chunk=chunk,
|
||
filename=filename,
|
||
knowledge_type=knowledge_type,
|
||
collection=expected_collection,
|
||
storage_key=key,
|
||
reviewer_id=reviewer_id,
|
||
now=now,
|
||
)
|
||
knowledge_ids.append(knowledge_id)
|
||
await self._enqueue_vector_sync(knowledge_id, now=now)
|
||
return knowledge_ids
|
||
|
||
def _collection_for(self, knowledge_type: str, collection: str | None) -> str:
|
||
"""白名单校验:类型必须已知,显式给定的集合必须等于该类型的唯一集合。"""
|
||
if knowledge_type not in ALLOWED_KNOWLEDGE_TYPES:
|
||
raise ValidationAgentError(f"knowledge_type 非法:{knowledge_type}")
|
||
expected = TYPE_TO_COLLECTION[knowledge_type]
|
||
if collection is not None and collection != expected:
|
||
raise ValidationAgentError(
|
||
f"集合名只能是 {expected},不接受调用方指定的 {collection}"
|
||
)
|
||
if expected not in ALLOWED_COLLECTIONS: # pragma: no cover - 契约层与本地映射的护栏
|
||
raise ValidationAgentError(f"集合不在契约白名单内:{expected}")
|
||
return expected
|
||
|
||
def _resolve_created_by(self, created_by: int | None) -> int:
|
||
resolved = created_by if created_by is not None else self._created_by
|
||
if resolved is None:
|
||
raise ValidationAgentError("缺少 created_by:入库必须记录导入人")
|
||
return resolved
|
||
|
||
@staticmethod
|
||
def _storage_key(knowledge_type: str, filename: str, content: bytes) -> str:
|
||
"""`kb/<type>/<内容前 16 位 sha256>-<文件名>`:同一份文档重导落到同一个 key。"""
|
||
digest = hashlib.sha256(content).hexdigest()
|
||
return f"kb/{knowledge_type}/{digest[:16]}-{Path(filename).name}"
|
||
|
||
async def _insert_knowledge_row(
|
||
self,
|
||
*,
|
||
chunk: Any,
|
||
filename: str,
|
||
knowledge_type: str,
|
||
collection: str,
|
||
storage_key: str,
|
||
reviewer_id: int,
|
||
now: datetime,
|
||
) -> int:
|
||
"""插入一行知识并返回它的自增 `id`。
|
||
|
||
## 与 `LAST_INSERT_ID()` 的硬耦合(已知脆弱点,**不要**在中间插语句)
|
||
|
||
`LAST_INSERT_ID()` 是 **MySQL 连接级的隐式状态**,不是"上一次 INSERT 的返回值":
|
||
它只在**同一个连接**上、且**中间没有其它 INSERT** 时才等于本行 id。
|
||
因此本方法把 `INSERT` 与紧随其后的 `SELECT LAST_INSERT_ID()` 放在**同一方法体内成对**执行,
|
||
中间的语句数是 **0**,并删掉了 `sqlalchemy.text` 之外的任何额外查询
|
||
(`test_last_insert_id_is_read_once_per_chunk_in_the_same_statement_order`
|
||
把"成对且相邻"的语句序列固定下来)。调用方在本方法返回前不得插入别的写操作。
|
||
|
||
**为什么不用"显式回读"替代**(本方案评估过、明确否决):
|
||
`fin_knowledge_meta` 的唯一键只有主键 `id`(DDL 见
|
||
`docs/evidence/20260909-before-constraint-fix.sql`:
|
||
`PRIMARY KEY (id)` + `KEY idx_knowledge_status(status)`),
|
||
`(knowledge_type, source_file, minio_path)` 上**没有唯一约束**。回读只能靠这几列近似匹配,
|
||
而同一份文档用同一 key 重导会产出**完全相同的行** → 回读 `LIMIT 1` 可能拿到上次导入的
|
||
旧行 id(`ORDER BY id` 也救不了:新行 id 更大,反而更易混淆),会静默写坏 outbox 事件。
|
||
"先算好外部 id 再插入"同样不可行:`id` 是 `bigint unsigned AUTO_INCREMENT`,要显式指定就得
|
||
自己分配全局唯一整数,等于引入新表/新序列(触碰基线),且要额外处理并发分配的唯一性。
|
||
所以在本 Schema 下 `LAST_INSERT_ID()` 是**唯一**无需改基线的可行做法。
|
||
|
||
## 失败路径:两个分支的文案必须区分真相
|
||
|
||
- `rowcount == 0`(INSERT 没生效):此时 `LAST_INSERT_ID()` 保持**上一次**的值、
|
||
新连接上是 0,所以真相是"插入没生效",
|
||
**不能**报成"取不到 id"(那会把排障引向错误方向)。
|
||
- `LAST_INSERT_ID() == 0` 而 INSERT 影响了 1 行:说明连接上还没有过自增插入就拿到了 0,
|
||
文案直接点出 `LAST_INSERT_ID=0`,并提示不要对它取整。
|
||
"""
|
||
result = await self._session.execute(INSERT_KNOWLEDGE_META, {
|
||
"knowledge_type": knowledge_type,
|
||
"title": chunk.text[:TITLE_MAX_LENGTH],
|
||
"source_file": filename,
|
||
"minio_path": storage_key,
|
||
"milvus_collection": collection,
|
||
"version": DEFAULT_VERSION,
|
||
"content_text": chunk.text,
|
||
# `fin_knowledge_meta.tags` 是 JSON 列:这里写 json 字符串(与
|
||
# `tools/import_knowledge_seed.py` 同口径),读回来是 dict。
|
||
"tags": json.dumps(
|
||
{"heading_path": list(chunk.heading_path), "chunk_index": chunk.index},
|
||
ensure_ascii=False,
|
||
),
|
||
"reviewer_id": reviewer_id,
|
||
"now": now,
|
||
})
|
||
# `rowcount` 可能是 -1(驱动不报告行数),所以只在"明确为 0"时判失败,不把 -1 当失败。
|
||
affected = getattr(result, "rowcount", None)
|
||
if isinstance(affected, int) and affected == 0:
|
||
raise ValidationAgentError(
|
||
"知识行插入未生效(INSERT 影响 0 行),未继续取 LAST_INSERT_ID()"
|
||
)
|
||
inserted = await self._session.scalar(SELECT_LAST_INSERT_ID)
|
||
if inserted is None: # pragma: no cover - 只在非 MySQL 后端上发生
|
||
raise ValidationAgentError("插入知识行后取不到 LAST_INSERT_ID()(返回 NULL)")
|
||
if int(inserted) == 0: # pragma: no cover - 真机上的异常连接状态
|
||
raise ValidationAgentError(
|
||
"插入知识行后 LAST_INSERT_ID=0:取不到本次自增 id,不能拿它当 knowledge_id"
|
||
)
|
||
return int(inserted)
|
||
|
||
async def _supersede_previous_version(
|
||
self, *, filename: str, collection: str, now: datetime
|
||
) -> list[int]:
|
||
"""把同一份文档的上一版下线,**并投出向量删除事件**;返回被下线的 knowledge_id。
|
||
|
||
## 为什么必须在这里做
|
||
|
||
知识写入侧只认 active(`knowledge_vector_worker.KNOWLEDGE_ACTIVE_STATUS`),
|
||
但**检索侧不看 status**(`knowledge_search_service` 只过滤 `visibility`)。
|
||
于是"重新入库一份同名文档"会留下旧版的向量,它们继续参与排序 —— 表现就是
|
||
**同一份内容有 N 个副本互相抢答**:top1 与 top2 都是同一主题的近似块,
|
||
分数天然贴近,把客服的"领先次优 ≥0.07"门槛永远卡死(实测:产品手册被重复入库 7 次,
|
||
175 条历史副本的向量让「风险等级 R1–R5」这类问题只能转人工)。
|
||
|
||
## 口径
|
||
|
||
- 只下线**同一 `source_file` + 同一集合**的 active 行:不同文档互不影响;
|
||
- 下线 = `status='expired'` + 每行投一条 `knowledge.vector_delete_requested`
|
||
(消费侧 `knowledge_vector_worker.remove` 按 `knowledge_id` 删向量);
|
||
- 历史遗留的 expired 行(当年没投过删除事件、向量还在)不在本方法范围内,
|
||
用管理端口的"向量残留清理"动作单独处理(`cleanup_vector`)。
|
||
"""
|
||
previous = list(await self._session.scalars(
|
||
select(KnowledgeMeta.id).where(
|
||
KnowledgeMeta.source_file == filename,
|
||
KnowledgeMeta.milvus_collection == collection,
|
||
KnowledgeMeta.status == ACTIVE_STATUS,
|
||
)
|
||
))
|
||
if not previous:
|
||
return []
|
||
ids = [int(row) for row in previous]
|
||
await self._session.execute(
|
||
update(KnowledgeMeta)
|
||
.where(KnowledgeMeta.id.in_(ids))
|
||
.values(status=EXPIRED_STATUS, updated_at=now)
|
||
)
|
||
for knowledge_id in ids:
|
||
await self._enqueue_vector_delete(knowledge_id, now=now)
|
||
logger.info(
|
||
"知识覆盖入库:下线上一版 %s 个块(source_file=%s)并投出向量删除事件",
|
||
len(ids), filename,
|
||
)
|
||
return ids
|
||
|
||
async def _enqueue_vector_sync(self, knowledge_id: int, *, now: datetime) -> None:
|
||
"""投一条向量同步事件;`aggregate_id` 与 payload 都用 `str(knowledge_id)`。
|
||
|
||
`domain_event_outbox` 上没有唯一约束,所以**重复导入会产生多条同 `aggregate_id` 的事件**:
|
||
这是已知且被接受的 —— 消费侧按 `aggregate_id` 覆盖写入而非跳过
|
||
(见 `knowledge_vector_worker` 的幂等口径),因此重复投递不会在 Milvus 里留下重复向量。
|
||
"""
|
||
await self._session.execute(INSERT_OUTBOX_EVENT, {
|
||
"event_id": str(uuid4()),
|
||
"event_type": VECTOR_SYNC_EVENT,
|
||
"aggregate_type": KNOWLEDGE_AGGREGATE_TYPE,
|
||
"aggregate_id": str(knowledge_id),
|
||
"trace_id": str(uuid4()),
|
||
"payload": json.dumps({"knowledge_id": str(knowledge_id)}),
|
||
"now": now,
|
||
})
|
||
|
||
async def _enqueue_vector_delete(self, knowledge_id: int, *, now: datetime) -> None:
|
||
"""投一条向量**删除**事件(覆盖入库时清掉上一版用)。
|
||
|
||
与同步事件同一套装配:消费侧 `knowledge_vector_worker.remove` 按 `knowledge_id`
|
||
把该块从 Milvus 里删掉。事件与本次入库**同事务**,因此不会出现"库里有新版、
|
||
向量库里旧版还在"的中间态。
|
||
"""
|
||
await self._session.execute(INSERT_OUTBOX_EVENT, {
|
||
"event_id": str(uuid4()),
|
||
"event_type": VECTOR_DELETE_EVENT,
|
||
"aggregate_type": KNOWLEDGE_AGGREGATE_TYPE,
|
||
"aggregate_id": str(knowledge_id),
|
||
"trace_id": str(uuid4()),
|
||
"payload": json.dumps({"knowledge_id": str(knowledge_id)}),
|
||
"now": now,
|
||
})
|