Python 在 AI 系统中的作用¶
这里不把 Python 当成“训练大模型”的语言教程,而是说明它怎样连接模型服务、解析文档、生成 Embedding、写入 Milvus、执行 RAG 检索和评测。Python 是各组件之间的胶水,模型、向量库和知识质量仍是独立层次。
文档入库:文件 → 解析 → 清洗 → 切片 → Embedding → Milvus
在线问答:问题 → Embedding → Milvus 召回 → Reranker → LLM → 引用答案
质量闭环:问题 + 候选 + 答案 + 反馈 → 评测 → 调整切片/模型/参数
项目环境¶
uv init --no-package rag-ops
cd rag-ops
uv python pin 3.12
uv add httpx pydantic pyyaml pypdf pymilvus tenacity
uv add --dev pytest ruff
uv run python main.py
推荐结构:
rag-ops/
├── pyproject.toml
├── uv.lock
├── config.example.yml
├── src/
│ ├── model_client.py
│ ├── document_parser.py
│ ├── chunker.py
│ ├── vector_store.py
│ ├── retrieval.py
│ └── evaluation.py
├── tests/
├── data/raw/
├── data/processed/
└── output/
模型 URL、模型名、超时可以进入配置;API Key、数据库密码和知识正文不能写入 Git 或普通日志。
调用 OpenAI 兼容模型 API¶
vLLM 等推理服务通常可以提供 OpenAI 兼容接口。下面直接使用 HTTP,便于理解请求、超时与返回结构:
import os
import httpx
class ModelClient:
def __init__(self) -> None:
self.base_url = os.environ["MODEL_BASE_URL"].rstrip("/")
self.api_key = os.environ["MODEL_API_KEY"]
self.model = os.environ["MODEL_NAME"]
self.client = httpx.Client(
timeout=httpx.Timeout(connect=5, read=120, write=30, pool=5),
limits=httpx.Limits(max_connections=20, max_keepalive_connections=10),
)
def chat(self, question: str, context: str) -> str:
response = self.client.post(
f"{self.base_url}/v1/chat/completions",
headers={"Authorization": f"Bearer {self.api_key}"},
json={
"model": self.model,
"temperature": 0.1,
"messages": [
{
"role": "system",
"content": "只依据提供的证据回答;证据不足时明确说明。",
},
{
"role": "user",
"content": f"问题:{question}\n\n证据:\n{context}",
},
],
},
)
response.raise_for_status()
payload = response.json()
return payload["choices"][0]["message"]["content"]
模型调用必须设置连接和读取超时。生成请求可能需要较长读取时间,但不能无限等待;服务端还应限制最大输入、最大输出和并发。
Embedding 调用¶
def embed_texts(client: httpx.Client, base_url: str, token: str, model: str, texts: list[str]) -> list[list[float]]:
response = client.post(
f"{base_url.rstrip('/')}/v1/embeddings",
headers={"Authorization": f"Bearer {token}"},
json={"model": model, "input": texts},
)
response.raise_for_status()
items = sorted(response.json()["data"], key=lambda item: item["index"])
vectors = [item["embedding"] for item in items]
if len(vectors) != len(texts):
raise ValueError("Embedding 返回数量与输入不一致")
return vectors
首次调用后记录向量维度:
vectors = embed_texts(client, base_url, token, model, ["SSH 登录失败"])
dimension = len(vectors[0])
print(dimension)
这个维度必须与 Milvus Collection 的 Vector Field 一致。文档和用户问题必须使用兼容的 Embedding 模型、相同归一化规则和距离度量;模型升级时按版本新建或重建向量数据。
文档解析¶
from pathlib import Path
from pypdf import PdfReader
def extract_pdf(path: Path) -> list[dict]:
reader = PdfReader(path)
pages = []
for page_number, page in enumerate(reader.pages, start=1):
text = (page.extract_text() or "").strip()
if text:
pages.append({
"source": str(path),
"page": page_number,
"text": text,
})
return pages
PDF 能提取出文字不代表结果可靠。扫描件需要 OCR;表格、多栏排版、页眉页脚和图片说明需要单独验证。入库前保留来源文件校验值、页码、章节路径和解析器版本,便于答案回溯。
切片不是简单按字符截断¶
每个 Chunk 至少保存:
chunk_id
document_id
title / section_path
page / source_url
text
permission_scope
updated_at
embedding_model_version
建议按标题和语义段落切分,再设置合理重叠。运维规程应尽量保留“现象 → 检查 → 原因 → 处理 → 验证”的完整链路;如果命令和适用条件被切到不同 Chunk,即使向量召回命中也无法安全执行。
写入 Milvus¶
import os
from pymilvus import MilvusClient
milvus = MilvusClient(
uri=os.environ["MILVUS_URI"],
token=os.environ["MILVUS_TOKEN"],
)
rows = []
for chunk, vector in zip(chunks, vectors, strict=True):
rows.append({
"chunk_id": chunk["chunk_id"],
"embedding": vector,
"document_id": chunk["document_id"],
"text": chunk["text"],
"source_url": chunk["source_url"],
"permission_scope": chunk["permission_scope"],
"embedding_model_version": embedding_version,
})
result = milvus.insert(
collection_name="ops_knowledge_chunks",
data=rows,
)
生产写入要分批、记录成功 ID 和失败批次,并设计幂等主键。文档更新时要让旧 Chunk 失效或删除,不能只不断追加新版。
完整召回流程¶
query_vector = embed_texts(
client,
base_url,
token,
embedding_model,
[question],
)[0]
hits = milvus.search(
collection_name="ops_knowledge_chunks",
data=[query_vector],
filter="permission_scope in ['netops'] and device_vendor == 'huawei'",
limit=30,
output_fields=["chunk_id", "text", "title", "source_url"],
)[0]
candidates = [hit["entity"] for hit in hits]
selected = rerank(question, candidates)[:5]
context = "\n\n".join(
f"[{index}] {item['text']}\n来源: {item['source_url']}"
for index, item in enumerate(selected, start=1)
)
answer = model_client.chat(question, context)
权限过滤必须由后端根据登录身份注入,不能相信用户在问题中声明的角色。Milvus 只做候选召回;Reranker 负责精排,LLM 负责基于最终证据组织答案。
重试和降级¶
| 故障 | 建议处理 |
|---|---|
| 模型接口 429/临时 5xx | 指数退避、有上限重试并记录排队时间 |
| 模型读取超时 | 取消请求;减少上下文/输出;检查推理队列和 GPU |
| Embedding 批次失败 | 保存批次 ID,只重跑失败批次,避免重复写入 |
| Milvus 不可用 | 不绕过检索让 LLM 凭空回答私有知识;返回降级说明 |
| Reranker 不可用 | 可按经过验证的策略降级为向量分数,但明确质量下降 |
认证失败、请求格式错误和向量维度错误不应自动重试,应立即停止并修正配置。
日志与评测¶
每次问答建议记录:
不要记录 API Key 和未脱敏的敏感正文。评测至少分三层:
| 层次 | 指标 |
|---|---|
| 召回 | 正确证据是否进入 Top-K,Recall@K |
| 重排 | 正确证据是否进入最终上下文 |
| 回答 | 是否基于证据、引用正确、无依据时是否拒答 |
先定位失败层次,再调整切片、Embedding、过滤、Top-K、Reranker 或 Prompt,不能只通过更换大模型解决所有问题。