320 lines
15 KiB
Python
320 lines
15 KiB
Python
"""把已审核的 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())
|