fix(worker): 知识写入 Milvus 必须探测字段名并补齐必填字段
背景:走 `POST /api/v1/knowledge/upload` 灌了 22 块场内基金知识,MySQL 全部写入成功、 向量事件也全部投递,却全部同步失败(重试 3 次进死信),客服检索不到新知识。 逐层定位出三个真问题,都在写入侧: 1) **字段名硬编码**。检索侧早已按 AGENTS.md 改用运行时探测 (`app/core/knowledge_schema.py`:`doc_id`↔`knowledge_id`、`content`↔`snippet`), 写侧却一直硬编码 `knowledge_id` / `snippet`。本机集合实际是 `doc_id`/`content`/`chapter`/`visibility`/…(架构师那套 schema),于是 `Attempt to insert an unexpected field knowledge_id` 整条失败。 现在两侧共用 `resolve_schema` 的同一份映射表,调用方只用**逻辑**字段名; 集合没有的字段(如另一套环境无 `intent`)跳过而不是报错。 2) **非 nullable 必填字段没给**。改完名字后报 `Insert missed an field chapter to collection without set nullable==true`: 集合里存在、调用方没提供的 VARCHAR 字段必须补值。VARCHAR 的判据用 `params.max_length`,不必 import pymilvus 的枚举。 3) **补值不能一律空串**:`visibility` 留空会被检索侧 `visibility == "public"` 的过滤整行排除。这条最隐蔽 —— `query` 查得到、`search` 查不到,表现为「入库成功但客服永远答不出新知识」,比写入直接失败更难定位。 现在按 `FIELD_DEFAULTS` 给有语义的字段默认值。 Worker 侧把 `fields` 的键改成逻辑名(`snippet` → `content`),上限表随之改名。 **测试**:原来的桩没有 `describe_collection`,所以这条路径从未覆盖到真机 schema (这正是缺陷长期存在的原因)。现在桩提供两套真实 schema 并参数化,另加 「跳过集合没有的字段」「探测不出必要字段必须失败关闭」两个用例。 **验证**:22 块重新同步后 `fin_product_collection` 236 行、`visibility` 全为 `public`; 问「南方沪深300ETF 的起投金额是多少」命中新块 `score=0.7932`(≥ `HIGH_SCORE` 0.75), 客服从「转人工」变为直接作答,端到端 3.34 秒。 **另**:新增 `docs/43-场内基金产品手册(知识库入库版).md` —— 从 `docs/42` 抽出 **客户可见**的纯净内容(去掉全部内部决策备注)供入库;`docs/42` 保留为草稿与决策记录。
This commit is contained in:
@@ -4,30 +4,101 @@
|
||||
`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` 的退避重试与死信机制处理(本层不写重试逻辑)。
|
||||
同一批集合名(`fin_faq_collection` 等)在不同环境里是**两套不同的 schema**(实测确认):
|
||||
|
||||
装配位置:本适配器的实例由**组装层**(`app/service/agent/bootstrap.py`)创建后注入
|
||||
`build_knowledge_handlers(...)`,再注册进 `WorkerRuntime.dispatch_one` —— 这件事**还没做**
|
||||
(Task 6 的独立待办;Task 5 不碰 `runtime.py`/`bootstrap.py`),所以目前只有
|
||||
`dispatch_knowledge_events(...)` 这条直接调用路径能用它。
|
||||
| 逻辑字段 | 一套环境 | 另一套环境 |
|
||||
|---|---|---|
|
||||
| 文档标识 | `knowledge_id` | `doc_id` |
|
||||
| 正文 | `snippet` | `content` |
|
||||
| 章节 / 可见性 / 来源文件 | 无 | `chapter` / `visibility` / `source_file` |
|
||||
|
||||
Milvus 对不存在的字段直接报错(`Attempt to insert an unexpected field`),而这些集合都
|
||||
`enable_dynamic_field=False`,所以**写错一个键名整条 upsert 就失败**。检索侧早已改用运行时
|
||||
探测(`app/core/knowledge_schema.py`),写侧此前一直硬编码 `knowledge_id`/`snippet` ——
|
||||
后果是:**只要环境不是这套名字,从接口上传的知识全部同步失败,而检索侧看不出异常**
|
||||
(读得到老知识,新知识静默缺席)。2026-09-13 实测踩到:22 块新知识全部
|
||||
`RecoverableAgentError: 知识向量写入失败`,事件重试 3 次后进死信。
|
||||
|
||||
现在两侧共用 `resolve_schema` 的同一份映射表,调用方只用**逻辑字段名**,
|
||||
由本模块映射到该集合的真实物理名;集合没有的逻辑字段(如另一套环境没有 `intent`)
|
||||
**跳过而不是报错**。
|
||||
|
||||
## 缺字段也要补:非 nullable 的标量字段是必填
|
||||
|
||||
字段名对上了还不够。改了名字之后的第二次实测报的是另一件事:
|
||||
|
||||
Insert missed an field `chapter` to collection without set nullable==true or set default_value
|
||||
|
||||
即集合里存在、但调用方没有提供的**非 nullable 标量字段**,Milvus 在 insert/upsert 时要求
|
||||
必须给出 —— 缺一个就整条失败。所以本模块会把「集合有、这一行没给」的 VARCHAR 字段补成
|
||||
空串;`params.max_length` 是判断 VARCHAR 的稳定依据(无需 import pymilvus 的枚举)。
|
||||
|
||||
## 幂等口径
|
||||
|
||||
Milvus 的 `upsert` 按主键**覆盖**同一实体的向量与标量字段,因此同一知识被重复投递时结果是
|
||||
"向量数仍为 1、内容是最后一次写入"。这正是重跑导入、或正文被 UPDATE 后重投时需要的行为 ——
|
||||
本适配器**不做**任何"已存在就跳过"的判断:跳过会让 Milvus 里留着旧正文对应的旧向量
|
||||
(检索命中旧答案)。
|
||||
|
||||
连接是**惰性**的:`__init__` 不连 Milvus,首次写入时才 `import pymilvus` 并建立
|
||||
`AsyncMilvusClient`;`pymilvus` 缺失/连不上统一转成 `RecoverableAgentError`,
|
||||
交 `OutboxWorker` 的退避重试与死信机制处理(本层不写重试逻辑)。
|
||||
"""
|
||||
|
||||
from collections.abc import Mapping, Sequence
|
||||
from typing import Any
|
||||
|
||||
from app.core.errors import RecoverableAgentError
|
||||
from app.core.knowledge_schema import CollectionSchema, SchemaCache, resolve_schema
|
||||
|
||||
#: 向量字段名必须与集合 schema 一致:`tools/setup_milvus_knowledge_collections.py`
|
||||
#: 用 `embedding` 建字段,且集合 `enable_dynamic_field=False` —— 写成别的键(例如 `vector`)
|
||||
#: 会让 `upsert` 直接失败,检索侧永远命中不到(真机已核实 schema 字段名为 `embedding`)。
|
||||
#: 向量字段名必须与集合 schema 一致。两套实测 schema 都叫 `embedding`,
|
||||
#: 且集合 `enable_dynamic_field=False` —— 写成别的键(例如 `vector`)会让 `upsert` 直接失败,
|
||||
#: 检索侧永远命中不到。
|
||||
VECTOR_FIELD = "embedding"
|
||||
|
||||
#: 主键字段的**逻辑**名(物理名由探测决定:`doc_id` 或 `knowledge_id`)。
|
||||
PRIMARY_LOGICAL_FIELD = "doc_id"
|
||||
#: 正文字段的逻辑名(物理名可能是 `content` 或 `snippet`)。
|
||||
CONTENT_LOGICAL_FIELD = "content"
|
||||
|
||||
#: 「集合里有、这一行没给」的字段该补什么值。
|
||||
#:
|
||||
#: 多数 VARCHAR 补空串即可,但**有语义的字段必须给对**:`visibility` 留空会让检索侧的
|
||||
#: `visibility == "public"` 过滤把整行排除掉。实测代价(2026-09-13):22 块知识全部
|
||||
#: 写进了 Milvus(`query` 能查到),却一条都检索不到(`search` 查不到)——
|
||||
#: 现象是"入库成功但客服永远答不出新知识",比写入失败更难查。
|
||||
FIELD_DEFAULTS: dict[str, str] = {
|
||||
"visibility": "public",
|
||||
}
|
||||
|
||||
|
||||
def _varchar_fields(description: Any) -> frozenset[str]:
|
||||
"""从 `describe_collection` 结果里挑出 VARCHAR 字段名。
|
||||
|
||||
判据用 `params.max_length`:Milvus 的 VarChar 字段必带它,而向量/数值字段不带。
|
||||
这样不必 import `pymilvus` 的 `DataType` 枚举(本模块刻意对它惰性依赖)。
|
||||
"""
|
||||
if not isinstance(description, Mapping):
|
||||
return frozenset()
|
||||
fields = description.get("fields")
|
||||
if fields is None:
|
||||
schema = description.get("schema")
|
||||
fields = schema.get("fields") if isinstance(schema, Mapping) else None
|
||||
if not isinstance(fields, Sequence):
|
||||
return frozenset()
|
||||
names: list[str] = []
|
||||
for item in fields:
|
||||
if not isinstance(item, Mapping):
|
||||
continue
|
||||
name = item.get("name")
|
||||
params = item.get("params")
|
||||
if isinstance(name, str) and name and isinstance(params, Mapping):
|
||||
if "max_length" in params:
|
||||
names.append(name)
|
||||
return frozenset(names)
|
||||
|
||||
|
||||
class MilvusKnowledgeWriter:
|
||||
"""知识向量的写边界:upsert(覆盖)/ delete,失败一律 `RecoverableAgentError`。"""
|
||||
@@ -36,6 +107,10 @@ class MilvusKnowledgeWriter:
|
||||
self._uri = uri
|
||||
self._token = token
|
||||
self._client: Any = None
|
||||
self._schemas = SchemaCache()
|
||||
#: 集合名 → (字段映射, 该集合的 VARCHAR 字段)。后者用于给「集合有、这一行没给」的
|
||||
#: 标量字段补空值,理由见模块 docstring。
|
||||
self._descriptions: dict[str, tuple[CollectionSchema, frozenset[str]]] = {}
|
||||
|
||||
async def _ensure(self) -> Any:
|
||||
if self._client is None:
|
||||
@@ -49,6 +124,38 @@ class MilvusKnowledgeWriter:
|
||||
raise RecoverableAgentError("Milvus 写客户端初始化失败") from exc
|
||||
return self._client
|
||||
|
||||
async def _describe(self, collection: str) -> tuple[CollectionSchema, frozenset[str]]:
|
||||
"""探测(并缓存)该集合的字段映射与 VARCHAR 字段集合。
|
||||
|
||||
与检索侧同一个 `resolve_schema`,但这里必须走**异步**调用:写侧用的是
|
||||
`AsyncMilvusClient`,而 `knowledge_schema.detect_schema` 是同步实现、只服务于
|
||||
检索侧的 `MilvusClient`。探测失败不抛异常,返回不可用的 schema 由调用方判定。
|
||||
"""
|
||||
cached = self._descriptions.get(collection)
|
||||
if cached is not None:
|
||||
return cached
|
||||
client = await self._ensure()
|
||||
try:
|
||||
description = await client.describe_collection(collection_name=collection)
|
||||
except Exception as exc: # 探测失败=该集合不可用,由调用方给出明确错误
|
||||
schema = CollectionSchema(
|
||||
collection=collection, fields={}, physical_names=frozenset(),
|
||||
missing_required=(PRIMARY_LOGICAL_FIELD, CONTENT_LOGICAL_FIELD),
|
||||
error=f"{type(exc).__name__}: {exc}",
|
||||
)
|
||||
result: tuple[CollectionSchema, frozenset[str]] = (schema, frozenset())
|
||||
else:
|
||||
schema = resolve_schema(collection, description)
|
||||
result = (schema, _varchar_fields(description))
|
||||
self._descriptions[collection] = result
|
||||
self._schemas.put(schema)
|
||||
return result
|
||||
|
||||
async def schema_for(self, collection: str) -> CollectionSchema:
|
||||
"""探测(并缓存)该集合的字段映射。"""
|
||||
schema, _ = await self._describe(collection)
|
||||
return schema
|
||||
|
||||
async def upsert(
|
||||
self,
|
||||
*,
|
||||
@@ -57,8 +164,10 @@ class MilvusKnowledgeWriter:
|
||||
vector: list[float],
|
||||
fields: dict[str, Any],
|
||||
) -> None:
|
||||
"""按 `knowledge_id` 覆盖写入一条向量(Milvus 主键 upsert,重复投递不产生重复向量)。
|
||||
"""按主键覆盖写入一条向量(Milvus 主键 upsert,重复投递不产生重复向量)。
|
||||
|
||||
`fields` 的键是**逻辑字段名**(`content`/`title`/`tags`/`version`/`intent`),
|
||||
由探测结果映射成物理名;集合没有的逻辑字段直接跳过。
|
||||
`fields` 里**可以没有** `intent` 键:知识契约把 `intent` 定为稀疏标签,
|
||||
无显式标签时由检索侧按集合名推断(见 `knowledge_vector_worker` 的约定说明)。
|
||||
"""
|
||||
@@ -66,12 +175,28 @@ class MilvusKnowledgeWriter:
|
||||
raise RecoverableAgentError("knowledge_id 不能为空")
|
||||
if not vector:
|
||||
raise RecoverableAgentError("向量不能为空")
|
||||
schema, varchar_fields = await self._describe(collection)
|
||||
primary = schema.resolve(PRIMARY_LOGICAL_FIELD)
|
||||
if primary is None or not schema.usable:
|
||||
missing = list(schema.missing_required) or [schema.error or "未知原因"]
|
||||
raise RecoverableAgentError(
|
||||
f"知识集合 {collection} 缺少必要字段({missing}),无法写入向量"
|
||||
)
|
||||
row: dict[str, Any] = {primary: knowledge_id, VECTOR_FIELD: vector}
|
||||
for logical, value in fields.items():
|
||||
physical = schema.resolve(logical)
|
||||
if physical is None:
|
||||
continue # 该集合没有这个字段(如 `intent`),跳过而不是让整条写入失败
|
||||
row[physical] = value
|
||||
# 集合里存在、但这一行没给的 VARCHAR 字段必须补值,否则 Milvus 报
|
||||
# `Insert missed an field ...`(非 nullable 且无默认值=必填)。主键已经填过,
|
||||
# 这里跳过它免得覆盖掉真正的 id;有语义的字段按 FIELD_DEFAULTS 给对值。
|
||||
for physical in varchar_fields:
|
||||
if physical != primary and physical not in row:
|
||||
row[physical] = FIELD_DEFAULTS.get(physical, "")
|
||||
client = await self._ensure()
|
||||
try:
|
||||
await client.upsert(
|
||||
collection_name=collection,
|
||||
data=[{"knowledge_id": knowledge_id, VECTOR_FIELD: vector, **fields}],
|
||||
)
|
||||
await client.upsert(collection_name=collection, data=[row])
|
||||
except RecoverableAgentError:
|
||||
raise
|
||||
except Exception as exc:
|
||||
|
||||
@@ -105,9 +105,11 @@ Handler = Callable[[dict[str, Any]], Awaitable[None]]
|
||||
#: 会超限并使整条 upsert 失败(实测:政策文档标题 346 字节 > 256,报
|
||||
#: `length of varchar field title exceeds max length`)。FAQ 语料标题短,掩盖了这个缺陷;
|
||||
#: 长标题的政策/产品文档一灌就炸。所以写入前必须逐个字段按字节截断。
|
||||
#: 键是**逻辑字段名**(见 `MilvusKnowledgeWriter`):物理名由集合 schema 运行时探测决定
|
||||
#: (两套环境分别叫 `content` 与 `snippet`),这里统一用 `content`。
|
||||
VECTOR_FIELD_LIMITS: dict[str, int] = {
|
||||
"title": 256,
|
||||
"snippet": 4000,
|
||||
"content": 4000,
|
||||
"tags": 512,
|
||||
"version": 16,
|
||||
"intent": 32,
|
||||
@@ -117,7 +119,7 @@ VECTOR_FIELD_LIMITS: dict[str, int] = {
|
||||
def _fit_varchar(value: object, max_bytes: int) -> str:
|
||||
"""把值转成字符串并按 **UTF-8 字节**截断到上限内(不切坏多字节字符)。
|
||||
|
||||
截断而非报错:`title`/`snippet` 只是检索辅助与展示字段,让整条向量同步因为字段过长
|
||||
截断而非报错:`title`/`content` 只是检索辅助与展示字段,让整条向量同步因为字段过长
|
||||
失败会阻塞知识入库;而正文的完整内容仍保存在 MySQL(权威源)。
|
||||
"""
|
||||
text = "" if value is None else str(value)
|
||||
@@ -204,9 +206,10 @@ def build_knowledge_handlers(
|
||||
raise RecoverableAgentError("嵌入维度与集合定义不一致")
|
||||
fields: dict[str, Any] = {
|
||||
"title": _fit_varchar(row.title, VECTOR_FIELD_LIMITS["title"]),
|
||||
# snippet 是检索返回给模型的正文;必须是**本次**读到的 content_text,
|
||||
# 否则正文更新后重投会留下旧答案(见模块 docstring 的幂等口径)。
|
||||
"snippet": _fit_varchar(row.content_text, VECTOR_FIELD_LIMITS["snippet"]),
|
||||
# `content` 是**逻辑**字段名,由 writer 按集合 schema 映射成 `content` 或 `snippet`。
|
||||
# 它必须是**本次**读到的 content_text,否则正文更新后重投会留下旧答案
|
||||
# (见模块 docstring 的幂等口径)。
|
||||
"content": _fit_varchar(row.content_text, VECTOR_FIELD_LIMITS["content"]),
|
||||
"tags": _fit_varchar(_tags_to_string(row.tags), VECTOR_FIELD_LIMITS["tags"]),
|
||||
"version": _fit_varchar(
|
||||
row.version or "", VECTOR_FIELD_LIMITS["version"]
|
||||
|
||||
Reference in New Issue
Block a user