Files
group_fqcd_jr/app/infrastructure/milvus_knowledge_writer.py

96 lines
4.6 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.
"""知识写路径的 Milvus 适配器。
**只被写路径(`app/worker/knowledge_vector_worker.py`)导入;检索侧(Task 6 的
`MilvusKnowledgeClient` / `KnowledgeRetrievalService`)不得导入本模块** —— 读写物理隔离:
检索进程永远不持有写客户端,向量库故障不能从写路径传染到问答主链路,反之亦然。
幂等口径(Task 5 裁定 3):Milvus 的 `upsert` 按主键 `knowledge_id` **覆盖**同一实体的向量与
标量字段,因此同一 `knowledge_id` 被重复投递时结果是"向量数仍为 1、内容是最后一次写入"。
这正是重跑导入、或正文被 UPDATE 后重投时需要的行为 —— 本适配器**不做**任何"已存在就跳过"
的判断:跳过会让 Milvus 里留着旧正文对应的旧向量(检索命中旧答案)。
连接是**惰性**的:`__init__` 不连 Milvus(Task 5 不碰真机,Task 6 才建集合),
首次写入时才 `import pymilvus` 并建立 `AsyncMilvusClient`;`pymilvus` 缺失/连不上统一
转成 `RecoverableAgentError`,交 `OutboxWorker` 的退避重试与死信机制处理(本层不写重试逻辑)。
装配位置:本适配器的实例由**组装层**(`app/service/agent/bootstrap.py`)创建后注入
`build_knowledge_handlers(...)`,再注册进 `WorkerRuntime.dispatch_one` —— 这件事**还没做**
(Task 6 的独立待办;Task 5 不碰 `runtime.py`/`bootstrap.py`),所以目前只有
`dispatch_knowledge_events(...)` 这条直接调用路径能用它。
"""
from typing import Any
from app.core.errors import RecoverableAgentError
#: 向量字段名必须与集合 schema 一致:`tools/setup_milvus_knowledge_collections.py`
#: 用 `embedding` 建字段,且集合 `enable_dynamic_field=False` —— 写成别的键(例如 `vector`)
#: 会让 `upsert` 直接失败,检索侧永远命中不到(真机已核实 schema 字段名为 `embedding`)。
VECTOR_FIELD = "embedding"
class MilvusKnowledgeWriter:
"""知识向量的写边界:upsert(覆盖)/ delete,失败一律 `RecoverableAgentError`。"""
def __init__(self, uri: str, token: str = "") -> None:
self._uri = uri
self._token = token
self._client: Any = None
async def _ensure(self) -> Any:
if self._client is None:
try:
from pymilvus import AsyncMilvusClient # type: ignore[import-untyped]
except ImportError as exc: # pragma: no cover - 依赖已声明,缺装是环境问题
raise RecoverableAgentError("pymilvus 未安装,无法写入知识向量") from exc
try:
self._client = AsyncMilvusClient(uri=self._uri, token=self._token or None)
except Exception as exc:
raise RecoverableAgentError("Milvus 写客户端初始化失败") from exc
return self._client
async def upsert(
self,
*,
collection: str,
knowledge_id: str,
vector: list[float],
fields: dict[str, Any],
) -> None:
"""按 `knowledge_id` 覆盖写入一条向量(Milvus 主键 upsert,重复投递不产生重复向量)。
`fields` 里**可以没有** `intent` 键:知识契约把 `intent` 定为稀疏标签,
无显式标签时由检索侧按集合名推断(见 `knowledge_vector_worker` 的约定说明)。
"""
if not knowledge_id:
raise RecoverableAgentError("knowledge_id 不能为空")
if not vector:
raise RecoverableAgentError("向量不能为空")
client = await self._ensure()
try:
await client.upsert(
collection_name=collection,
data=[{"knowledge_id": knowledge_id, VECTOR_FIELD: vector, **fields}],
)
except RecoverableAgentError:
raise
except Exception as exc:
raise RecoverableAgentError("知识向量写入失败") from exc
async def delete(self, *, collection: str, knowledge_id: str) -> None:
"""按主键删除该知识的向量(幂等:主键不存在时 Milvus 视为无操作)。"""
if not knowledge_id:
raise RecoverableAgentError("knowledge_id 不能为空")
client = await self._ensure()
try:
await client.delete(collection_name=collection, ids=[knowledge_id])
except RecoverableAgentError:
raise
except Exception as exc:
raise RecoverableAgentError("知识向量删除失败") from exc
async def close(self) -> None:
if self._client is not None:
client, self._client = self._client, None
await client.close()