Files
group_fqcd_jr/tools/import_knowledge_seed.py

320 lines
15 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""把已审核的 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())