"""把已审核的 QA 源文件一次性导入 `fin_knowledge_meta` + `agent_faq_synonym`,并投递向量同步事件。 为什么是一个独立的种子脚本、而不是"跑一次 SQL":这 105 条(+1 条验收补充)是客服 Agent RAG 检索的**全部**知识来源。它们必须能被重复导入——Task 2 改治理层、Task 4 建文档解析、 Task 5 写 Outbox Worker 都可能要重放这次导入;一旦脚本不幂等,重放就会产生重复知识、 重复同义问法,并在 Milvus 里留下同一段文本的多个向量(检索结果会重复命中同一条)。 幂等口径(控制者裁定 4): 1. 知识行按 `tags.qa_id` 定位——存在则 UPDATE、不存在则 INSERT。`qa_id` 是源文件的稳定编号, 不依赖自增主键,因此重复导入不会新增行。 2. 同义问法先按 `(knowledge_id, source_type='import')` 删除再插入:`import` 来源的整批 由本脚本独占,删除-重插能得到源文件的精确镜像;手工/会话/工单来源的同义问法不受影响。 3. 去重按 `phrase_hash`(SHA-256(归一化问法))而非原文字符串——`uk_faq_synonym (knowledge_id, phrase_hash)` 就是这么定义的,且「问题」本身也作为一个 phrase 写入 (控制者裁定 5),所以标准问法与相似问法之间的重复也在同一套哈希下去重。 本脚本**不写 Milvus**:向量由 `knowledge.vector_sync_requested` 事件的 Worker(Task 5)落库。 这里只负责把事件投进 `domain_event_outbox`。 写入的固定值(控制者裁定 1/2/3,全部有据可查): - `milvus_collection='fin_faq_collection'`(与 Task 6 的三集合白名单一致)。 - `effective_date`/`expire_date` 留空 = 长期有效;读路径对空值按"长期有效"处理。 - `review_status='published'` + `status='active'`(源文件本身即 `approved_candidate` 且已人工审核); 同义问法 `status='approved'` + `source_type='import'`(`docs/02` §7.3 DDL 的合法枚举,带 CHECK)。 - `created_by`/`reviewer_id` = `sys_user.id`(现库 `9003`)——`agent_faq_synonym.created_by` 非空且带外键 `fk_synonym_created_by`,写不存在的用户会直接失败。 执行(可重复执行): .\\.venv\\Scripts\\python.exe tools\\import_knowledge_seed.py \\ --source "C:\\...\\客服Agent知识库_QA问答对_v5_RAG发布候选版.txt" --created-by 9003 """ from __future__ import annotations import argparse import asyncio import json import sys from datetime import UTC, datetime from pathlib import Path from uuid import uuid4 ROOT = Path(__file__).resolve().parents[1] if str(ROOT) not in sys.path: sys.path.insert(0, str(ROOT)) from sqlalchemy import text # noqa: E402 from sqlalchemy.ext.asyncio import AsyncSession # noqa: E402 from app.infrastructure.db import SessionFactory # noqa: E402 from tools.qa_source_parser import ( # noqa: E402 QaRecord, normalize_phrase, parse_qa_source, phrase_hash, ) #: 三个 Milvus 集合里的 FAQ 集合(Task 6 的白名单口径)。 FAQ_COLLECTION = "fin_faq_collection" #: `fin_knowledge_meta.knowledge_type`:FAQ 文本知识。 KNOWLEDGE_TYPE = "faq" #: 源文件版本,写入 `version` 与 `tags.source_version`,便于回溯"这条来自哪一版"。 SOURCE_VERSION = "v5.8" #: 向量同步事件类型(Task 5 的 Worker 消费它)。 VECTOR_SYNC_EVENT = "knowledge.vector_sync_requested" #: 本脚本独占的 `agent_faq_synonym.source_type`;重导时只删这一来源的行。 IMPORT_SOURCE_TYPE = "import" #: 老师 Phase 1 验收点名的测试问题「基金申购后多久确认」**不在 105 条内** #: (源文件里只有 `RAG-CONFIG-006`「申购和赎回什么时候确认、到账?」,措辞不同), #: 这里补一条同义的验收条目,使验收用例能真正命中知识库而不是落到兜底话术。 ACCEPTANCE_SUPPLEMENT = QaRecord( qa_id="SUP-001", question="基金申购后多久确认?", synonyms=("申购多久确认份额?", "买基金几天能确认?", "基金申购确认时间"), answer=( "交易日 15:00 前提交的申购申请,T+1 日确认份额(QDII 基金为 T+2 日);" "15:00 后提交则顺延至下一交易日。非交易日(周末及法定节假日)提交的申请顺延至" "下一交易日处理。确认后即可查看持仓。具体以产品说明书和交易页面显示为准。" ), ) SELECT_KNOWLEDGE_ID_BY_QA_ID = text( "SELECT id FROM fin_knowledge_meta" " WHERE JSON_UNQUOTE(JSON_EXTRACT(tags, '$.qa_id')) = :qa_id" " ORDER BY id LIMIT 1" ) INSERT_KNOWLEDGE_META = text( "INSERT INTO fin_knowledge_meta" " (knowledge_type, title, source_file, milvus_collection, version, content_text," " tags, reviewer_id, review_status, status, created_at, updated_at)" " VALUES (:knowledge_type, :title, :source_file, :collection, :version, :content_text," " :tags, :reviewer_id, 'published', 'active', :now, :now)" ) UPDATE_KNOWLEDGE_META = text( "UPDATE fin_knowledge_meta SET title=:title, source_file=:source_file," " milvus_collection=:collection, version=:version, content_text=:content_text," " tags=:tags, reviewer_id=:reviewer_id, review_status='published', status='active'," " updated_at=:now WHERE id=:id" ) SELECT_LAST_INSERT_ID = text("SELECT LAST_INSERT_ID()") DELETE_IMPORTED_SYNONYMS = text( "DELETE FROM agent_faq_synonym WHERE knowledge_id=:knowledge_id" " AND source_type=:source_type" ) INSERT_SYNONYM = text( "INSERT INTO agent_faq_synonym" " (knowledge_id, phrase, normalized_phrase, phrase_hash, language_code, source_type," " hit_count, status, created_by, reviewer_id, reviewed_at, created_at, updated_at)" " VALUES (:knowledge_id, :phrase, :normalized_phrase, :phrase_hash, 'zh-CN'," " :source_type, 0, 'approved', :created_by, :created_by, :now, :now, :now)" ) 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, 'knowledge_meta', :aggregate_id, :trace_id, :payload," " 'pending', 0, :now, :now, :now)" ) def _knowledge_tags(record: QaRecord) -> str: return json.dumps( {"qa_id": record.qa_id, "phase": "phase_1", "source_version": SOURCE_VERSION}, ensure_ascii=False, ) async def upsert_knowledge_meta( session: AsyncSession, record: QaRecord, *, source_file: str, created_by: int, now: datetime, ) -> int: """按 `tags.qa_id` 幂等写入知识行,返回 `fin_knowledge_meta.id`。""" params = { "knowledge_type": KNOWLEDGE_TYPE, "title": record.question, "source_file": source_file, "collection": FAQ_COLLECTION, "version": SOURCE_VERSION, "content_text": record.answer, "tags": _knowledge_tags(record), "reviewer_id": created_by, "now": now, } existing = await session.scalar(SELECT_KNOWLEDGE_ID_BY_QA_ID, {"qa_id": record.qa_id}) if existing is None: await session.execute(INSERT_KNOWLEDGE_META, params) inserted = await session.scalar(SELECT_LAST_INSERT_ID) if inserted is None: # pragma: no cover - 只在非 MySQL 后端上发生 raise RuntimeError(f"{record.qa_id} 插入后取不到 LAST_INSERT_ID()") return int(inserted) knowledge_id = int(existing) await session.execute(UPDATE_KNOWLEDGE_META, {**params, "id": knowledge_id}) return knowledge_id async def replace_imported_synonyms( session: AsyncSession, knowledge_id: int, record: QaRecord, *, created_by: int, now: datetime, ) -> int: """重建该知识的 `import` 来源同义问法(含标准问法本身),返回写入行数。""" await session.execute( DELETE_IMPORTED_SYNONYMS, {"knowledge_id": knowledge_id, "source_type": IMPORT_SOURCE_TYPE}, ) digests: set[str] = set() written = 0 for phrase in (record.question, *record.synonyms): digest = phrase_hash(phrase) if digest in digests: continue digests.add(digest) await session.execute(INSERT_SYNONYM, { "knowledge_id": knowledge_id, "phrase": phrase, "normalized_phrase": normalize_phrase(phrase), "phrase_hash": digest, "source_type": IMPORT_SOURCE_TYPE, "created_by": created_by, "now": now, }) written += 1 return written async def enqueue_vector_sync( session: AsyncSession, knowledge_id: int, *, now: datetime, ) -> None: """投递向量同步事件;Milvus 的写入由 Task 5 的 Outbox Worker 负责。""" await session.execute(INSERT_OUTBOX_EVENT, { "event_id": str(uuid4()), "event_type": VECTOR_SYNC_EVENT, "aggregate_id": str(knowledge_id), "trace_id": str(uuid4()), "payload": json.dumps({"knowledge_id": str(knowledge_id)}), "now": now, }) async def import_records( records: list[QaRecord], *, source_file: str, created_by: int, now: datetime, ) -> tuple[int, int]: """单事务导入全部记录,返回 `(知识行数, 同义问法行数)`。""" synonyms_written = 0 async with SessionFactory() as session, session.begin(): for record in records: knowledge_id = await upsert_knowledge_meta( session, record, source_file=source_file, created_by=created_by, now=now, ) synonyms_written += await replace_imported_synonyms( session, knowledge_id, record, created_by=created_by, now=now, ) await enqueue_vector_sync(session, knowledge_id, now=now) return len(records), synonyms_written async def verify(expected_knowledge: int, expected_synonyms: int) -> None: """按读路径的**同一口径**复查落库结果:任一项不符就 `SystemExit`(非零退出码)。 四项断言,全部是"== 期望值":知识行数、同义问法行数、被向量同步事件覆盖的 `aggregate_id` 数、验收补充条目可查。**不打印了事**——见函数体内的说明。 """ async with SessionFactory() as session: published = await session.scalar(text( "SELECT COUNT(*) FROM fin_knowledge_meta" " WHERE review_status='published' AND status='active'" " AND milvus_collection=:collection" " AND (effective_date IS NULL OR effective_date <= UTC_DATE())" " AND (expire_date IS NULL OR expire_date > UTC_DATE())" ), {"collection": FAQ_COLLECTION}) imported_synonyms = await session.scalar(text( "SELECT COUNT(*) FROM agent_faq_synonym" " WHERE source_type=:source_type AND status='approved'" ), {"source_type": IMPORT_SOURCE_TYPE}) events = await session.scalar(text( "SELECT COUNT(*) FROM domain_event_outbox WHERE event_type=:event_type" ), {"event_type": VECTOR_SYNC_EVENT}) covered = await session.scalar(text( "SELECT COUNT(DISTINCT aggregate_id) FROM domain_event_outbox" " WHERE event_type=:event_type" ), {"event_type": VECTOR_SYNC_EVENT}) supplement = await session.scalar(text( "SELECT id FROM fin_knowledge_meta WHERE title=:title" " AND knowledge_type=:knowledge_type AND milvus_collection=:collection" " AND review_status='published' AND status='active'" ), {"title": ACCEPTANCE_SUPPLEMENT.question, "knowledge_type": KNOWLEDGE_TYPE, "collection": FAQ_COLLECTION}) # 断言而不是打印:`verify()` 存在的意义就是"每轮导入都确认真的能被读路径查到"。 # 只打印不断言的话,一轮投递 0 个事件、知识行被别的流程删掉,脚本照样打印 # `verified:` 并退出 0——「静默失效」比报错难查得多(本 Task 审查时真的发生过: # 底座配置发布测试的 `finally` 清理删 Outbox 时漏了 `event_type` 条件, # `aggregate_id` 是字符串列而 release_id 也是小整数,数值碰撞把知识事件误删了)。 print(f" published+active+有效 知识行: {published}(本次导入 {expected_knowledge} 条)") print(f" approved+import 同义问法行: {imported_synonyms}(本次写入 {expected_synonyms} 条)") print(f" {VECTOR_SYNC_EVENT} 事件行: {events},覆盖 {covered} 个 aggregate_id" f"(本次投递 {expected_knowledge} 条)") # 用 `==` 而不是 `>=`:多出来的行同样是坏消息(别的流程往同一集合里写了不在 # 本次导入范围内的知识),放过去等于把这个集合的"精确镜像"性质丢掉。 if published is None or int(published) != expected_knowledge: raise SystemExit(f"知识行不符:期望 == {expected_knowledge},实际 {published}") if imported_synonyms is None or int(imported_synonyms) != expected_synonyms: raise SystemExit(f"同义问法行不符:期望 {expected_synonyms},实际 {imported_synonyms}") # 本脚本对 `import_records` 的每条知识都投一条事件(见 `import_records` 循环), # 所以"事件覆盖的 aggregate_id 数 == 知识行数"在本设计下必须成立: # 少一个就说明有知识行没进入向量同步队列(就是那个被误删的场景)。 if covered is None or int(covered) != expected_knowledge: raise SystemExit( f"向量同步事件未覆盖全部知识行:期望 {expected_knowledge} 个 aggregate_id," f"实际 {covered}" ) if supplement is None: raise SystemExit( f"验收补充条目查不到:{ACCEPTANCE_SUPPLEMENT.question}" f"(knowledge_type={KNOWLEDGE_TYPE}, collection={FAQ_COLLECTION}, published+active)" ) print(f" verified: 验收补充条目 id={supplement} 可查;" f"{expected_knowledge} 个知识行全部有向量同步事件") async def main() -> None: parser = argparse.ArgumentParser(description="导入已审核 QA 知识库(幂等,可重复执行)") parser.add_argument("--source", required=True, type=Path, help="已审核 QA 源文件") parser.add_argument("--created-by", required=True, type=int, help="sys_user.id(外键)") parser.add_argument("--dry-run", action="store_true", help="只解析、不写库") args = parser.parse_args() records = parse_qa_source(args.source.read_text(encoding="utf-8")) records.append(ACCEPTANCE_SUPPLEMENT) print(f"parsed {len(records)} records" f" (含 1 条验收补充:{ACCEPTANCE_SUPPLEMENT.qa_id})") if args.dry_run: print("dry-run: 未写库") return now = datetime.now(UTC).replace(tzinfo=None) knowledge_count, synonym_count = await import_records( records, source_file=str(args.source), created_by=args.created_by, now=now, ) print(f"imported {knowledge_count} knowledge rows, {synonym_count} synonym rows;" f" {knowledge_count} vector sync events queued") await verify(knowledge_count, synonym_count) if __name__ == "__main__": asyncio.run(main())