refactor(rag): extract retrieval and ingest to py-rag service

- replace in-process Milvus stack with PyRagClient + PyRagKnowledgeSearchAdapter behind KnowledgeSearchPort (RERANK score passthrough)
- move document ingest to py-rag /documents:ingest; DocumentManagementService keeps MySQL ledger + local files
- sink L0 query understanding to py-rag; drop KnowledgeQueryTransformer, single UNFILTERED_VECTOR attempt
- remove Milvus deps, config classes, dead demo services and obsolete rebuild scripts
- compose/Makefile reduced to MySQL/Redis; add pyrag.* config
This commit is contained in:
zhuyongxin
2026-09-30 17:03:21 +08:00
parent 83193bdf4a
commit 9cf162482d
74 changed files with 837 additions and 8285 deletions
-106
View File
@@ -1,106 +0,0 @@
# 重建 hybrid 知识库(dense + BM25)
面向当前 `knowledge_base/` 目录文档,**清空并重建**配置中的 Milvus collection(默认 **`biz`**)。
## 前提
1. 应用已启动(默认 `http://localhost:9900`)
2. `MILVUS_TOKEN` 等连接配置可用
3. `application.yml` 已配置:
```yaml
milvus:
collection: biz
retrieval:
search:
mode: hybrid
knowledge:
base-path: knowledge_base/
```
## 一键脚本(Python)
在项目根目录执行:
```bash
python scripts/rebuild_hybrid_knowledge.py --confirm REBUILD
```
指定服务地址:
```bash
python scripts/rebuild_hybrid_knowledge.py --base-url http://127.0.0.1:9900 --confirm REBUILD
```
跳过前后 stats:
```bash
python scripts/rebuild_hybrid_knowledge.py --confirm REBUILD --skip-stats
```
依赖:Python 3.9+ 标准库即可(无需 pip 包)。
## 脚本会做什么
| 步骤 | 动作 |
|---|---|
| 1 | 检查 `/milvus/health` |
| 2 | 打印重建前 `/api/knowledge/stats` |
| 3 | `POST /api/knowledge/rebuild-hybrid?confirm=REBUILD` |
| 4 | 打印重建后 stats |
服务端 `rebuild-hybrid` 内部顺序:
1. **Drop + recreate** Milvus collection(`milvus.collection`,默认 `biz`)
- 原有向量数据会被删除
- 按 dense + BM25 schema 重建
2. **清空** MySQL `api_document`
3. **清空** 内存 L0 索引
4. **扫描** `knowledge_base/**/*.md`(跳过 `README.md`)并 force 全量导入
- 写 MySQL 元数据
- 切片
- 写 dense 向量 + BM25 `search_text`
- 更新 L0
## 不会做什么
- **不会**动 `knowledge_base/` 源文件
- **不会**在未传 `--confirm REBUILD` 时执行
## 手动 curl 等价命令
```bash
# 重建(危险:会清空 biz collection + api_document)
curl -X POST "http://localhost:9900/api/knowledge/rebuild-hybrid?confirm=REBUILD"
# 仅强制导入(不 drop collection)
curl -X POST "http://localhost:9900/api/knowledge/init?force=true"
# 统计
curl "http://localhost:9900/api/knowledge/stats"
```
## 成功判据
响应中大致应有:
```json
{
"success": true,
"collection": "biz",
"inserted": 15,
"failed": 0,
"milvus": { "recreated": true, "loaded": true }
}
```
然后用一条知识库里真实存在的术语/故障词走 `lookup_knowledge` 或 chat 验证 hybrid 命中。
## 失败排查
| 现象 | 可能原因 |
|---|---|
| connect / token 错误 | `MILVUS_TOKEN`、host、database |
| BM25 / analyzer 相关报错 | 云端 Milvus/Zilliz 版本不支持 BM25 Function |
| inserted=0 | `knowledge_base` 路径不对,或 md 缺 frontmatter/title |
| failed>0 | 看响应 `details` 与应用日志 |
-292
View File
@@ -1,292 +0,0 @@
#!/usr/bin/env python3
"""Live acceptance runner for post-reindex RAG retrieval checks.
This script calls the running Spring Boot retrieval endpoint. It is intentionally
separate from the offline fixture baseline because it depends on live service and
Milvus/Zilliz state.
"""
from __future__ import annotations
import argparse
import json
import sys
import urllib.error
import urllib.parse
import urllib.request
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
DEFAULT_BASE_URL = "http://127.0.0.1:9900"
DEFAULT_JSON_REPORT = Path("eval/rag-retrieval/reports/live-post-reindex.json")
DEFAULT_MD_REPORT = Path("eval/rag-retrieval/reports/live-post-reindex.md")
DEFAULT_CASES: list[dict[str, Any]] = [
{
"caseId": "breadcrumb-rag-chunk-context",
"query": "If a long RAG section is split into multiple chunks, how do we keep retrieval context?",
"topK": 5,
"purpose": "Breadcrumb-sensitive RAG chunk context retrieval.",
},
{
"caseId": "breadcrumb-diagnosis-flow",
"query": "What is the standard troubleshooting flow for an application incident?",
"topK": 5,
"purpose": "Process-style retrieval where section path matters.",
},
{
"caseId": "core-err-timeout",
"query": "ERR_TIMEOUT",
"topK": 3,
"purpose": "Exact error-code retrieval should remain stable.",
},
{
"caseId": "core-mysql-connection-pool",
"query": "MySQL connection pool is exhausted. How should I diagnose it?",
"topK": 3,
"purpose": "Core infrastructure troubleshooting retrieval.",
},
{
"caseId": "aiops-payment-latency",
"query": "Alert HighLatency on payment-service with p95 latency above threshold",
"topK": 3,
"purpose": "AIOps alert-style retrieval.",
},
]
@dataclass
class LiveCase:
case_id: str
query: str
top_k: int
purpose: str
category: str | None = None
@classmethod
def from_json(cls, raw: dict[str, Any]) -> "LiveCase":
return cls(
case_id=str(raw["caseId"]),
query=str(raw["query"]),
top_k=int(raw.get("topK") or 3),
purpose=str(raw.get("purpose") or raw.get("notes") or ""),
category=(
str(raw.get("category"))
if raw.get("category") not in (None, "")
else None
),
)
def load_cases(path: Path | None) -> list[LiveCase]:
if path is None:
return [LiveCase.from_json(item) for item in DEFAULT_CASES]
with path.open("r", encoding="utf-8") as handle:
payload = json.load(handle)
raw_cases = payload.get("cases", payload)
return [LiveCase.from_json(item) for item in raw_cases]
def write_json(path: Path, payload: Any) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
with path.open("w", encoding="utf-8", newline="\n") as handle:
json.dump(payload, handle, ensure_ascii=False, indent=2)
handle.write("\n")
def write_text(path: Path, content: str) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
with path.open("w", encoding="utf-8", newline="\n") as handle:
handle.write(content)
def request_case(base_url: str, case: LiveCase, timeout_seconds: float) -> dict[str, Any]:
endpoint = base_url.rstrip("/") + "/api/search/similar"
params: dict[str, str] = {
"query": case.query,
"topK": str(case.top_k),
}
if case.category:
params["category"] = case.category
url = endpoint + "?" + urllib.parse.urlencode(params)
started_at = datetime.now(timezone.utc)
try:
with urllib.request.urlopen(url, timeout=timeout_seconds) as response:
body = response.read().decode("utf-8")
payload = json.loads(body)
status = int(getattr(response, "status", 200))
except (urllib.error.URLError, TimeoutError, json.JSONDecodeError) as exc:
return {
"caseId": case.case_id,
"query": case.query,
"topK": case.top_k,
"category": case.category,
"purpose": case.purpose,
"url": url,
"ok": False,
"error": str(exc),
"resultCount": 0,
"topCandidates": [],
"rawResponse": None,
"startedAt": started_at.isoformat(),
}
data = payload.get("data") if isinstance(payload, dict) else None
if not isinstance(data, list):
data = []
ok = status == 200 and payload.get("code") == 200
return {
"caseId": case.case_id,
"query": case.query,
"topK": case.top_k,
"category": case.category,
"purpose": case.purpose,
"url": url,
"ok": ok,
"httpStatus": status,
"responseCode": payload.get("code"),
"responseMessage": payload.get("message"),
"resultCount": len(data),
"topCandidates": [summarize_candidate(item, index + 1) for index, item in enumerate(data)],
"rawResponse": payload,
"startedAt": started_at.isoformat(),
}
def summarize_candidate(raw: dict[str, Any], rank: int) -> dict[str, Any]:
metadata = parse_metadata(raw.get("metadata"))
return {
"rank": rank,
"id": raw.get("id"),
"title": metadata.get("title"),
"breadcrumb": metadata.get("breadcrumb"),
"category": metadata.get("category"),
"source": metadata.get("_source") or metadata.get("source"),
"score": raw.get("score"),
"rawScore": raw.get("rawScore"),
"scoreLabel": raw.get("scoreLabel"),
"contentPreview": preview(raw.get("content")),
}
def parse_metadata(value: Any) -> dict[str, Any]:
if isinstance(value, dict):
return value
if isinstance(value, str) and value.strip():
try:
parsed = json.loads(value)
return parsed if isinstance(parsed, dict) else {}
except json.JSONDecodeError:
return {}
return {}
def preview(value: Any, limit: int = 180) -> str:
text = " ".join(str(value or "").split())
if len(text) <= limit:
return text
return text[: limit - 3] + "..."
def render_markdown(report: dict[str, Any]) -> str:
lines = [
"# RAG Live Post-Reindex Acceptance",
"",
f"Generated at: `{report['generatedAt']}`",
f"Base URL: `{report['baseUrl']}`",
"",
"> Reindex prerequisite: this report only reflects breadcrumb-aware embedding if the knowledge base was reindexed after the embedding-text change.",
"",
"## Summary",
"",
"| Metric | Value |",
"|---|---:|",
f"| Cases | {report['caseCount']} |",
f"| Successful calls | {report['successfulCalls']} |",
f"| Empty result cases | {report['emptyResultCases']} |",
"",
"## Cases",
"",
"| Case | Purpose | Results | Top Candidates |",
"|---|---|---:|---|",
]
for item in report["results"]:
top = "<br>".join(format_candidate(candidate) for candidate in item["topCandidates"])
if not top and item.get("error"):
top = "ERROR: " + str(item["error"])
lines.append(
"| {case} | {purpose} | {count} | {top} |".format(
case=item["caseId"],
purpose=item.get("purpose") or "",
count=item["resultCount"],
top=top,
)
)
lines.append("")
return "\n".join(lines)
def format_candidate(candidate: dict[str, Any]) -> str:
label = candidate.get("title") or candidate.get("source") or candidate.get("id") or ""
breadcrumb = candidate.get("breadcrumb") or ""
score_label = candidate.get("scoreLabel") or ""
score = candidate.get("score")
raw_score = candidate.get("rawScore")
details = f"score={score}"
if raw_score is not None:
details += f", raw={raw_score}"
if score_label:
details += f", label={score_label}"
if breadcrumb:
return f"{candidate['rank']}. {label} ({breadcrumb}; {details})"
return f"{candidate['rank']}. {label} ({details})"
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--base-url", default=DEFAULT_BASE_URL)
parser.add_argument("--cases", type=Path, default=None)
parser.add_argument("--json-report", type=Path, default=DEFAULT_JSON_REPORT)
parser.add_argument("--markdown-report", type=Path, default=DEFAULT_MD_REPORT)
parser.add_argument("--timeout-seconds", type=float, default=10.0)
return parser.parse_args()
def main() -> int:
args = parse_args()
cases = load_cases(args.cases)
results = [
request_case(args.base_url, case, args.timeout_seconds)
for case in cases
]
successful = [item for item in results if item["ok"]]
empty = [item for item in results if item["ok"] and item["resultCount"] == 0]
report = {
"generatedAt": datetime.now(timezone.utc).isoformat(),
"baseUrl": args.base_url,
"caseCount": len(results),
"successfulCalls": len(successful),
"emptyResultCases": len(empty),
"reindexPrerequisite": "Run or trigger knowledge-base reindex before treating this as breadcrumb-aware embedding evidence.",
"results": results,
}
write_json(args.json_report, report)
write_text(args.markdown_report, render_markdown(report))
print(
"Ran {total} live cases: successful={successful}, empty={empty}".format(
total=len(results),
successful=len(successful),
empty=len(empty),
)
)
return 1 if len(successful) != len(results) else 0
if __name__ == "__main__":
raise SystemExit(main())
-176
View File
@@ -1,176 +0,0 @@
#!/usr/bin/env python3
"""Rebuild knowledge into the configured Milvus collection (default: biz).
Clears:
- milvus.collection (drop + recreate dense+BM25 schema)
- MySQL api_document
- in-memory L0 index
Then force-imports all markdown under server-side knowledge.base-path
(default: knowledge_base/).
Usage:
# start Spring Boot first, then:
python scripts/rebuild_hybrid_knowledge.py --confirm REBUILD
python scripts/rebuild_hybrid_knowledge.py --base-url http://127.0.0.1:9900 --confirm REBUILD
"""
from __future__ import annotations
import argparse
import json
import sys
import urllib.error
import urllib.request
from typing import Any
DEFAULT_BASE_URL = "http://localhost:9900"
def http_json(method: str, url: str, timeout: float = 3600.0) -> tuple[int, Any]:
req = urllib.request.Request(url=url, method=method.upper())
req.add_header("Accept", "application/json")
try:
with urllib.request.urlopen(req, timeout=timeout) as resp:
raw = resp.read().decode("utf-8", errors="replace")
status = getattr(resp, "status", 200)
if not raw.strip():
return status, None
return status, json.loads(raw)
except urllib.error.HTTPError as exc:
raw = exc.read().decode("utf-8", errors="replace")
body: Any
try:
body = json.loads(raw) if raw.strip() else None
except json.JSONDecodeError:
body = raw
raise RuntimeError(f"HTTP {method} {url} failed status={exc.code}: {body}") from exc
except urllib.error.URLError as exc:
raise RuntimeError(f"HTTP {method} {url} failed: {exc}") from exc
def pretty(obj: Any) -> str:
return json.dumps(obj, ensure_ascii=False, indent=2)
def step(title: str) -> None:
print()
print(f"==> {title}")
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(
description="Drop/recreate milvus.collection (default biz), clear MySQL api_document + L0, "
"and reimport knowledge_base markdown into dense+BM25."
)
parser.add_argument(
"--base-url",
default=DEFAULT_BASE_URL,
help=f"Service base URL (default: {DEFAULT_BASE_URL})",
)
parser.add_argument(
"--confirm",
required=True,
choices=["REBUILD"],
help="Must be REBUILD to execute destructive rebuild",
)
parser.add_argument(
"--skip-stats",
action="store_true",
help="Skip before/after /api/knowledge/stats",
)
parser.add_argument(
"--timeout",
type=float,
default=7200.0,
help="Rebuild request timeout seconds (default: 7200)",
)
args = parser.parse_args(argv)
base_url = args.base_url.rstrip("/")
print("Hybrid knowledge rebuild")
print(f" BaseUrl : {base_url}")
print(f" Confirm : {args.confirm}")
print(" Source : knowledge_base/ (server-side knowledge.base-path)")
print()
print("This will DESTROY data in:")
print(" - Milvus collection milvus.collection (default: biz)")
print(" - MySQL table api_document")
print(" - In-memory L0 knowledge index")
print("Then re-import all markdown under knowledge_base.")
print()
# 1) health
step("Check service health")
try:
status, body = http_json("GET", f"{base_url}/milvus/health", timeout=30)
print(f" milvus health status={status}")
print(pretty(body))
except Exception as exc: # noqa: BLE001 - ops script should continue on soft health failure
print(f" WARN: /milvus/health failed: {exc}")
print(" Continue if app is up but milvus health endpoint has issues.")
# 2) stats before
if not args.skip_stats:
step("Knowledge stats (before)")
try:
_, body = http_json("GET", f"{base_url}/api/knowledge/stats", timeout=30)
print(pretty(body))
except Exception as exc: # noqa: BLE001
print(f" WARN: stats before failed: {exc}")
# 3) rebuild
step("POST /api/knowledge/rebuild-hybrid?confirm=REBUILD")
rebuild_url = f"{base_url}/api/knowledge/rebuild-hybrid?confirm={args.confirm}"
try:
status, body = http_json("POST", rebuild_url, timeout=args.timeout)
except RuntimeError as exc:
print(str(exc))
return 1
print(f" HTTP {status}")
print(pretty(body))
if not isinstance(body, dict):
print("Unexpected rebuild response type", file=sys.stderr)
return 1
inserted = int(body.get("inserted") or 0)
failed = int(body.get("failed") or 0)
success = bool(body.get("success"))
if not success:
if inserted <= 0:
print()
print("Rebuild reported failure and inserted=0. Inspect details above.", file=sys.stderr)
return 2
print()
print(f"Rebuild finished with failed={failed} inserted={inserted}. Review details.")
else:
print()
print(f"Rebuild OK: inserted={inserted}, failed={failed}")
# 4) stats after
if not args.skip_stats:
step("Knowledge stats (after)")
try:
_, body = http_json("GET", f"{base_url}/api/knowledge/stats", timeout=30)
print(pretty(body))
except Exception as exc: # noqa: BLE001
print(f" WARN: stats after failed: {exc}")
print()
print("Done.")
print("Next:")
print(" 1) Ensure application.yml has:")
print(" milvus.collection: biz")
print(" retrieval.search.mode: hybrid")
print(" 2) Smoke test lookup_knowledge / chat with a known doc query")
return 0 if success or inserted > 0 else 2
if __name__ == "__main__":
raise SystemExit(main())