diff --git a/app/api/controllers/knowledge_management.py b/app/api/controllers/knowledge_management.py index e619ca9..f1099f6 100644 --- a/app/api/controllers/knowledge_management.py +++ b/app/api/controllers/knowledge_management.py @@ -1,13 +1,16 @@ -"""知识库管理接口(Task 11):上传 / 查询 / 删除文档。 +"""知识库管理接口(Task 11):上传 / 查询 / 删除文档(+ 补投向量清理)。 老师《需求文档-修改版》Phase 1 验收第 7 条要求"知识库管理接口可正常上传/查询/删除文档", F1.2 点名了这三个端点。既有 `app/api/controllers/knowledge.py` 只有 `/api/v1/knowledge-references` (引用解析,K001),与此处是**两个不同的资源面**,因此单独一个模块、一个 router: 不动既有路由,也不把"引用解析"和"文档管理"混在同一个 prefix 下。 +第四个端点 `POST /{knowledge_id}/vector-cleanups` 是**运维补口**(不在老师点名的三个之内): +只给已经过期的历史行补投向量删除事件,理由见 Service 的 `cleanup_vector`。 + ## 鉴权(硬约束,不得绕过) -三个端点都声明 `Depends(build_request_context)`(并且 router 级叠了 `enforce_rate_limit`, +全部端点都声明 `Depends(build_request_context)`(并且 router 级叠了 `enforce_rate_limit`, 它自身又依赖 `build_request_context`,所以认证一定先于限流)。底座**没有**匿名路径, `context.user_id` 是入库 `created_by` 的唯一来源。 @@ -120,3 +123,22 @@ async def delete_document( 不存在的 id 返回 404(`SESSION_NOT_FOUND` 语义),**不静默成功**。 """ return await knowledge_management_service().delete_document(context, knowledge_id) + + +@router.post("/{knowledge_id}/vector-cleanups") +async def cleanup_document_vectors( + knowledge_id: int = Path(gt=0), + context: RequestContext = Depends(build_request_context), # noqa: B008 +) -> dict[str, Any]: + """补投一次向量清理:给**已经过期**的历史行补发 `knowledge.vector_delete_requested`。 + + 与 `DELETE` 的分工:`DELETE` 是"下线一份在用文档"(标记 + 投事件,一步到位), + 本端点**不改状态**,只把"这条已下线的知识,向量可能还在 Milvus 里"这件事补投出去。 + 没有它,历史上那些"被别的途径置为 expired、从未投过删除事件"的行在管理端口上无路可走 + (`DELETE` 对它们一律 404,见 Service docstring)。 + + - id 不存在 → 404; + - 行仍是 `active` → 422(请改用 `DELETE`,否则会出现"行是 active、向量已删"的不一致); + - 成功 → 200 + `vector_delete_event`,调用方能确认事件真的投了出去。 + """ + return await knowledge_management_service().cleanup_vector(context, knowledge_id) diff --git a/app/service/knowledge_ingest_service.py b/app/service/knowledge_ingest_service.py index 240f86c..99d1f10 100644 --- a/app/service/knowledge_ingest_service.py +++ b/app/service/knowledge_ingest_service.py @@ -42,20 +42,34 @@ Outbox handler(`app/worker/knowledge_vector_worker.py`)完成。这里只把 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 text +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_AGGREGATE_TYPE, VECTOR_SYNC_EVENT +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] = { @@ -132,6 +146,14 @@ class KnowledgeIngestService: # 只在存储里留一个未引用的对象,比留一堆无事件的孤儿知识行安全。 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( @@ -245,6 +267,51 @@ class KnowledgeIngestService: ) 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)`。 @@ -261,3 +328,20 @@ class KnowledgeIngestService: "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, + }) diff --git a/app/service/knowledge_management_service.py b/app/service/knowledge_management_service.py index 723def3..0738b79 100644 --- a/app/service/knowledge_management_service.py +++ b/app/service/knowledge_management_service.py @@ -1,4 +1,4 @@ -"""知识库管理服务(Task 11):上传 / 查询 / 删除文档。 +"""知识库管理服务(Task 11):上传 / 查询 / 删除文档 / 补投向量清理。 ## 为什么单独一个 Service,而不是塞进 `KnowledgeIngestService` @@ -23,9 +23,14 @@ 且**尽力而为**:归档失败只记日志,不回滚已经提交的删除 —— 否则会出现"库里还是 active、 存储里已经归档"的不一致,比"库已 expired、文件还在"危险得多(后者只浪费磁盘)。 +删除只覆盖"**本次调用**造成的 active→expired"。历史上已经过期、但从未投过删除事件的行 +(早期脚本改库、导入侧幂等上线前的重复副本)需要 `cleanup_vector` 补投一次 —— 为什么不能 +让 `delete_document` 顺带兼容它们:那会让"重复删除"变成静默成功,调用方再也分不清 +"这一次真的下线了"和"它早就过期了,我什么都没改"。 + ## 权限 -三个端点统一要求 `knowledge:manage`(admin 角色持有;权限码在 `tools/seed_test_rbac.py` +四个端点统一要求 `knowledge:manage`(admin 角色持有;权限码在 `tools/seed_test_rbac.py` 的 `PERMISSIONS` 里,id 9019)。复用既有 `knowledge:query`(9018)会让**客户角色**看见 整库文档清单与正文摘要,知识库的后台维护面不该向客户开放。 """ @@ -201,6 +206,55 @@ class KnowledgeManagementService: "vector_delete_event": VECTOR_DELETE_EVENT, } + async def cleanup_vector( + self, context: RequestContext, knowledge_id: int + ) -> dict[str, Any]: + """给**已经过期**的历史行补投一次向量删除事件(清掉向量侧残留)。 + + ## 为什么需要它 + + `delete_document` 只在"active → expired"这**一次状态转换**上投事件。历史上有一批行 + 是**别的途径**变成 `expired` 的(早期脚本直接改库、重复入库被覆盖、导入侧幂等上线前的 + 重复副本),它们**从来没投过**删除事件,向量还留在 Milvus 里继续参与排序; + 而 `delete_document` 对这些行**刻意返回 404**("不静默成功"),于是管理端口上 + 根本没有一个能清理它们的动作 —— 这是行为与运维需求之间的缺口,本方法就是补这个缺口。 + + ## 口径(三条,都是"宁可报错也不假装成功") + + 1. id 不存在 → 404(`SESSION_NOT_FOUND` 语义),与删除接口一致; + 2. 行**仍是 active** → 校验失败(422):在用的文档不该走"只清向量不标状态"这条路, + 那会让检索侧读不到正文而知识行看起来还好好的,请改用 DELETE(它同时做两件事); + 3. 成功 → 投出一条**新的** `knowledge.vector_delete_requested`(`event_id` 是新 uuid4, + 消费侧删除幂等,重复清理同一行只是多删一次空集,不会报错)。 + + 返回 `status`(当前状态)与 `vector_delete_event`(投出的事件名), + 调用方据此能确认"事件确实被投出去了",而不是只看一个 200。 + """ + await AuthorizationService.require(context, REQUIRED_PERMISSION) + async with self._session_factory() as session, session.begin(): + row = await session.get(KnowledgeMeta, knowledge_id) + if row is None: + raise GenericResourceNotFoundError("知识文档不存在") + status = str(row.status) + if status != EXPIRED_STATUS: + raise ValidationAgentError( + f"知识文档当前状态为 {status},只有已过期/已下线的文档才需要补投向量清理;" + "在用的文档请改用 DELETE 接口(它同时标记过期并投出删除事件)" + ) + now = datetime.now(UTC).replace(tzinfo=None) + self._enqueue_vector_delete( + session, knowledge_id, trace_id=context.trace_id, now=now + ) + logger.info( + "补投向量清理事件:knowledge_id=%s status=%s trace_id=%s", + knowledge_id, status, context.trace_id, + ) + return { + "knowledge_id": knowledge_id, + "status": status, + "vector_delete_event": VECTOR_DELETE_EVENT, + } + # --- 内部 --------------------------------------------------------------- @staticmethod diff --git a/docs/evidence/knowledge-vectors-reconcile.json b/docs/evidence/knowledge-vectors-reconcile.json new file mode 100644 index 0000000..2a20f8c --- /dev/null +++ b/docs/evidence/knowledge-vectors-reconcile.json @@ -0,0 +1,898 @@ +{ + "collections": { + "fin_faq_collection": { + "milvus_vectors": 126, + "mysql_rows": 42, + "mysql_active": 1, + "active_files": 1, + "orphan_vectors": 0, + "stale_vectors": 0, + "missing_vectors": 0, + "unbacked_seed_vectors": 125, + "duplicate_content_groups": 0, + "duplicate_content_vectors": 0, + "tiny_vectors": 59, + "heading_only_vectors": 0, + "duplicated_active_files": 0, + "detail": { + "orphan_vectors": [], + "stale_vectors_ids": [], + "missing_vectors": [], + "unbacked_seed_vectors": [ + "COMP-001", + "COMP-001-01", + "COMP-001-02", + "COMP-001-03", + "COMP-001-04", + "COMP-001-05", + "COMP-001-06", + "COMP-001-07", + "COMP-001-08", + "COMP-001-09", + "COMP-001-10", + "COMP-001-11", + "COMP-001-12", + "COMP-001-13", + "COMP-001-14", + "COMP-002", + "COMP-003", + "COMP-004", + "COMP-004-01", + "COMP-004-02", + "COMP-004-03", + "COMP-004-04", + "COMP-004-05", + "COMP-004-06", + "COMP-004-07", + "COMP-005", + "COMP-005-01", + "COMP-005-02", + "COMP-005-03", + "COMP-005-04", + "COMP-005-05", + "COMP-006", + "COMP-006-01", + "COMP-006-02", + "COMP-006-03", + "COMP-006-04", + "COMP-006-05", + "COMP-006-06", + "COMP-006-07", + "COMP-007", + "COMP-008", + "COMP-009", + "COMP-010", + "COMP-011", + "COMP-012", + "COMP-013", + "COMP-014", + "COMP-015", + "COMP-016", + "COMP-017", + "COMP-017-01", + "COMP-017-02", + "COMP-017-03", + "COMP-018", + "COMP-019", + "COMP-019-01", + "COMP-019-02", + "COMP-019-03", + "COMP-019-04", + "COMP-019-05", + "COMP-019-06", + "COMP-020", + "COMP-021", + "COMP-022", + "COMP-022-01", + "COMP-022-02", + "COMP-022-03", + "COMP-022-04", + "COMP-022-05", + "COMP-022-06", + "COMP-022-07", + "COMP-022-08", + "COMP-022-09", + "COMP-022-10", + "COMP-022-11", + "COMP-022-12", + "COMP-022-13", + "COMP-022-14", + "COMP-022-15", + "COMP-022-16", + "COMP-022-17", + "FAQ-0001", + "FAQ-0002", + "FAQ-0003", + "FAQ-0004", + "FAQ-0005", + "FAQ-0006", + "FAQ-0007", + "FAQ-0008", + "FAQ-0009", + "FAQ-0010", + "FAQ-0011", + "FAQ-0012", + "FAQ-0013", + "FAQ-0014", + "FAQ-0015", + "FAQ-0016", + "FAQ-0017", + "FAQ-0018", + "FAQ-0019", + "FAQ-0020", + "FAQ-0021", + "FAQ-0022", + "FAQ-0023", + "FAQ-0024", + "FAQ-0025", + "FAQ-0026", + "FAQ-0027", + "FAQ-0028", + "FAQ-0029", + "FAQ-0030", + "FAQ-0031", + "FAQ-0032", + "FAQ-0033", + "FAQ-0034", + "FAQ-0035", + "FAQ-0036", + "FAQ-0037", + "FAQ-0038", + "FAQ-0039", + "FAQ-0040", + "FAQ-0041", + "FAQ-0042", + "FAQ-0043", + "FAQ-0044" + ], + "tiny_vectors": [ + "COMP-001-01", + "COMP-001-02", + "COMP-001-03", + "COMP-001-04", + "COMP-001-05", + "COMP-001-06", + "COMP-001-07", + "COMP-001-08", + "COMP-001-09", + "COMP-001-10", + "COMP-001-11", + "COMP-001-12", + "COMP-001-13", + "COMP-001-14", + "COMP-004-01", + "COMP-004-02", + "COMP-004-03", + "COMP-004-04", + "COMP-004-05", + "COMP-004-06", + "COMP-004-07", + "COMP-005-01", + "COMP-005-02", + "COMP-005-03", + "COMP-005-04", + "COMP-005-05", + "COMP-006-01", + "COMP-006-02", + "COMP-006-03", + "COMP-006-04", + "COMP-006-05", + "COMP-006-06", + "COMP-006-07", + "COMP-017-01", + "COMP-017-02", + "COMP-017-03", + "COMP-019-01", + "COMP-019-02", + "COMP-019-03", + "COMP-019-04", + "COMP-019-05", + "COMP-019-06", + "COMP-022-01", + "COMP-022-02", + "COMP-022-03", + "COMP-022-04", + "COMP-022-05", + "COMP-022-06", + "COMP-022-07", + "COMP-022-08", + "COMP-022-09", + "COMP-022-10", + "COMP-022-11", + "COMP-022-12", + "COMP-022-13", + "COMP-022-14", + "COMP-022-15", + "COMP-022-16", + "COMP-022-17" + ], + "heading_only_vectors": [], + "duplicated_active_files": {}, + "duplicate_content_groups": {} + } + }, + "fin_policy_collection": { + "milvus_vectors": 297, + "mysql_rows": 0, + "mysql_active": 0, + "active_files": 0, + "orphan_vectors": 0, + "stale_vectors": 0, + "missing_vectors": 0, + "unbacked_seed_vectors": 297, + "duplicate_content_groups": 1, + "duplicate_content_vectors": 15, + "tiny_vectors": 208, + "heading_only_vectors": 0, + "duplicated_active_files": 0, + "detail": { + "orphan_vectors": [], + "stale_vectors_ids": [], + "missing_vectors": [], + "unbacked_seed_vectors": [ + "POL-AST-001", + "POL-AST-002", + "POL-AST-003", + "POL-AST-004", + "POL-AST-004-01", + "POL-AST-004-02", + "POL-AST-005", + "POL-AST-005-01", + "POL-AST-005-02", + "POL-AST-005-03", + "POL-AST-005-04", + "POL-AST-006", + "POL-AST-006-01", + "POL-AST-006-02", + "POL-AST-006-03", + "POL-AST-006-04", + "POL-AST-006-05", + "POL-AST-007", + "POL-AST-007-01", + "POL-AST-007-02", + "POL-AST-007-03", + "POL-AST-007-04", + "POL-AST-007-05", + "POL-AST-007-06", + "POL-AST-008", + "POL-AST-009", + "POL-AST-009-01", + "POL-AST-009-02", + "POL-AST-009-03", + "POL-AST-009-04", + "POL-AST-009-05", + "POL-AST-009-06", + "POL-AST-009-07", + "POL-AST-009-08", + "POL-AST-009-09", + "POL-AST-009-10", + "POL-AST-009-11", + "POL-AST-009-12", + "POL-AST-009-13", + "POL-AST-009-14", + "POL-AST-009-15", + "POL-AST-009-16", + "POL-AST-009-17", + "POL-AST-009-18", + "POL-AST-009-19", + "POL-AST-009-20", + "POL-AST-009-21", + "POL-AST-009-22", + "POL-AST-009-23", + "POL-AST-009-24", + "POL-AST-009-25", + "POL-AST-009-26", + "POL-AST-009-27", + "POL-AST-009-28", + "POL-AST-009-29", + "POL-AST-009-30", + "POL-AST-009-31", + "POL-AST-009-32", + "POL-AST-009-33", + "POL-AST-009-34", + "POL-AST-009-35", + "POL-AST-009-36", + "POL-AST-009-37", + "POL-AST-009-38", + "POL-AST-009-39", + "POL-AST-009-40", + "POL-AST-009-41", + "POL-AST-009-42", + "POL-AST-009-43", + "POL-AST-009-44", + "POL-AST-009-45", + "POL-AST-009-46", + "POL-AST-009-47", + "POL-AST-009-48", + "POL-AST-009-49", + "POL-AST-009-50", + "POL-AST-009-51", + "POL-AST-009-52", + "POL-AST-009-53", + "POL-AST-009-54", + "POL-AST-009-55", + "POL-AST-009-56", + "POL-AST-009-57", + "POL-AST-009-58", + "POL-AST-009-59", + "POL-AST-009-60", + "POL-AST-009-61", + "POL-AST-009-62", + "POL-AST-009-63", + "POL-AST-009-64", + "POL-AST-009-65", + "POL-AST-009-66", + "POL-AST-009-67", + "POL-AST-009-68", + "POL-AST-009-69", + "POL-AST-009-70", + "POL-AST-009-71", + "POL-AST-009-72", + "POL-AST-009-73", + "POL-AST-009-74", + "POL-AST-009-75", + "POL-AST-009-76", + "POL-AST-009-77", + "POL-AST-009-78", + "POL-AST-009-79", + "POL-AST-009-80", + "POL-AST-009-81", + "POL-AST-009-82", + "POL-AST-009-83", + "POL-AST-009-84", + "POL-AST-009-85", + "POL-AST-009-86", + "POL-AST-009-87", + "POL-AST-009-88", + "POL-AST-009-89", + "POL-AST-009-90", + "POL-AST-009-91", + "POL-AST-009-92", + "POL-AST-009-93", + "POL-AST-010", + "POL-AST-010-01", + "POL-AST-010-02", + "POL-AST-010-03", + "POL-AST-010-04", + "POL-AST-010-05", + "POL-AST-011", + "POL-AST-011-01", + "POL-AST-011-02", + "POL-AST-011-03", + "POL-AST-011-04", + "POL-AST-011-05", + "POL-AST-012", + "POL-AST-012-01", + "POL-AST-012-02", + "POL-AST-012-03", + "POL-AST-012-04", + "POL-AST-012-05", + "POL-AST-013", + "POL-AST-014", + "POL-AST-015", + "POL-AST-016", + "POL-AST-017", + "POL-AST-018", + "POL-AST-019", + "POL-AST-020", + "POL-AST-021", + "POL-AST-021-01", + "POL-AST-021-02", + "POL-AST-021-03", + "POL-AST-021-04", + "POL-AST-021-05", + "POL-AST-022", + "POL-AST-023", + "POL-AST-024", + "POL-AST-025", + "POL-AST-026", + "POL-AST-026-01", + "POL-AST-026-02", + "POL-AST-026-03", + "POL-AST-026-04", + "POL-AST-026-05", + "POL-AST-026-06", + "POL-AST-026-07", + "POL-AST-027", + "POL-AST-028", + "POL-AST-029", + "POL-AST-030", + "POL-SPM-001", + "POL-SPM-002", + "POL-SPM-003", + "POL-SPM-004", + "POL-SPM-005", + "POL-SPM-006", + "POL-SPM-006-01", + "POL-SPM-006-02", + "POL-SPM-006-03", + "POL-SPM-006-04", + "POL-SPM-006-05", + "POL-SPM-007", + "POL-SPM-007-01", + "POL-SPM-007-02", + "POL-SPM-007-03", + "POL-SPM-008", + "POL-SPM-009", + "POL-SPM-010", + "POL-SPM-010-01", + "POL-SPM-010-02", + "POL-SPM-010-03", + "POL-SPM-010-04", + "POL-SPM-010-05", + "POL-SPM-010-06", + "POL-SPM-010-07", + "POL-SPM-011", + "POL-SPM-012", + "POL-SPM-012-01", + "POL-SPM-012-02", + "POL-SPM-012-03", + "POL-SPM-012-04", + "POL-SPM-012-05", + "POL-SPM-012-06" + ], + "tiny_vectors": [ + "POL-AST-005-02", + "POL-AST-005-03", + "POL-AST-006-01", + "POL-AST-006-02", + "POL-AST-006-03", + "POL-AST-006-04", + "POL-AST-006-05", + "POL-AST-007-01", + "POL-AST-007-02", + "POL-AST-007-03", + "POL-AST-007-04", + "POL-AST-007-05", + "POL-AST-007-06", + "POL-AST-009-01", + "POL-AST-009-02", + "POL-AST-009-03", + "POL-AST-009-04", + "POL-AST-009-05", + "POL-AST-009-06", + "POL-AST-009-07", + "POL-AST-009-08", + "POL-AST-009-09", + "POL-AST-009-10", + "POL-AST-009-11", + "POL-AST-009-12", + "POL-AST-009-13", + "POL-AST-009-14", + "POL-AST-009-15", + "POL-AST-009-16", + "POL-AST-009-17", + "POL-AST-009-18", + "POL-AST-009-19", + "POL-AST-009-20", + "POL-AST-009-21", + "POL-AST-009-22", + "POL-AST-009-23", + "POL-AST-009-24", + "POL-AST-009-25", + "POL-AST-009-26", + "POL-AST-009-27", + "POL-AST-009-28", + "POL-AST-009-29", + "POL-AST-009-30", + "POL-AST-009-32", + "POL-AST-009-33", + "POL-AST-009-34", + "POL-AST-009-35", + "POL-AST-009-36", + "POL-AST-009-37", + "POL-AST-009-38", + "POL-AST-009-39", + "POL-AST-009-40", + "POL-AST-009-41", + "POL-AST-009-42", + "POL-AST-009-43", + "POL-AST-009-44", + "POL-AST-009-45", + "POL-AST-009-46", + "POL-AST-009-47", + "POL-AST-009-48", + "POL-AST-009-49", + "POL-AST-009-50", + "POL-AST-009-51", + "POL-AST-009-52", + "POL-AST-009-53", + "POL-AST-009-54", + "POL-AST-009-55", + "POL-AST-009-57", + "POL-AST-009-58", + "POL-AST-009-59", + "POL-AST-009-60", + "POL-AST-009-61", + "POL-AST-009-63", + "POL-AST-009-64", + "POL-AST-009-65", + "POL-AST-009-66", + "POL-AST-009-67", + "POL-AST-009-68", + "POL-AST-009-69", + "POL-AST-009-70", + "POL-AST-009-71", + "POL-AST-009-72", + "POL-AST-009-73", + "POL-AST-009-74", + "POL-AST-009-75", + "POL-AST-009-76", + "POL-AST-009-77", + "POL-AST-009-78", + "POL-AST-009-79", + "POL-AST-009-80", + "POL-AST-009-81", + "POL-AST-009-82", + "POL-AST-009-83", + "POL-AST-009-84", + "POL-AST-009-86", + "POL-AST-009-87", + "POL-AST-009-88", + "POL-AST-009-89", + "POL-AST-009-90", + "POL-AST-010-01" + ], + "heading_only_vectors": [], + "duplicated_active_files": {}, + "duplicate_content_groups": { + "534022938b4b": [ + "POL-AST-009-07", + "POL-AST-009-12", + "POL-AST-009-19", + "POL-AST-009-27", + "POL-AST-009-32", + "POL-AST-009-39", + "POL-AST-009-45", + "POL-AST-009-51", + "POL-AST-009-57", + "POL-AST-009-63", + "POL-AST-009-69", + "POL-AST-009-74", + "POL-AST-009-79", + "POL-AST-009-84", + "POL-AST-009-89" + ] + } + } + }, + "fin_product_collection": { + "milvus_vectors": 238, + "mysql_rows": 161, + "mysql_active": 24, + "active_files": 1, + "orphan_vectors": 0, + "stale_vectors": 0, + "missing_vectors": 0, + "unbacked_seed_vectors": 214, + "duplicate_content_groups": 0, + "duplicate_content_vectors": 0, + "tiny_vectors": 184, + "heading_only_vectors": 3, + "duplicated_active_files": 0, + "detail": { + "orphan_vectors": [], + "stale_vectors_ids": [], + "missing_vectors": [], + "unbacked_seed_vectors": [ + "HNW-001", + "HNW-001-01", + "HNW-001-02", + "HNW-001-03", + "HNW-001-04", + "HNW-002", + "HNW-003", + "HNW-003-01", + "HNW-003-02", + "HNW-003-03", + "HNW-003-04", + "HNW-004", + "HNW-005", + "HNW-006", + "HNW-007", + "PROD-001", + "PROD-001-01", + "PROD-001-02", + "PROD-001-03", + "PROD-001-04", + "PROD-001-05", + "PROD-001-06", + "PROD-001-07", + "PROD-001-08", + "PROD-001-09", + "PROD-001-10", + "PROD-001-11", + "PROD-001-12", + "PROD-001-13", + "PROD-002", + "PROD-002-01", + "PROD-002-02", + "PROD-002-03", + "PROD-002-04", + "PROD-002-05", + "PROD-002-06", + "PROD-002-07", + "PROD-002-08", + "PROD-002-09", + "PROD-002-10", + "PROD-002-11", + "PROD-002-12", + "PROD-002-13", + "PROD-002-14", + "PROD-002-15", + "PROD-003", + "PROD-003-01", + "PROD-003-02", + "PROD-003-03", + "PROD-003-04", + "PROD-003-05", + "PROD-003-06", + "PROD-003-07", + "PROD-003-08", + "PROD-003-09", + "PROD-003-10", + "PROD-003-11", + "PROD-003-12", + "PROD-003-13", + "PROD-003-14", + "PROD-003-15", + "PROD-003-16", + "PROD-004", + "PROD-004-01", + "PROD-004-02", + "PROD-004-03", + "PROD-004-04", + "PROD-004-05", + "PROD-004-06", + "PROD-004-07", + "PROD-004-08", + "PROD-004-09", + "PROD-004-10", + "PROD-004-11", + "PROD-004-12", + "PROD-004-13", + "PROD-004-14", + "PROD-004-15", + "PROD-004-16", + "PROD-004-17", + "PROD-004-18", + "PROD-004-19", + "PROD-004-20", + "PROD-004-21", + "PROD-004-22", + "PROD-004-23", + "PROD-005", + "PROD-005-01", + "PROD-005-02", + "PROD-005-03", + "PROD-005-04", + "PROD-005-05", + "PROD-005-06", + "PROD-005-07", + "PROD-005-08", + "PROD-005-09", + "PROD-005-10", + "PROD-005-11", + "PROD-005-12", + "PROD-005-13", + "PROD-005-14", + "PROD-005-15", + "PROD-005-16", + "PROD-005-17", + "PROD-005-18", + "PROD-005-19", + "PROD-006", + "PROD-006-01", + "PROD-006-02", + "PROD-006-03", + "PROD-006-04", + "PROD-006-05", + "PROD-006-06", + "PROD-006-07", + "PROD-006-08", + "PROD-006-09", + "PROD-006-10", + "PROD-006-11", + "PROD-006-12", + "PROD-006-13", + "PROD-006-14", + "PROD-006-15", + "PROD-007", + "PROD-007-01", + "PROD-007-02", + "PROD-007-03", + "PROD-007-04", + "PROD-007-05", + "PROD-007-06", + "PROD-007-07", + "PROD-007-08", + "PROD-007-09", + "PROD-007-10", + "PROD-007-11", + "PROD-007-12", + "PROD-008", + "PROD-008-01", + "PROD-008-02", + "PROD-008-03", + "PROD-008-04", + "PROD-008-05", + "PROD-008-06", + "PROD-008-07", + "PROD-008-08", + "PROD-008-09", + "PROD-008-10", + "PROD-008-11", + "PROD-008-12", + "PROD-009", + "PROD-009-01", + "PROD-009-02", + "PROD-009-03", + "PROD-009-04", + "PROD-009-05", + "PROD-009-06", + "PROD-009-07", + "PROD-009-08", + "PROD-009-09", + "PROD-009-10", + "PROD-010", + "PROD-010-01", + "PROD-010-02", + "PROD-010-03", + "PROD-010-04", + "PROD-010-05", + "PROD-010-06", + "PROD-010-07", + "PROD-010-08", + "PROD-010-09", + "PROD-010-10", + "PROD-011", + "PROD-011-01", + "PROD-011-02", + "PROD-011-03", + "PROD-011-04", + "PROD-011-05", + "PROD-011-06", + "PROD-011-07", + "PROD-011-08", + "PROD-011-09", + "PROD-011-10", + "PROD-011-11", + "PROD-011-12", + "PROD-011-13", + "PROD-011-14", + "PROD-011-15", + "PROD-012", + "PROD-012-01", + "PROD-012-02", + "PROD-012-03", + "PROD-012-04", + "PROD-012-05", + "PROD-012-06", + "PROD-012-07", + "PROD-012-08", + "PROD-012-09", + "PROD-012-10", + "PROD-012-11", + "PROD-013", + "PROD-014" + ], + "tiny_vectors": [ + "177", + "189", + "195", + "HNW-001-01", + "HNW-001-02", + "HNW-001-03", + "HNW-001-04", + "HNW-003-01", + "HNW-003-02", + "HNW-003-04", + "PROD-001-01", + "PROD-001-02", + "PROD-001-03", + "PROD-001-04", + "PROD-001-05", + "PROD-001-06", + "PROD-001-07", + "PROD-001-08", + "PROD-001-09", + "PROD-001-10", + "PROD-001-11", + "PROD-001-12", + "PROD-001-13", + "PROD-002-01", + "PROD-002-02", + "PROD-002-03", + "PROD-002-04", + "PROD-002-05", + "PROD-002-06", + "PROD-002-07", + "PROD-002-08", + "PROD-002-09", + "PROD-002-10", + "PROD-002-11", + "PROD-002-12", + "PROD-002-13", + "PROD-002-14", + "PROD-003-01", + "PROD-003-02", + "PROD-003-03", + "PROD-003-04", + "PROD-003-05", + "PROD-003-06", + "PROD-003-07", + "PROD-003-08", + "PROD-003-09", + "PROD-003-10", + "PROD-003-11", + "PROD-003-12", + "PROD-003-13", + "PROD-003-14", + "PROD-003-15", + "PROD-004-01", + "PROD-004-02", + "PROD-004-03", + "PROD-004-04", + "PROD-004-05", + "PROD-004-06", + "PROD-004-07", + "PROD-004-08", + "PROD-004-09", + "PROD-004-10", + "PROD-004-11", + "PROD-004-12", + "PROD-004-13", + "PROD-004-14", + "PROD-004-16", + "PROD-004-17", + "PROD-004-18", + "PROD-004-19", + "PROD-004-20", + "PROD-004-21", + "PROD-004-22", + "PROD-004-23", + "PROD-005-01", + "PROD-005-02", + "PROD-005-03", + "PROD-005-04", + "PROD-005-05", + "PROD-005-06", + "PROD-005-07", + "PROD-005-08", + "PROD-005-09", + "PROD-005-10", + "PROD-005-11", + "PROD-005-12", + "PROD-005-14", + "PROD-005-15", + "PROD-005-16", + "PROD-005-17", + "PROD-005-18", + "PROD-005-19", + "PROD-006-01", + "PROD-006-02", + "PROD-006-03", + "PROD-006-04", + "PROD-006-05", + "PROD-006-06", + "PROD-006-07", + "PROD-006-08" + ], + "heading_only_vectors": [ + "177", + "189", + "195" + ], + "duplicated_active_files": {}, + "duplicate_content_groups": {} + } + } + }, + "errors": {}, + "mysql_rows_total": 203, + "collections_checked": [ + "fin_faq_collection", + "fin_policy_collection", + "fin_product_collection" + ] +} \ No newline at end of file diff --git a/docs/演示用/知识库向量对账与清理-2026-09-15.md b/docs/演示用/知识库向量对账与清理-2026-09-15.md new file mode 100644 index 0000000..488b4cf --- /dev/null +++ b/docs/演示用/知识库向量对账与清理-2026-09-15.md @@ -0,0 +1,214 @@ +# 知识库「向量 ↔ 元数据」对账 + 导入幂等 + 清理入口(2026-09-15) + +> 承接 `docs/演示用/知识库问答诊断-2026-09-14.md`(那份是**现象与根因诊断**,这份是**动手修 + 对账实测**)。 +> 本文所有数字都是本机真机实测,命令都在 §7,可逐条复现。 + +## 摘要(先看这四句) + +1. 三件事都做完了并真机验证通过:**① 对账工具、② 导入侧幂等、③ 清理入口**。 +2. 对账实测:**孤儿向量 0 条、死向量 0 条、缺向量 0 条** —— 也就是说 §4 之前担心的 + "向量库和元数据严重不一致"**并不存在**,那批 `FAQ-0013` 式的向量是**设计如此**(见 §5 更正)。 +3. 但查出一个**更严重、也更具体**的问题:**661 条向量里有 451 条正文不到 40 字**(占 68%), + 它们是灌库时被切碎的表格行/标题行。实测「风险评估问卷怎么评分」一问,**前 4 名是 4 条 + 一模一样的 19 字碎片**(0.7378 / 0.7354 / 0.7350 / 0.7350),真正有内容的 2828 字那条排第 5 + —— gap 只有 **0.0024**,客服**必然**判「候选并列」转人工。**这条问题不是靠补内容能修的。** +4. 需要你拍板的只有一件事:**要不要修灌库脚本的切分**(§6)。其余都已闭环。 + +--- + +## 1. 交付了什么 + +| # | 东西 | 位置 | 性质 | +|---|---|---|---| +| ① | 向量-元数据对账工具 | `tools/reconcile_knowledge_vectors.py` | **只读**,不改任何东西 | +| ② | 导入侧幂等(同文档重传 = 覆盖上一版) | `app/service/knowledge_ingest_service.py` | 写路径行为变更 | +| ③ | 清理入口:给已过期行补投向量删除 | `POST /api/v1/knowledge/{id}/vector-cleanups`
`app/service/knowledge_management_service.py: cleanup_vector` | **新增端点**(老师点名的三个端点一个没动) | +| ③ | 批量补投脚本(驱动的就是上面那个端点) | `tools/purge_expired_knowledge_vectors.py` | 默认 dry-run | + +### ② 为什么是"覆盖上一版"而不是"按内容哈希去重" + +写入侧只认 `active`(`knowledge_vector_worker.KNOWLEDGE_ACTIVE_STATUS`),**但检索侧不看 `status`** +(`knowledge_search_service` 只过滤 `visibility`)。所以"重传一份同名文档"必须**主动把上一版下线 +并投出向量删除事件**,否则旧向量继续参与排序(同一内容多副本抢答)。 + +口径:只下线**同一 `source_file` + 同一集合**的 `active` 行(不同文档互不影响), +逐行投一条 `knowledge.vector_delete_requested`,与本次入库**同一事务**。 + +⚠️ 实现时发现并修掉一个真 bug:**先判 `chunks` 非空再下线**。否则传一份解析出 0 块的文档 +(空文件/不支持的扩展名)会把上一版下架、新版一行不写 —— 这份文档在检索侧**凭空消失**。 +已加用例 `test_empty_document_must_not_expire_the_previous_version`。 + +### ③ 为什么不能顺手让 `DELETE` 兼容已过期行 + +`DELETE` 对已 `expired` 的行**刻意返回 404**("不静默成功")。如果让它对已过期行也返回 200, +调用方就再也分不清"这一次真的下线了"和"它早就过期了、我什么都没改"。 +所以新开一个**语义明确**的端点: + +| 情况 | 行为 | +|---|---| +| id 不存在 | **404** `SESSION_NOT_FOUND` | +| 行仍是 `active` | **422**,文案让调用方改用 `DELETE`(否则会出现"行是 active、向量已删"的不一致) | +| 行已 `expired` | **200** + 返回 `vector_delete_event`,调用方能确认事件**真的投出去了**(不是只看一个 200) | + +清理动作**不改状态**,只补投事件;真正删除由 Worker 消费事件完成 —— 权限、审计、 +Outbox 幂等语义一个都不绕。 + +--- + +## 2. 真机验证证据(②③ 端到端) + +脚本:`%TEMP%\verify_ingest_idempotency.py`(临时脚本,流程已固化在下面的命令里)。实测输出: + +``` +① 首次上传 + upload HTTP 201 {"knowledge_ids":[204],...} + 等 Worker 投向量 … 向量已存在:[True] + +② 同文件名重传(应触发覆盖) + upload HTTP 201 {"knowledge_ids":[205],...} + 库里该文件名的行:{204: 'expired', 205: 'active'} ← 旧行确实被下线 + 旧 id 的事件:[('knowledge.vector_sync_requested', 'published'), + ('knowledge.vector_delete_requested', 'published')] ← 删除事件已投且已消费 + 等 Worker 删除旧向量 … 旧向量是否还在:已全部删除 ✅ ← 向量真从 Milvus 没了 + +③ 端点 vector-cleanups + expired 行 → HTTP 200 {"knowledge_id":204,"status":"expired","vector_delete_event":"knowledge.vector_delete_requested"} + active 行 → HTTP 422 AGENT_INPUT_INVALID「知识文档当前状态为 active,只有已过期/已下线的文档才需要补投向量清理;在用的文档请改用 DELETE 接口…」 + 不存在 id → HTTP 404 SESSION_NOT_FOUND「知识文档不存在」 + +④ 收尾:删掉本次上传的 active 行 → 200;本次残留向量:无 ✅ +``` + +**两个问句回归**(改动后仍正常回答,未受影响): + +``` +问:r1到r5分别代表什么 → 处置:直接回答 top1=0.7387 次优=0.6621 gap=0.0766(≥0.07 可答) +问:基金申购后多久能确认 → 处置:直接回答 top1=0.8453 次优=0.6336 gap=0.2117(高置信) +``` + +--- + +## 3. 对账结果:`python tools/reconcile_knowledge_vectors.py` + +``` +集合 向量 MySQL行 active 文档数 孤儿 死向量 缺向量 无元数据种子 重复正文 碎片 纯标题 +------------------------------------------------------------------------------------------------------------------- +fin_faq_collection 126 42 1 1 0 0 0 125 0 59 0 +fin_policy_collection 297 0 0 0 0 0 0 297 15 208 0 +fin_product_collection 238 161 24 1 0 0 0 214 0 184 3 +``` + +怎么读(工具在 stdout 里也印了这几行): + +> **快照口径**:上表取自 2026-09-15 端到端验证**之前**(`fin_product_collection` MySQL 161 行)。 +> 验证过程上传了 2 行测试数据并在收尾时下线(**均无向量**),所以现在重跑会看到该列是 **163**, +> 其余各列不变 —— `死向量` 仍是 0 就是这件事的直接证据。 + +- **孤儿向量 0 / 死向量 0 / 缺向量 0** → 当前**没有**"向量与元数据不一致"的问题。健康。 +- **无元数据种子向量 636 条** → `FAQ-0013` / `POL-AST-009-07` / `COMP-001` 这类**人工语义 id**, + 由灌库脚本直接写 Milvus,**本来就不该出现在 MySQL**。工具**单独归一类、不建议删** + —— 答得最好的内容恰恰在这批里,当孤儿清理会把唯一正确的答案删掉。 +- **碎片 451 条**(<40 字)、**纯标题 3 条** → 见 §4,这是本次最重要的发现。 +- **重复正文 15 条向量**:policy 集合里 `POL-AST-009-07/-12/-19/…` 共 15 条**正文逐字相同** + (都是那句 19 字的「第九条 问卷内容及评分标准:选项 分值」),检索时必然互相打平。 +- 「向量」列取自 Milvus `query` 的**实际行数**。早前记的 159/397/297 是 + `get_collection_stats` 的 `row_count`,它**把已软删、尚未 compaction 的行也算进去** —— + 两个数不等是正常的,**以对账工具这一列为准**。 + +--- + +## 4. 新发现(必须看):68% 的向量是"回答不了任何问题"的碎片 + +三个集合的正文长度分布(实测): + +| 集合 | 向量数 | 中位长度 | <20 字 | 20–39 字 | 40–99 | 100–299 | 300+ | +|---|---|---|---|---|---|---|---| +| `fin_faq_collection` | 126 | 57 | 10 | 49 | 14 | 50 | 3 | +| `fin_policy_collection` | 297 | **28** | 59 | 149 | 22 | 47 | 20 | +| `fin_product_collection` | 238 | **23** | 47 | 137 | 11 | 16 | 27 | + +典型样本(就是这些内容在进排序): + +``` +'## 一、交易规则' ← 纯标题,9 字 +'分层体系:金卡 财富管理客户' ← 表格被切碎成"标签 值" +'第九条 问卷内容及评分标准:选项 分值' ← 19 字,且有 15 条完全一样 +'第五条 专业投资者认定标准:年收入 最近3年个人年均收入不低于50万元人民币' +``` + +### 实测后果(`风险评估问卷怎么评分`,同一套线上检索) + +``` +1. score=0.7378 len= 19 第九条 问卷内容及评分标准:选项 分值 ← 碎片 +2. score=0.7354 len= 19 第九条 问卷内容及评分标准:选项 分值 ← 碎片(内容与前一条逐字相同) +3. score=0.7350 len= 19 第九条 问卷内容及评分标准:选项 分值 ← 碎片 +4. score=0.7350 len= 19 第九条 问卷内容及评分标准:选项 分值 ← 碎片 +5. score=0.6640 len=2828 第九条 问卷内容及评分标准(完整一节) ← **真正的答案在这里** +→ top1=0.7378 次优=0.7354 gap=0.0024 → 客服判「候选并列」,转人工 +``` + +**这条问题的根因不是"知识库里没有内容",而是"碎片把答案挤到第 5 名、gap 被抹平"。** +补内容、改问法都救不了它 —— 因为挡在前面的是 4 条**内容完全相同**的 19 字碎片。 + +### 结论与建议 + +- 这是**灌库脚本的切分缺陷**(按行切 markdown,把表格的"标签 值"拆成独立块、 + 把标题单独成块),不是检索层的错,也**不需要改客服用例的阈值**。 +- 建议的修法(**未做,等你拍板**):灌库脚本改为"按章节切、块内保留完整语义段落", + 并把 <40 字的块**合并到相邻块**;重建受影响的两个集合后跑一次对账确认碎片归零。 +- 反例警告:**不要**为了这个问题去放宽 `MIN_GAP`(0.07)—— 那会让"4 条一样的碎片" + 里随便一条被当作答案吐给客户,比转人工危险得多。 + +--- + +## 5. 对 2026-09-14 诊断的两处更正(重要) + +那份诊断有两处推断被本次实测**推翻**,留在这里以免继续按错的结论做事: + +| 原结论 | 实测事实 | 说明 | +|---|---|---| +| §4.1「向量库与元数据**严重不一致**(大量孤儿向量)」 | **孤儿 0 条** | 当时比较的是 `get_collection_stats` 的 `row_count`(含已软删行)与 MySQL 行数,差额被误读成孤儿。实际上 `FAQ-*`/`POL-*`/`COMP-*` 是灌库脚本**有意**写进 Milvus 的语义 id,两边**本来就不该一一对应**。 | +| §4.3「175 行历史副本的向量让「R1–R5」卡死」 | 那 **175 行没有向量**(死向量 0 条) | 它们只是 MySQL 里的历史行,从未投过向量同步事件 → 从未参与过排序。R1–R5 当时 gap 不够的真实原因是**同一主题的多条短碎片**(当时实测次优 0.6412 / 20 字,正是 policy 集合的表格碎片)。 | + +不变的结论:**检索侧不看 `status`**(§4.2)依然成立且依然要修 —— 只是"重传同名文档" +这条路径现在已经被 ② 堵上了;历史遗留行则由 ③ 提供补投入口。 + +--- + +## 6. 没做的 / 需要你拍板的 + +1. **修灌库脚本的切分**(§4):影响面是重建 `fin_policy_collection` / `fin_product_collection`, + 属"要重灌知识库"的动作,**必须你同意再动**。 +2. **孤儿向量的物理删除**:孤儿(Milvus 有、MySQL 连行都没有)没有可挂事件的载体, + 真删只能直连 Milvus,会绕过权限/审计/Outbox。**对账工具只报告、不处理**; + 当前实测 0 条,所以也不急。 +3. **阈值策略**(`MIN_GAP=0.07`、"同一主题多来源不算并列"):**一行没改**。 + 本期所有修复都在内容侧与数据侧,策略仍归你决定。 +4. **那 38 行上传链路测试垃圾**(`e2e-*.md` / `check-*.md` / `verify-*.md` / `diag-test.md`): + 现在可以用 ③ 的端点按 id 清理了,但清理它们**不影响**上面任何结论,属可选卫生工作。 + +--- + +## 7. 复现命令清单 + +```bash +# ① 对账(只读;明细落 docs/evidence/knowledge-vectors-reconcile.json) +python tools/reconcile_knowledge_vectors.py +python tools/reconcile_knowledge_vectors.py --json --no-write # 只看不落盘 + +# ③ 清理(默认 dry-run;--apply 才真投,每条一次 POST) +python tools/purge_expired_knowledge_vectors.py +python tools/purge_expired_knowledge_vectors.py --apply --wait 30 + +# 两个问句回归 +python tools/ask_customer_service.py "r1到r5分别代表什么" +python tools/ask_customer_service.py "基金申购后多久能确认" + +# 单测 +python -m pytest tests/unit/service/test_knowledge_ingest_service.py \ + tests/unit/service/test_knowledge_management_service.py -q +``` + +⚠️ 前置:**API 与 Worker 都要在跑**(`python -m uvicorn app.main:app --port 8000`、 +`python -m app.worker`)。本机可直接双击 `启动金融Agent平台.bat`。 +改动新增了端点,**API 必须重启**才会生效(uvicorn 没开 `--reload`)。 diff --git a/docs/演示用/知识库问答诊断-2026-09-14.md b/docs/演示用/知识库问答诊断-2026-09-14.md index e6956cb..b4e404a 100644 --- a/docs/演示用/知识库问答诊断-2026-09-14.md +++ b/docs/演示用/知识库问答诊断-2026-09-14.md @@ -70,7 +70,14 @@ tools/seed_knowledge_r1r5_faq.py # 幂等上传(默认 dry-run,--appl ### 4.1 向量库与元数据严重不一致(大量"孤儿向量") -| 集合 | Milvus 向量数 | MySQL `fin_knowledge_meta` 行数 | +> ❌ **2026-09-15 更正:本节结论不成立,实测孤儿向量 0 条。** 下表左列的数是 +> `get_collection_stats` 的 `row_count`,它**把已软删、尚未 compaction 的行也算进去**, +> 拿来和 MySQL 行数相减得到的"差额"并不是孤儿。详见 +> `docs/演示用/知识库向量对账与清理-2026-09-15.md` §5。 +> **仍然成立的部分**:`faq/高频问答对.txt` 这类内容确实只在 Milvus 里 —— +> 但那是灌库脚本**有意**用的语义 id(`FAQ-0013`),不是"不一致"。 + +| 集合 | Milvus 向量数(当时记的) | MySQL `fin_knowledge_meta` 行数 | |---|---|---| | `fin_faq_collection` | **159** | 41 | | `fin_product_collection` | **397** | 161 | @@ -92,6 +99,11 @@ tools/seed_knowledge_r1r5_faq.py # 幂等上传(默认 dry-run,--appl ### 4.3 重复与测试垃圾都在抢答,而且没有清理入口 +> ❌ **2026-09-15 部分更正**:这 175 行历史副本**从来没有向量**(实测"死向量 0 条"), +> 因此**从未参与过排序** —— 它们不是当时 gap 不够的原因(真原因是 §4.1 之外的**短碎片**, +> 见新文档 §4)。「没有清理入口」这一条**属实**,已由新端点 +> `POST /api/v1/knowledge/{id}/vector-cleanups` + `tools/purge_expired_knowledge_vectors.py` 补上。 + - 《场内基金产品手册》被**重复种了 7 遍**:MySQL 199 行里只有 24 行 `active`(最新一套 176–199), 其余 175 行是历史副本; - `fin_faq_collection` 的 41 行里有 **38 行是上传链路测试留下的**(`e2e-*.md` / `check-*.md` / @@ -101,10 +113,32 @@ tools/seed_knowledge_r1r5_faq.py # 幂等上传(默认 dry-run,--appl ## 5. 遗留建议(未做,等决策) +> ⚠️ **2026-09-15 更新:下面 3 条已全部实现,且 §4.1 / §4.3 有两处推断**被实测推翻**。 +> 见 `docs/演示用/知识库向量对账与清理-2026-09-15.md`。要点先记在这里:** +> +> - **§4.1「大量孤儿向量」不成立**:实测孤儿向量 **0 条**。当时拿 +> `get_collection_stats` 的 `row_count`(**含已软删、未 compaction 的行**)去比 MySQL 行数, +> 差额被误读成孤儿;`FAQ-*`/`POL-*`/`COMP-*` 是灌库脚本**有意**写进 Milvus 的语义 id, +> 两边本来就不该一一对应。 +> - **§4.3「175 行历史副本的向量让 R1–R5 卡死」不成立**:那 175 行**没有向量**(死向量 0 条), +> 它们从未参与排序。当时 gap 不够的真实原因是**同一主题的多条短碎片**(次优那条 20 字, +> 正是 policy 集合被切碎的表格行)。 +> - **新查出更严重的问题**:661 条向量里 **451 条正文不到 40 字**(占 68%)。 +> 实测「风险评估问卷怎么评分」前排是 **4 条一模一样的 19 字碎片**(gap 0.0024), +> 真正 2828 字的答案排第 5 → 客服必然转人工。**这是灌库切分缺陷,补内容救不了, +> 也不该靠放宽 `MIN_GAP` 去"解决"。** + 1. **向量与元数据对账**:以 MySQL 为权威,补一个对账工具(列出"有向量无元数据""有元数据无向量"), 并把对账结果作为"知识库健康度"的一项输出。**这是上述所有问题的公共根因。** + → ✅ 已做:`tools/reconcile_knowledge_vectors.py`(只读) 2. **导入侧幂等**:知识入库按 `(source_file, content_hash)` upsert,避免同一份文档反复灌出 7 份副本。 + → ✅ 已做:同 `source_file` + 集合重传即**覆盖上一版**(下线旧行 + 逐行投向量删除事件), + 实际口径见新文档 §1(比 `content_hash` upsert 更贴合"重导一份文档"的真实操作)。 3. **清理入口**:给知识删除接口加"可删除已过期行"的口径(或提供一个管理端"清理历史副本"动作), 否则重复副本只能靠直接改库清理。 + → ✅ 已做:`POST /api/v1/knowledge/{id}/vector-cleanups`(仅对已 `expired` 行;`DELETE` 的 + 404 口径**保持不变**,因为"重复删除静默成功"会让调用方分不清真相) 4. (可选)**门槛策略**:如果业务接受"同一主题多来源不算并列",可在判定里加"同主题"识别; **本轮未改**——内容是能修的,策略一动影响面太大,需业务拍板。 + → ⏸ **仍然未改**(一行没动)。且新证据表明:**这个问题不该靠改阈值解决**(见上面的警告)。 + diff --git a/tests/unit/api/test_controller_routing_contract.py b/tests/unit/api/test_controller_routing_contract.py index 7cb4df4..7a0e416 100644 --- a/tests/unit/api/test_controller_routing_contract.py +++ b/tests/unit/api/test_controller_routing_contract.py @@ -36,6 +36,7 @@ PROTECTED_POST = [ "/api/v1/agent-runs/run-x/cancellations", "/api/v1/conversation-messages/1/feedback", "/api/v1/knowledge/upload", + "/api/v1/knowledge/1/vector-cleanups", ] PROTECTED_DELETE = [ diff --git a/tests/unit/service/test_knowledge_ingest_service.py b/tests/unit/service/test_knowledge_ingest_service.py index bdd8d50..430aa2a 100644 --- a/tests/unit/service/test_knowledge_ingest_service.py +++ b/tests/unit/service/test_knowledge_ingest_service.py @@ -43,17 +43,29 @@ class FakeParser: class FakeSession: - """记录每条带参数的语句;`SELECT LAST_INSERT_ID()` 按预设序列返回。""" + """记录每条带参数的语句;`SELECT LAST_INSERT_ID()` 按预设序列返回。 - def __init__(self, ids: list[Any] | None = None) -> None: + `previous_ids` 是"同一 source_file + 集合上已经 active 的行"(导入侧幂等的替身输入): + `_supersede_previous_version` 用 `scalars()` 读它们,本替身直接把它当查询结果返回。 + """ + + def __init__(self, ids: list[Any] | None = None, previous_ids: list[int] | None = None) -> None: self._ids = list(ids if ids is not None else []) + self.previous_ids = list(previous_ids if previous_ids is not None else []) self.statements: list[tuple[str, dict[str, Any]]] = [] self.commits = 0 async def execute(self, statement: Any, params: Any = None) -> None: - self.statements.append((str(statement), dict(params or {}))) + # SQLAlchemy 构造式语句(UPDATE ... WHERE id IN (...))的绑定值不在 `params` 里, + # 而是编译进语句本身;不取出来就看不见"到底把哪些行改成了什么状态"。 + values = dict(params or {}) or _compiled_params(statement) + self.statements.append((str(statement), values)) return None + async def scalars(self, statement: Any, params: Any = None) -> Any: + self.statements.append((str(statement), dict(params or {}) or _compiled_params(statement))) + return list(self.previous_ids) + async def scalar(self, statement: Any, params: Any = None) -> Any: self.statements.append((str(statement), dict(params or {}))) return self._ids.pop(0) if self._ids else None @@ -62,6 +74,14 @@ class FakeSession: self.commits += 1 +def _compiled_params(statement: Any) -> dict[str, Any]: + """取构造式语句编译后的绑定参数(纯 SQL 字符串没有可编译对象)。""" + compile_ = getattr(statement, "compile", None) + if compile_ is None: + return {} + return dict(compile_().params) + + class RecordingEmbedder: async def embed(self, endpoints: Any, text: str, *, max_attempts: int = 2) -> Any: raise AssertionError("入库服务不得自己调用 embedding(向量由 Outbox Worker 负责)") @@ -84,9 +104,10 @@ def _service( ids: list[Any] | None = None, chunks: list[ParsedChunk] | None = None, created_by: int | None = 9003, + previous_ids: list[int] | None = None, ) -> tuple[KnowledgeIngestService, FakeSession, FakeStorage, FakeParser]: """默认给 2 个 chunk 配 2 个自增 id(`LAST_INSERT_ID()` 的替身)。""" - session = FakeSession(ids if ids is not None else [11, 12]) + session = FakeSession(ids if ids is not None else [11, 12], previous_ids) storage = FakeStorage() parser = FakeParser(chunks if chunks is not None else _chunks()) service = KnowledgeIngestService( @@ -290,10 +311,13 @@ async def test_last_insert_id_is_read_once_per_chunk_in_the_same_statement_order "INSERT INTO fin_knowledge_meta" if sql.startswith("INSERT INTO fin_knowledge_meta") else "SELECT LAST_INSERT_ID()" if "LAST_INSERT_ID" in sql else "INSERT INTO domain_event_outbox" if sql.startswith("INSERT INTO domain_event_outbox") + else "SELECT previous active ids" if "fin_knowledge_meta.id" in sql else sql for sql, _ in session.statements ] assert shape == [ + # 导入侧幂等的第一步:查同 source_file + 集合的 active 旧版(本次为空,故没有 UPDATE)。 + "SELECT previous active ids", "INSERT INTO fin_knowledge_meta", "SELECT LAST_INSERT_ID()", "INSERT INTO domain_event_outbox", @@ -301,3 +325,96 @@ async def test_last_insert_id_is_read_once_per_chunk_in_the_same_statement_order "SELECT LAST_INSERT_ID()", "INSERT INTO domain_event_outbox", ] + + +# --- 导入侧幂等:同一份文档重传 = 覆盖上一版 --------------------------------------- + + +def _vector_events(session: FakeSession, event_type: str) -> list[dict[str, Any]]: + return [ + params + for _, params in session.statements + if params.get("event_type") == event_type + ] + + +async def test_reingesting_the_same_file_expires_the_previous_version() -> None: + """重传同名文档必须**先下线上一版并逐块投出向量删除事件**。 + + 不这么做的话,旧版向量会留在 Milvus 里继续参与排序(检索侧不看 status), + 表现就是"同一份内容有 N 个副本互相抢答"、客服的"领先次优 ≥0.07"门槛被永远卡死。 + """ + service, session, _, _ = _service(ids=[11, 12], previous_ids=[3, 4, 5]) + + await service.ingest(filename="faq.md", content=b"x", knowledge_type="faq") + + updates = [ + (sql, params) for sql, params in session.statements if sql.startswith("UPDATE") + ] + assert len(updates) == 1 + _, params = updates[0] + assert params["status"] == "expired" + # `id.in_([...])` 编译后是一个列表参数(参数名由 SQLAlchemy 生成,不写死名字)。 + id_lists = [value for value in params.values() if isinstance(value, list)] + assert id_lists == [[3, 4, 5]] + + deletes = _vector_events(session, "knowledge.vector_delete_requested") + assert [event["aggregate_id"] for event in deletes] == ["3", "4", "5"] + assert [json.loads(event["payload"]) for event in deletes] == [ + {"knowledge_id": "3"}, {"knowledge_id": "4"}, {"knowledge_id": "5"} + ] + assert {event["aggregate_type"] for event in deletes} == {"knowledge_meta"} + # 新版自己的同步事件不受影响。 + assert [event["aggregate_id"] + for event in _vector_events(session, "knowledge.vector_sync_requested")] == ["11", "12"] + + +async def test_supersede_happens_before_the_new_rows_are_written() -> None: + """顺序是硬约束:先下线旧版(含投删除事件)、再写新版。 + + 反过来的话,"新行插入失败"会留下一份"旧版已下线、新版没写成"的文档(检索侧彻底查不到); + 代价不对称,所以顺序要在测试里固定下来。 + """ + service, session, _, _ = _service(ids=[11, 12], previous_ids=[3]) + + await service.ingest(filename="faq.md", content=b"x", knowledge_type="faq") + + def kind(index: int) -> str: + sql = session.statements[index][0] + if sql.startswith("UPDATE"): + return "update" + if "fin_knowledge_meta.id" in sql: + return "select_previous" + if sql.startswith("INSERT INTO fin_knowledge_meta"): + return "insert_row" + if sql.startswith("INSERT INTO domain_event_outbox"): + return "insert_event" + return "other" + + kinds = [kind(index) for index in range(len(session.statements))] + assert kinds[0] == "select_previous" + assert kinds.index("update") < kinds.index("insert_row") + # 旧版那条删除事件也在新版第一行之前投出。 + delete_index = next( + index for index, (_, params) in enumerate(session.statements) + if params.get("event_type") == "knowledge.vector_delete_requested" + ) + assert delete_index < kinds.index("insert_row") + + +async def test_first_import_enqueues_no_vector_delete() -> None: + """首次导入(没有旧版)不得投出任何删除事件——否则会把别人的向量删掉。""" + service, session, _, _ = _service(ids=[11, 12], previous_ids=[]) + + await service.ingest(filename="faq.md", content=b"x", knowledge_type="faq") + + assert _vector_events(session, "knowledge.vector_delete_requested") == [] + assert not [sql for sql, _ in session.statements if sql.startswith("UPDATE")] + + +async def test_empty_document_must_not_expire_the_previous_version() -> None: + """空文档(解析出 0 块)不得触发下线:否则"旧版下架了、新版一行没写",文档凭空消失。""" + service, session, _, _ = _service(ids=[], chunks=[], previous_ids=[3, 4]) + + assert await service.ingest(filename="empty.md", content=b"", knowledge_type="faq") == [] + assert session.statements == [] diff --git a/tests/unit/service/test_knowledge_management_service.py b/tests/unit/service/test_knowledge_management_service.py index da40439..8acc03a 100644 --- a/tests/unit/service/test_knowledge_management_service.py +++ b/tests/unit/service/test_knowledge_management_service.py @@ -402,6 +402,59 @@ async def test_archive_key_comes_from_the_knowledge_row() -> None: assert storage.archived == ["kb/faq/deadbeef-faq.md"] +# --- ④ 补投向量清理(expired 行) -------------------------------------------- + + +async def test_cleanup_vector_enqueues_delete_for_an_expired_row() -> None: + """已过期的历史行必须能补投删除事件——这正是 `DELETE` 覆盖不到的那批行。""" + row = meta_row(7, status=EXPIRED_STATUS) + service, session, _, _ = build_service(rows=[row]) + + result = await service.cleanup_vector(context(), 7) + + assert row.status == EXPIRED_STATUS # 不改状态:本动作只补投事件 + events = delete_events(session) + assert len(events) == 1 + assert events[0].aggregate_id == "7" + assert events[0].payload == {"knowledge_id": "7"} + assert events[0].trace_id == "trace-task11" + # 事件 id 必须是新的 uuid4:与历史事件同 id 会被 outbox 的唯一键吞掉,等于没投。 + assert events[0].event_id + assert result == {"knowledge_id": 7, "status": EXPIRED_STATUS, + "vector_delete_event": "knowledge.vector_delete_requested"} + + +async def test_cleanup_vector_refuses_an_active_row() -> None: + """在用文档不得被"只删向量不标状态":那会造成"行是 active、向量已没"的不一致。""" + service, session, _, _ = build_service(rows=[meta_row(7, status="active")]) + + with pytest.raises(ValidationAgentError): + await service.cleanup_vector(context(), 7) + + assert delete_events(session) == [] + + +async def test_cleanup_vector_unknown_id_is_not_found() -> None: + service, session, _, _ = build_service(rows=[meta_row(7, status=EXPIRED_STATUS)]) + + with pytest.raises(GenericResourceNotFoundError): + await service.cleanup_vector(context(), 999) + + assert delete_events(session) == [] + + +async def test_cleanup_vector_requires_the_management_permission() -> None: + service, _, _, _ = build_service(rows=[meta_row(7, status=EXPIRED_STATUS)]) + + def explode() -> None: + raise AssertionError("未授权请求不得打开数据库会话") + + service._session_factory = explode # type: ignore[assignment] + + with pytest.raises(ForbiddenAgentError): + await service.cleanup_vector(context(permissions=()), 7) + + # --- 装配与路由 ------------------------------------------------------------- @@ -420,10 +473,12 @@ def test_router_registers_the_three_required_paths() -> None: assert ("/api/v1/knowledge/upload", "POST") in paths assert ("/api/v1/knowledge/list", "GET") in paths assert ("/api/v1/knowledge/{knowledge_id}", "DELETE") in paths + # 运维补口:给已过期的历史行补投向量清理。 + assert ("/api/v1/knowledge/{knowledge_id}/vector-cleanups", "POST") in paths def test_every_endpoint_depends_on_the_authentication_gate() -> None: - """硬约束:三个端点都必须依赖 `build_request_context`,禁止匿名入口。""" + """硬约束:每个端点都必须依赖 `build_request_context`,禁止匿名入口。""" from app.api.controllers.knowledge_management import router for route in router.routes: diff --git a/tools/purge_expired_knowledge_vectors.py b/tools/purge_expired_knowledge_vectors.py new file mode 100644 index 0000000..6a82b89 --- /dev/null +++ b/tools/purge_expired_knowledge_vectors.py @@ -0,0 +1,158 @@ +"""批量补投「向量残留清理」:把**已过期但向量还在**的知识逐条送去清理(走管理接口)。 + + python tools/purge_expired_knowledge_vectors.py # 只读:列出要清理哪些 + python tools/purge_expired_knowledge_vectors.py --apply # 真投(每条一个 HTTP 调用) + python tools/purge_expired_knowledge_vectors.py --apply --wait 30 # 投完等 Worker 清完再复核 + +## 它解决什么 + +`fin_knowledge_meta.status = 'expired'` 只说明"这行下架了",**不代表 Milvus 里的向量没了**。 +向量删除靠一条 `knowledge.vector_delete_requested` 事件,历史上有一批行是**别的途径**变成 +expired 的(早期脚本直接改库、导入侧幂等上线前的重复副本),它们**从来没投过**这条事件: + +- 检索侧只过滤 `visibility`、**不看 `status`**(`app/service/knowledge_search_service.py`), + 所以这些死向量**继续参与排序**,和同题的活块抢答; +- 而 `DELETE /api/v1/knowledge/{id}` 对已经 expired 的行**刻意返回 404**("不静默成功"), + 于是这些行在管理端口上**无路可走** —— 本脚本 + 新端点 `POST /{id}/vector-cleanups` + 就是补这个缺口。 + +## 与对账工具的关系 + +先用 `python tools/reconcile_knowledge_vectors.py` 看全局,再用本脚本动手。本脚本**只清理 +对账认定"确实还有向量"的那些 id**(`stale_vectors_ids`)—— 对一堆早就干净的 expired 行 +盲投事件,只会往 `domain_event_outbox` 里灌没有意义的删除事件,让"事件堆积"变成常态噪音。 + +## 为什么不直接连 Milvus 删 + +那样会**同时绕过**权限、审计与 Outbox 幂等语义(删除是不可逆动作,必须留痕、必须有唯一入口)。 +孤儿向量(MySQL 里连行都没有的那种)没有可挂事件的载体,本脚本**默认不处理**, +只在对账报告里点出来 —— 真要删属于"直接操作向量库"的运维动作,须单独授权。 +""" + +from __future__ import annotations + +import argparse +import importlib.util +import sys +import time +import uuid +from pathlib import Path +from typing import Any + +import httpx + +if hasattr(sys.stdout, "reconfigure"): + sys.stdout.reconfigure(errors="replace") + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT)) + +BASE_DEFAULT = "http://127.0.0.1:8000" + + +def _load_reconcile_module() -> Any: + """按文件路径加载对账工具(`tools/` 不是包,没有 `__init__.py`,不能 `import tools.x`)。""" + path = ROOT / "tools" / "reconcile_knowledge_vectors.py" + spec = importlib.util.spec_from_file_location("reconcile_knowledge_vectors", path) + if spec is None or spec.loader is None: # pragma: no cover - 只在文件缺失时发生 + raise SystemExit(f"找不到对账工具:{path}") + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +def _headers(token: str) -> dict[str, str]: + return {"Authorization": f"Bearer {token}", "Idempotency-Key": uuid.uuid4().hex} + + +def _login(client: httpx.Client, base: str, username: str, password: str) -> str: + response = client.post( + f"{base}/api/v1/auth/tokens", json={"username": username, "password": password} + ) + if response.status_code != 200: + raise SystemExit(f"登录失败:HTTP {response.status_code} {response.text[:160]}") + return str(response.json()["data"]["access_token"]) + + +def stale_ids_by_collection() -> dict[str, list[str]]: + """本机对账一次,返回「还有向量的已过期行」:`{集合: [knowledge_id, ...]}`。""" + reconcile = _load_reconcile_module() + report = reconcile.reconcile( + reconcile.mysql_rows(), reconcile.milvus_vectors(sorted(_collections())) + ) + return { + name: list(item["detail"]["stale_vectors_ids"]) + for name, item in report.get("collections", {}).items() + if item["detail"]["stale_vectors_ids"] + } + + +def _collections() -> list[str]: + from app.core.knowledge_contracts import ALLOWED_COLLECTIONS + + return sorted(ALLOWED_COLLECTIONS) + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description="补投已过期知识的向量清理(走管理接口)") + parser.add_argument("--base", default=BASE_DEFAULT, help="平台地址") + parser.add_argument("--username", default="admin_t") + parser.add_argument("--password", default="88888888") + parser.add_argument("--apply", action="store_true", help="真正投递(默认只列出要清理的)") + parser.add_argument("--wait", type=int, default=0, metavar="秒", + help="投完等 Worker 消费,再复核一次残留数") + parser.add_argument("--limit", type=int, default=0, help="最多处理多少条(0=不限)") + args = parser.parse_args(argv) + + stale = stale_ids_by_collection() + total = sum(len(ids) for ids in stale.values()) + if total == 0: + print("对账结果:没有任何「已过期但向量还在」的知识行,无需清理。") + print("(要看全局请跑 python tools/reconcile_knowledge_vectors.py)") + return 0 + + print(f"对账发现 {total} 条已过期知识仍有向量(检索侧不看 status,它们仍在参与排序):") + targets: list[str] = [] + for name, ids in stale.items(): + print(f" · {name}: {len(ids)} 条 → {ids[:12]}{' …' if len(ids) > 12 else ''}") + targets.extend(ids) + if args.limit > 0: + targets = targets[: args.limit] + print(f"(--limit {args.limit}:本次只处理前 {len(targets)} 条)") + + if not args.apply: + print("\n[dry-run] 未投递。加 --apply 真投(每条一次 POST /api/v1/knowledge/{id}/vector-cleanups)。") + return 0 + + with httpx.Client(base_url=args.base, timeout=60) as client: + token = _login(client, args.base, args.username, args.password) + ok = 0 + failed: list[tuple[str, str]] = [] + for knowledge_id in targets: + response = client.post( + f"{args.base}/api/v1/knowledge/{knowledge_id}/vector-cleanups", + headers=_headers(token), + ) + if response.status_code == 200: + ok += 1 + else: + failed.append((knowledge_id, f"HTTP {response.status_code} {response.text[:120]}")) + print(f"\n投递完成:成功 {ok} / 失败 {len(failed)}") + for knowledge_id, reason in failed[:10]: + print(f" ✗ {knowledge_id}: {reason}") + + if args.wait > 0: + print(f"等 {args.wait} 秒让 Worker 消费删除事件 …") + time.sleep(args.wait) + remaining = stale_ids_by_collection() + left = sum(len(ids) for ids in remaining.values()) + print(f"复核:仍有向量的已过期行 {left} 条" + + ("" if left == 0 else f" → { {k: v[:8] for k, v in remaining.items()} }")) + if left > 0: + print(" 提示:确认 Worker 在跑(python -m app.worker),并看 domain_event_outbox 里" + " knowledge.vector_delete_requested 是否卡在 failed/dead。") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tools/reconcile_knowledge_vectors.py b/tools/reconcile_knowledge_vectors.py new file mode 100644 index 0000000..4bbd44e --- /dev/null +++ b/tools/reconcile_knowledge_vectors.py @@ -0,0 +1,373 @@ +"""知识库「向量 ↔ 元数据」对账(**只读**,不写 Milvus、不改 MySQL)。 + + python tools/reconcile_knowledge_vectors.py # 人读摘要 + python tools/reconcile_knowledge_vectors.py --json # 机器读(含全部 id 明细) + python tools/reconcile_knowledge_vectors.py --out docs/evidence/xxx.json + +## 为什么要这个工具 + +知识的**元数据在 MySQL**(`fin_knowledge_meta`),**正文与向量在 Milvus**,两侧靠 +`id`(Milvus 侧的物理字段名是 `doc_id` 或 `knowledge_id`,因环境而异,**运行时探测**)关联。 +这个关联**没有任何数据库层约束**,因此会以四种方式坏掉,而**每一种在接口层都不报错**: + +| 分类 | 现象 | 为什么会发生 | +|---|---|---| +| `orphan_vectors` | 向量在 Milvus、MySQL 里没有这个 id | 知识行被物理删除过(早期脚本) | +| `stale_vectors` | 向量在 Milvus、MySQL 行是 `expired` | 行是"别的途径"变成 expired 的,**从没投过删除事件** —— 这正是管理端口 `POST /api/v1/knowledge/{id}/vector-cleanups` 要清理的那批 | +| `missing_vectors` | MySQL 行是 `active`、Milvus 里没有向量 | 同步事件没被消费(Worker 没跑 / 事件 `dead`),或向量集合被重建过 | +| `duplicate_groups` | 同一 `source_file` + 集合有多条 active 行 / 正文完全相同的多行 | 导入侧幂等上线前重复入库;手工灌库 | + +**为什么必须查这个而不是"看检索结果对不对"**:检索侧只过滤 `visibility`、**不看 `status`**, +所以 `stale_vectors` 会**继续参与排序**并和其他副本抢答 —— 表现出来只是"客服答得不好/总转人工", +从任何一个接口都看不出向量库里多了一堆死向量(实测:产品手册被重复入库 7 次, +175 条历史副本把「风险等级 R1–R5」这类问题卡在"领先次优不够"的门槛下)。 + +## 口径(三条,缺一条结论就会误导人) + +1. **非数字 id 单独一类(`unbacked_seed_vectors`)**:`faq` / `policy` 两个集合里有一批 + `FAQ-0013` 之类的**人工语义 id**,它们**本来就不该出现在 MySQL**(由 + `tools/load_knowledge_milvus.py` / 灌库脚本直接写 Milvus)。把它们算成 `orphan_vectors` + 会让人"清理孤儿"时把**唯一正确的答案**删掉 —— 本工具单独列出,且**不建议删**。 +2. **只处理三个知识集合**:长期记忆向量(`user_long_term_memory_v1`)不属于知识库, + 不在这里的统计范围内。 +3. **只读**:本工具不修任何东西。发现 `stale_vectors` 后,用管理端口的 vector-cleanups + 补投事件(见 `tools/purge_expired_knowledge_vectors.py`),由 Worker 真正删除 —— + 工具自己去删向量会绕过权限、审计与 Outbox 语义。 +""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import sys +from collections import defaultdict +from pathlib import Path +from typing import Any + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT)) + +from app.core.knowledge_contracts import ALLOWED_COLLECTIONS # noqa: E402 + +#: Milvus `query` 的分页上限(各环境默认 16384);超过就分页取。 +PAGE_SIZE = 16384 + +#: 「低信息量向量」的正文长度门槛:一个向量整篇正文不到这么长,不可能回答任何问题, +#: 却照样进排序。实测 `fin_policy_collection` 里有 20 条正文都是 19 字的 +#: 「第九条 问卷内容及评分标准:选项 分值」(灌库时把 markdown 表格切碎了) +#: —— 它们对任何"问卷/评分"类提问都是**同分并列**,直接把客服的"领先次优 ≥0.07"卡死。 +TINY_CONTENT_LENGTH = 40 + +#: 默认的只读证据落盘位置(与人读摘要无关,摘要始终打到 stdout)。 +DEFAULT_OUT = Path("docs/evidence/knowledge-vectors-reconcile.json") + + +def is_heading_only(content: str) -> bool: + """整段正文只是一行 markdown 标题(`## 一、交易规则` 这种)。 + + 这类向量**不可能回答任何问题**,却会在任何提到这几个字的提问上拿到不低的相似度 + —— 它把"领先次优"的差值抹平,也让真正有内容的块排不到第一。 + 实测 `fin_product_collection` 里有 47 条正文不到 20 字、其中不少就是标题行, + 来源是灌库脚本按行切 markdown 时把标题单独切成了块。 + """ + stripped = content.strip() + if not stripped.startswith("#"): + return False + return len(stripped.lstrip("#").strip()) < 20 + + +def mysql_rows() -> list[dict[str, Any]]: + """读 `fin_knowledge_meta` 的全部行(只取对账需要的列,**不读正文**)。 + + 不读 `content_text`:这份表里存着整篇正文,全表拉回来既慢又没必要 —— + 正文重复度由 Milvus 侧算(那里本来就有正文/摘要字段)。 + """ + from urllib.parse import unquote, urlparse + + import pymysql # type: ignore[import-untyped] + + from app.core.config import get_settings + + parsed = urlparse(get_settings().mysql_dsn.replace("mysql+asyncmy://", "mysql+pymysql://")) + connection = pymysql.connect( + host=parsed.hostname or "127.0.0.1", + port=parsed.port or 3306, + user=unquote(parsed.username or ""), + password=unquote(parsed.password or ""), + database=(parsed.path or "/").lstrip("/"), + charset="utf8mb4", + cursorclass=pymysql.cursors.DictCursor, + ) + try: + with connection.cursor() as cursor: + cursor.execute( + "SELECT id, knowledge_type, source_file, milvus_collection, version, status," + " review_status, created_at FROM fin_knowledge_meta" + ) + return [dict(row) for row in cursor.fetchall()] + finally: + connection.close() + + +def milvus_vectors(collections: list[str]) -> dict[str, Any]: + """按集合取 (向量 id, 正文哈希) 明细,物理字段名**运行时探测**。 + + 返回 `{"collections": {...}, "errors": {...}}`;单集合失败不中断其余 + (与检索侧一贯口径一致:一个集合探测不了,不该让整份报告消失)。 + """ + from app.core.config import get_settings + from app.core.knowledge_schema import detect_schema + from pymilvus import MilvusClient # type: ignore[import-untyped] + + settings = get_settings() + client = MilvusClient(uri=settings.milvus_uri, token=settings.milvus_token or None) + report: dict[str, Any] = {"collections": {}, "errors": {}} + for name in collections: + try: + schema = detect_schema(client, name) + if not schema.usable: + report["errors"][name] = ( + schema.error or f"缺少必需字段:{','.join(schema.missing_required)}" + ) + continue + id_field = schema.resolve("doc_id") + content_field = schema.resolve("content") + assert id_field is not None and content_field is not None # schema.usable 已保证 + entries: list[dict[str, Any]] = [] + offset = 0 + while True: + rows = client.query( + collection_name=name, + filter="", + output_fields=[id_field, content_field], + limit=PAGE_SIZE, + offset=offset, + ) + if not rows: + break + for row in rows: + content = str(row.get(content_field) or "") + entries.append({ + "vector_id": str(row.get(id_field) or ""), + "content_sha1": hashlib.sha1( + content.encode("utf-8") + ).hexdigest()[:12], + "content_length": len(content), + "heading_only": is_heading_only(content), + }) + if len(rows) < PAGE_SIZE: + break + offset += len(rows) + report["collections"][name] = { + "id_field": id_field, + "content_field": content_field, + "entries": entries, + } + except Exception as exc: # 单集合失败不影响其余 + report["errors"][name] = f"{type(exc).__name__}: {exc}" + return report + + +def reconcile( + rows: list[dict[str, Any]], vectors: dict[str, Any] +) -> dict[str, Any]: + """把两侧数据算成一份对账结论(纯函数:不连库、不连 Milvus,便于单测)。 + + `rows` 是 `fin_knowledge_meta` 的行;`vectors` 是 `milvus_vectors()` 的返回。 + """ + by_collection: dict[str, list[dict[str, Any]]] = defaultdict(list) + for row in rows: + by_collection[str(row.get("milvus_collection") or "")].append(row) + + result: dict[str, Any] = {"collections": {}, "errors": dict(vectors.get("errors", {}))} + for name, payload in vectors.get("collections", {}).items(): + entries = payload["entries"] + mysql_for_collection = by_collection.get(name, []) + status_by_id = {int(row["id"]): str(row.get("status") or "") for row in mysql_for_collection} + vector_ids = [entry["vector_id"] for entry in entries] + + numeric_ids = [vid for vid in vector_ids if vid.isdigit()] + unbacked = [vid for vid in vector_ids if not vid.isdigit()] + orphan = [vid for vid in numeric_ids if int(vid) not in status_by_id] + stale = [ + vid for vid in numeric_ids + if status_by_id.get(int(vid), "") not in {"", "active"} + ] + active_ids = { + int(row["id"]) for row in mysql_for_collection + if str(row.get("status") or "") == "active" + } + missing = sorted(active_ids - {int(vid) for vid in numeric_ids}) + + # 正文完全相同的向量(同一份内容被灌了多次)——**不管状态**:检索侧不看 status, + # 所以一条 expired 行的向量照样参与排序,和别人的副本抢答。 + sha_groups: dict[str, list[str]] = defaultdict(list) + for entry in entries: + sha_groups[entry["content_sha1"]].append(entry["vector_id"]) + duplicate_content = { + sha: ids for sha, ids in sha_groups.items() if len(ids) > 1 + } + + # **同一份文档还有几个 active 版本**:只看行数会误判 —— 一份文档正常会切成十几个 chunk + # (实测 `场内基金产品手册.md` 一个版本就是 24 行)。真正的重复信号是 + # **该文件的 active 行里出现了正文完全相同的两块**:那只能是同一份文档被入库了多次。 + sha_by_vector = {entry["vector_id"]: entry["content_sha1"] for entry in entries} + active_files: dict[str, list[str]] = defaultdict(list) + for row in mysql_for_collection: + if str(row.get("status") or "") != "active": + continue + active_files[str(row.get("source_file") or "")].append( + sha_by_vector.get(str(int(row["id"])), "") + ) + duplicated_active_files = { + file: sorted({sha for sha in shas if sha and shas.count(sha) > 1}) + for file, shas in active_files.items() + if any(sha and shas.count(sha) > 1 for sha in shas) + } + + # 低信息量向量:正文短到不可能回答任何问题,却照样参与排序(见 TINY_CONTENT_LENGTH)。 + tiny = [ + entry["vector_id"] for entry in entries + if entry["content_length"] < TINY_CONTENT_LENGTH + ] + heading_only = [ + entry["vector_id"] for entry in entries if entry["heading_only"] + ] + + result["collections"][name] = { + "milvus_vectors": len(entries), + "mysql_rows": len(mysql_for_collection), + "mysql_active": len(active_ids), + "active_files": len(active_files), + "orphan_vectors": len(orphan), + "stale_vectors": len(stale), + "missing_vectors": len(missing), + "unbacked_seed_vectors": len(unbacked), + "duplicate_content_groups": len(duplicate_content), + "duplicate_content_vectors": sum(len(ids) for ids in duplicate_content.values()), + "tiny_vectors": len(tiny), + "heading_only_vectors": len(heading_only), + "duplicated_active_files": len(duplicated_active_files), + "detail": { + "orphan_vectors": orphan[:200], + "stale_vectors_ids": stale[:200], + "missing_vectors": missing[:200], + "unbacked_seed_vectors": unbacked[:200], + "tiny_vectors": tiny[:100], + "heading_only_vectors": heading_only[:100], + "duplicated_active_files": { + file: shas[:20] for file, shas in list(duplicated_active_files.items())[:50] + }, + "duplicate_content_groups": { + sha: ids[:20] for sha, ids in list(duplicate_content.items())[:50] + }, + }, + } + return result + + +def summarize(report: dict[str, Any]) -> str: + """人读摘要。每条都带"下一步该做什么",避免只给一堆数字。""" + lines: list[str] = ["知识库「向量 ↔ 元数据」对账(只读)", ""] + header = ( + f"{'集合':<24}{'向量':>7}{'MySQL行':>9}{'active':>8}{'文档数':>8}" + f"{'孤儿':>7}{'死向量':>8}{'缺向量':>8}{'无元数据种子':>13}{'重复正文':>9}{'碎片':>7}{'纯标题':>7}" + ) + lines.append(header) + lines.append("-" * len(header)) + for name, item in report.get("collections", {}).items(): + lines.append( + f"{name:<24}{item['milvus_vectors']:>7}{item['mysql_rows']:>9}" + f"{item['mysql_active']:>8}{item.get('active_files', 0):>8}" + f"{item['orphan_vectors']:>7}" + f"{item['stale_vectors']:>8}{item['missing_vectors']:>8}" + f"{item['unbacked_seed_vectors']:>13}" + f"{item.get('duplicate_content_vectors', 0):>9}{item.get('tiny_vectors', 0):>7}" + f"{item.get('heading_only_vectors', 0):>7}" + ) + for name, error in report.get("errors", {}).items(): + lines.append(f"⚠️ {name} 探测失败:{error}") + + lines.append("") + totals = report.get("collections", {}).values() + stale_total = sum(item["stale_vectors"] for item in totals) + orphan_total = sum(item["orphan_vectors"] for item in totals) + missing_total = sum(item["missing_vectors"] for item in totals) + dup_files = sum(item.get("duplicated_active_files", 0) for item in totals) + dup_vectors = sum(item.get("duplicate_content_vectors", 0) for item in totals) + tiny_total = sum(item.get("tiny_vectors", 0) for item in totals) + heading_total = sum(item.get("heading_only_vectors", 0) for item in totals) + lines.append("怎么读这几列:") + lines.append( + f" · 死向量 {stale_total} 条:MySQL 里已是 expired,向量却还在,**会继续参与检索排序**。" + ) + lines.append(" 处理:python tools/purge_expired_knowledge_vectors.py --apply(走管理端口补投删除事件)") + lines.append( + f" · 孤儿向量 {orphan_total} 条:MySQL 里连行都没有(历史物理删除的残留),处理同「死向量」。" + ) + lines.append( + f" · 缺向量 {missing_total} 条:active 行没有向量 → 这份知识**检索永远命中不到**。" + ) + lines.append(" 处理:确认 Worker 在跑,再重传一次该文档让同步事件重投(最省事)。") + lines.append( + f" · 重复正文 {dup_vectors} 条向量:正文逐字相同的多份,检索时**必然互相打平**," + "是「领先次优 ≥0.07」被卡死的直接来源。" + ) + lines.append( + f" (其中「同一份文档还有多个 active 版本」{dup_files} 份 → 重传一次该文件即可收敛)" + ) + lines.append( + f" · 碎片 {tiny_total} 条:整篇正文不到 {TINY_CONTENT_LENGTH} 字的向量,回答不了任何问题," + "却照样进排序、并和同题其它块打平(实测来源是灌库时把 markdown 表格切碎)。" + ) + lines.append(" 处理:属**灌库脚本的切分缺陷**,要在灌库侧修;不要为此改检索层的阈值。") + lines.append( + f" · 纯标题 {heading_total} 条:整段正文就是一行 markdown 标题(`## 一、交易规则`)。" + " 这类向量**不可能回答任何问题**,却会在提到这几个字的提问上拿到不低的相似度。" + ) + lines.append( + " 「文档数」列是一个集合里有多少份 active 文档(不是行数):一份文档正常切成十几块。" + ) + lines.append( + " · 无元数据种子向量:`FAQ-0013` 这类人工语义 id,**本来就不在 MySQL**,不要当孤儿清理。" + ) + lines.append( + " · 「向量」列来自 Milvus `query` 的**实际行数**;`get_collection_stats` 的 row_count" + " 会把已软删、尚未 compaction 的行也算进去,故两者不等是正常的,以本列为准。" + ) + return "\n".join(lines) + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description="知识库向量-元数据对账(只读)") + parser.add_argument("--json", action="store_true", help="打印完整 JSON(含 id 明细)") + parser.add_argument("--out", type=Path, default=DEFAULT_OUT, help="证据落盘路径") + parser.add_argument("--no-write", action="store_true", help="不落盘,只打印") + args = parser.parse_args(argv) + + rows = mysql_rows() + vectors = milvus_vectors(sorted(ALLOWED_COLLECTIONS)) + report = reconcile(rows, vectors) + report["mysql_rows_total"] = len(rows) + report["collections_checked"] = sorted(ALLOWED_COLLECTIONS) + + if not args.no_write: + args.out.parent.mkdir(parents=True, exist_ok=True) + args.out.write_text( + json.dumps(report, ensure_ascii=False, indent=2), encoding="utf-8" + ) + + if args.json: + print(json.dumps(report, ensure_ascii=False, indent=2)) + else: + print(summarize(report)) + if not args.no_write: + print(f"\n明细已写入 {args.out}") + + # 退出码只表达"能不能判定",不表达"有没有问题":对账发现残留是常态,不该让 CI 红。 + return 1 if report.get("errors") else 0 + + +if __name__ == "__main__": + raise SystemExit(main())