导入向量数据库节点
本文档详细介绍知识库导入流程中的导入向量数据库节点(ImportMilvusNode),该节点负责将向量化后的切片数据批量导入 Milvus 向量数据库,包括集合创建、Schema 定义、索引构建和数据插入等完整流程。
学习理念:向量数据入库是导入流程的"终点"——前面的切片、向量化、商品名识别等所有工作,最终都要写入 Milvus 才能被检索到。ImportMilvusNode 的核心是幂等性设计(集合不存在才创建)+ Schema + 索引分离构建 + 主键自动回填。
海外对标:ImportMilvusNode 的"Schema 定义 → 索引构建 → 批量插入 → 主键回填"流程对标 Pinecone 的
create_index+upsertAPI,以及 Weaviate 的schema.create_class+batch.add模式。AUTOINDEX + SPARSE_INVERTED_INDEX 的双索引配置是 Milvus 官方推荐的混合检索最佳实践。
本节 AI 替代率:~85% | 人工干预率:~15%
| 角色 | 能力范围 |
|---|---|
| 🤖 AI 擅长 | Schema 定义、索引参数配置、批量插入代码、主键回填逻辑、集合创建代码 |
| 👤 人类需理解 | Schema 字段设计(哪些字段需要存储、字段类型选择)、双索引策略(AUTOINDEX vs SPARSE_INVERTED_INDEX 的适用场景) |
阅读指引
| 颜色 | 章节 | AI 替代率 | 人工干预 | 说明 |
|---|---|---|---|---|
| 🟡 | §1 任务目标 | ~95% | ~5% | 学习目标明确 |
| 🟢 | §2 核心概念扫盲 | ~95% | ~5% | 全部是图片说明 |
| 🟡 | §3 整体流程 | ~90% | ~10% | 理解 Schema 结构即可 |
| 🟠 | §4 分步实现 | ~80% | ~20% | 9 步流程中 Schema 定义 + 索引构建是核心 |
| 🔴 | §4.4 主代码 | ~75% | ~25% | ImportMilvusNode + _insert_and_backfill_ids |
| 🟢 | §5 测试运行 | ~95% | ~5% | 看预期输出 + Attu 验证即可 |
| 🟡 | §6 总结 | ~90% | ~10% | 设计要点回顾 |
技术栈健康度标签体系
| 技术 | 健康度 | 建议 |
|---|---|---|
| Milvus | 🟢 稳定 | 核心向量数据库,Schema + AUTOINDEX + SPARSE_INVERTED_INDEX 是成熟方案。PyMilvus 2.5+ API 较稳定,但 Collection Schema 定义在不同版本间有微调。 |
| AUTOINDEX | 🔥 巅峰 | Milvus 自动索引选择算法,无需手动指定 IVF/HNSW 参数。适合学习和原型开发。生产环境可考虑手动调优 HNSW。 |
| SPARSE_INVERTED_INDEX | 🔥 巅峰 | 稀疏向量倒排索引,搭配 DAAT_MAXSCORE 剪枝算法。Milvus 2.4+ 新增功能,是混合检索的核心支撑。 |
| PyMilvus DataType | 🟢 稳定 | Milvus Python SDK 的数据类型定义。INT64 / VARCHAR / FLOAT_VECTOR / SPARSE_FLOAT_VECTOR 等类型映射固定。 |
| auto_id 主键 | 🟢 稳定 | Milvus 自动主键生成机制。INT64 类型,全局唯一。插入后需回填到业务数据中。 |
体系说明:🟢🟡🟠🔴 标识学习优先级 / AI 替代率;🔥🟢⏳⚠️💀 标识技术栈健康度。
中英文对照表
| English | 中文 | 本质 |
|---|---|---|
| Collection | 集合 | Milvus 中存储向量数据的基本单位,类似 SQL 的表 |
| Schema | 模式 | 定义集合中字段的名称、类型和约束的结构 |
| Index | 索引 | 加速向量检索的数据结构,不同索引类型适用不同场景 |
| AUTOINDEX | 自动索引 | Milvus 自动选择最优索引算法的配置 |
| SPARSE_INVERTED_INDEX | 稀疏倒排索引 | 稀疏向量专用的倒排索引,支持 DAAT_MAXSCORE 剪枝 |
| Primary Key (PK) | 主键 | 唯一标识每条记录的字段,Milvus 支持自动生成 |
| auto_id | 自动 ID | Schema 配置,让 Milvus 自动生成主键而非手动传入 |
| IP (Inner Product) | 内积 | 向量相似度度量,配合 L2 归一化等价于余弦相似度 |
| DAAT_MAXSCORE | 动态剪枝 | 稀疏向量检索的高效算法,在查询时动态剪枝跳过低分文档 |
💡 程序员比喻
- ImportMilvusNode 就像
pg_dump+pg_restore——把内存中的结构化数据(chunks)持久化到数据库(Milvus),并返回生成的 ID(chunk_id)。- Schema 定义 就像 SQL 的
CREATE TABLE——定义每列的名称、类型和约束,enable_dynamic_fields=jsonb灵活列。- AUTOINDEX 就像 PostgreSQL 的
auto_explain——自动选择最优查询计划,你不用管用什么索引算法。- SPARSE_INVERTED_INDEX + DAAT_MAXSCORE 就像 Elasticsearch 的
search+terminate_after——找到了足够多的好结果就提前停止。- 主键回填 就像
INSERT ... RETURNING id——插入一行数据后立刻拿到自动生成的 ID。
1. 任务目标
1.1 本章目标
通过本章学习,你将掌握:
- Milvus 集合管理:学会创建集合、定义 Schema、配置索引
- Schema 设计:理解标量字段与向量字段的定义方式
- 索引类型选择:掌握 AUTOINDEX 和 SPARSE_INVERTED_INDEX 的适用场景
- 数据插入与主键回填:学会批量插入数据并获取自动生成的主键
- 幂等性设计:理解"集合不存在才创建"的设计模式
1.2 涉及文件
knowledge/
├── processor/import_process/nodes/
│ └── import_milvus.py # 导入向量数据库节点(本章重点)
│
└── tools/
└── milvus_utils.py # Milvus 连接管理1.3 节点在流程中的位置

2. 核心概念扫盲
2.1 Milvus 集合(Collection)
Milvus 中的**集合(Collection)**类似于关系数据库中的表,是存储和检索向量数据的基本单位:

2.2 Schema 设计
Schema 定义了集合中的字段结构,包括字段名称、数据类型、约束等:

2.3 向量索引类型
索引是加速向量检索的关键,不同索引类型适用于不同场景:

2.4 度量类型(Metric Type)
度量类型决定了如何计算向量间的相似度:

2.5 主键自动生成与回填
Milvus 支持自动生成主键,插入后需要将主键回填到业务数据中:

3. 导入向量数据库业务处理流程(总)
3.1 整体流程概述

3.2 集合 Schema 结构

4. 导入向量数据库业务处理流程(分)
4.1 目标
将向量化后的文档切片数据持久化存储到 Milvus 向量数据库,支持后续的混合检索(语义检索 + 关键词检索),并将 Milvus 自动生成的主键回填到业务数据中,为后续的知识图谱构建提供关联依据。
4.2 需求分析
| 需求项 | 说明 | 解决方案 |
|---|---|---|
| 集合自动创建 | 首次导入时集合可能不存在 | 检测集合是否存在,不存在则创建 |
| Schema 灵活性 | 需要同时存储标量和向量数据 | 定义完整的 Schema,包含 8 个字段 |
| 混合检索支持 | 需要同时支持稠密和稀疏向量检索 | 为两种向量分别建立索引 |
| 主键唯一性 | 每条切片需要唯一标识 | 使用 auto_id 自动生成主键 |
| 主键可追溯 | 后续节点需要使用 chunk_id | 插入后回填主键到业务数据 |
| 幂等性 | 重复运行不应创建重复集合 | has_collection 检查后再创建 |
4.3 实现流程
4.3.1 实现流程图

4.3.2 具体实现步骤
Step1:获取并校验切片数据
从状态中获取切片列表,验证数据有效性并提取向量维度。
输入:
state: 图状态字典,包含向量化后的切片数据
处理逻辑:
- 从 state 中获取 "chunks" 字段
- 如果 chunks 为空,记录警告并直接返回(幂等设计)
- 从第一条切片中提取 dense_vector 的维度
- 如果向量维度为 0,抛出
MilvusError异常
输出:
- 验证通过的切片列表
chunks - 向量维度
vector_dim
代码片段:
🟡 【P1 看注释就行】 Step 1 校验代码——chunks 为空则跳过(幂等设计),
_get_vector_dim从首条 chunk 提取稠密向量维度。
def process(self, state: ImportGraphState) -> ImportGraphState:
config = get_config()
chunks = state.get("chunks", [])
# 1. 校验
if not chunks:
self.logger.warning("chunks 为空,跳过导入")
return state
vector_dim = self._get_vector_dim(chunks)
self.log_step("step_1", f"准备导入 {len(chunks)} 条数据,向量维度: {vector_dim}")
@staticmethod
def _get_vector_dim(chunks: List[Dict]) -> int:
"""从首条 chunk 中获取稠密向量维度。"""
dim = len(chunks[0].get("dense_vector", []))
if dim == 0:
raise MilvusError("切片数据不包含 dense_vector", node_name="import_milvus")
return dimStep2:连接 Milvus
获取 Milvus 客户端连接实例。
处理逻辑:
- 调用
get_milvus_client()获取客户端单例 - 客户端根据环境变量
MILVUS_URL配置连接地址 - 首次调用时建立连接,后续调用复用
环境变量配置:
MILVUS_URL: Milvus 服务地址,例如 "http://localhost:19530"
代码片段:
🟡 【P1 看注释就行】 Step 2 连接代码——
get_milvus_client()获取单例客户端。
try:
client = get_milvus_client()
collection_name = config.chunks_collection # 从配置获取集合名Step3:检查并创建集合
检查目标集合是否存在,不存在则创建。
处理逻辑:
- 调用
client.has_collection()检查集合是否存在 - 如果不存在,调用
_create_collection()创建集合 - 创建过程包括:构建 Schema → 构建索引参数 → 创建集合
幂等性保证:
- 只有在集合不存在时才创建
- 重复执行不会报错或创建重复集合
代码片段:
🟡 【P1 看注释就行】 Step 3 幂等创建——
has_collection检查 → 不存在才创建。避免重复运行报错。
# 3. 创建集合(如不存在)
if not client.has_collection(collection_name=collection_name):
self.log_step("step_2", f"创建集合: {collection_name}")
self._create_collection(client, collection_name, vector_dim)Step4:构建集合 Schema
定义集合的字段结构。
字段定义:
- chunk_id: INT64 主键,自动生成
- content: VARCHAR,存储切片文本内容
- title: VARCHAR,存储切片标题
- parent_title: VARCHAR,存储父级标题
- part: INT8,存储切片在同标题下的序号
- file_title: VARCHAR,存储源文件标题
- item_name: VARCHAR,存储商品名称
- sparse_vector: SPARSE_FLOAT_VECTOR,存储稀疏向量
- dense_vector: FLOAT_VECTOR,存储稠密向量
代码片段:
🔥 【P0 必须要学】 Step 4 Schema 定义是 Milvus 的核心设计。注意字段类型映射:
INT64= 主键(auto_id),VARCHAR= 文本字段(max_length=65535),FLOAT_VECTOR= 稠密向量(dim=vector_dim),SPARSE_FLOAT_VECTOR= 稀疏向量。enable_dynamic_fields=True允许动态添加字段。
def _build_schema(self, client, vector_dim: int):
"""构建集合 Schema。"""
# 1. 定义 schema
schema = client.create_schema(enable_dynamic_fields=True)
# 2. 添加主键(自增)
schema.add_field(
field_name="chunk_id",
datatype=DataType.INT64,
is_primary=True,
auto_id=True,
)
# 3. 添加标量字段
schema.add_field(field_name="content", datatype=DataType.VARCHAR, max_length=65535)
schema.add_field(field_name="title", datatype=DataType.VARCHAR, max_length=65535)
schema.add_field(field_name="parent_title", datatype=DataType.VARCHAR, max_length=65535)
schema.add_field(field_name="part", datatype=DataType.INT8)
schema.add_field(field_name="file_title", datatype=DataType.VARCHAR, max_length=65535)
schema.add_field(field_name="item_name", datatype=DataType.VARCHAR, max_length=65535)
# 4. 添加向量字段
schema.add_field(field_name="sparse_vector", datatype=DataType.SPARSE_FLOAT_VECTOR)
schema.add_field(field_name="dense_vector", datatype=DataType.FLOAT_VECTOR, dim=vector_dim)
# 5. 返回 schema
return schemaStep5:构建索引参数
为稠密向量和稀疏向量分别配置索引。
索引配置:
稠密向量索引:
- 类型:AUTOINDEX(自动选择最优算法)
- 度量:IP(内积)
稀疏向量索引:
- 类型:SPARSE_INVERTED_INDEX(倒排索引)
- 度量:IP(内积)
- 算法:DAAT_MAXSCORE(动态剪枝,高效查询)
代码片段:
🔥 【P0 必须要学】 Step 5 索引配置——稠密向量用
AUTOINDEX(自动选择 HNSW/IVF 等),稀疏向量用SPARSE_INVERTED_INDEX+DAAT_MAXSCORE剪枝。两者都用 IP 度量(配合 L2 归一化后 IP = 余弦相似度)。
def _build_index_params(self, client):
"""构建索引参数。"""
# 1. 获取索引对象
index_params = client.prepare_index_params()
# 2. 稠密向量添加索引
index_params.add_index(
field_name="dense_vector",
index_name="dense_vector_index",
index_type="AUTOINDEX",
metric_type="IP",
)
# 3. 稀疏向量添加索引
index_params.add_index(
field_name="sparse_vector",
index_name="sparse_inverted_index",
index_type="SPARSE_INVERTED_INDEX",
metric_type="IP",
params={"inverted_index_algo": "DAAT_MAXSCORE"},
)
return index_paramsStep6:创建集合
组合 Schema 和索引参数,创建集合。
处理逻辑:
- 调用
_build_schema()构建 Schema - 调用
_build_index_params()构建索引参数 - 调用
client.create_collection()创建集合
代码片段:
🟡 【P1 看注释就行】 Step 6 创建集合——组合 schema + index_params 一次性创建。注意
_build_schema+_build_index_params的职责分离。
def _create_collection(self, client, collection_name: str, vector_dim: int):
"""创建 Milvus 集合(schema + 索引)。"""
# 1. 构建 schema
schema = self._build_schema(client, vector_dim)
# 2. 构建索引
index_params = self._build_index_params(client)
# 3. 创建集合
client.create_collection(
collection_name=collection_name,
schema=schema,
index_params=index_params,
)
self.logger.info(f"集合 {collection_name} 创建成功")Step7:批量插入数据
将切片数据批量插入到 Milvus 集合中。
输入:
chunks: 包含向量的切片数据列表
处理逻辑:
- 调用
client.insert()批量插入数据 - Milvus 会自动为���条数据生成 chunk_id
- 插入结果包含
insert_count和ids
注意事项:
- 插入时不需要提供 chunk_id 字段(auto_id=True)
- Milvus 会自动处理数据类型转换
代码片段:
🟡 【P1 看注释就行】 Step 7 插入数据——
client.insert(collection_name, data=chunks)批量插入。auto_id=True时不需要传 chunk_id。
def _insert_and_backfill_ids(self, client, collection_name: str, chunks: List[Dict]):
"""插入数据并将 Milvus 返回的主键回填到每条 chunk。"""
result = client.insert(collection_name=collection_name, data=chunks)
insert_count = result.get("insert_count", 0)
self.logger.info(f"成功插入 {insert_count} 条数据")Step8:回填主键到业务数据
将 Milvus 返回的主键回填到原始切片数据中。
处理逻辑:
- 从插入结果中获取
ids列表 - 验证返回的 ID 数量与切片数量一致
- 使用
zip()将 ID 逐一回填到对应的 chunk - 将 ID 转换为字符串存储(便于后续 JSON 序列化和跨系统传递)
代码片段:
🔥 【P0 必须要学】 Step 8 主键回填是关键设计。
result.get("ids", [])获取 Milvus 自动生成的 ID,zip(chunks, inserted_ids)逐一回填,转str(chunk_id)便于 JSON 序列化。回填失败则无法关联知识图谱。
# 回填 chunk_id
inserted_ids = result.get("ids", [])
if inserted_ids and len(inserted_ids) == len(chunks):
for chunk, chunk_id in zip(chunks, inserted_ids):
chunk["chunk_id"] = str(chunk_id)
else:
self.logger.warning(
f"回填 chunk_id 失败: 返回 {len(inserted_ids)} 个 ID,"
f"期望 {len(chunks)} 个"
)Step9:更新状态并返回
将处理结果写回状态并返回。
处理逻辑:
- 将已回填 chunk_id 的 chunks 写回 state
- 返回更新后的 state
代码片段:
🟢 【P2 后面可以查】 Step 9 更新 state——
state["chunks"] = chunks。
# 5. 更新 chunks 状态
state["chunks"] = chunks
return state4.4 代码实现
🔥 【P0 必须要学】 ImportMilvusNode 是导入流程的"终点"——所有处理完的切片最终落库到 Milvus。重点理解:主流程(校验→集合创建→插入→回填)+ 幂等设计(
has_collection)+ Schema/索引分离构建。_insert_and_backfill_ids的zip回填是数据关联的关键。
"""
Milvus 导入节点
将向量化后的切片数据批量导入 Milvus,自动创建集合和索引(如果不存在)。
代码结构:
1. 主流程 (process)
2. 集合创建(schema 构建 + 索引构建 + 创建集合)
3. 数据插入与 chunk_id 回填
4. 测试入口
"""
import json
import os
from typing import List, Dict
from pymilvus import DataType
from knowledge.processor.import_process.base import BaseNode, setup_logging
from knowledge.processor.import_process.state import ImportGraphState
from knowledge.processor.import_process.config import get_config
from knowledge.processor.import_process.exceptions import MilvusError
from knowledge.tools.milvus_utils import get_milvus_client
class ImportMilvusNode(BaseNode):
"""Milvus 导入节点。
将包含稠密/稀疏向量的切片数据批量导入 Milvus,
集合不存在时自动创建 schema 和索引。
"""
name = "import_milvus"
# ================================================================== #
# 1. 主流程 #
# ================================================================== #
def process(self, state: ImportGraphState) -> ImportGraphState:
"""执行 Milvus 导入。
流程:校验数据 → 连接 Milvus → 创建集合(如需) → 插入数据 → 回填 chunk_id。
Args:
state: 图状态,需包含带向量的 chunks 列表。
Returns:
更新后的状态(chunks 中包含 chunk_id)。
Raises:
MilvusError: 数据无效或 Milvus 操作失败时抛出。
"""
config = get_config()
chunks = state.get("chunks", [])
# 1. 校验
if not chunks:
self.logger.warning("chunks 为空,跳过导入")
return state
vector_dim = self._get_vector_dim(chunks)
self.log_step("step_1", f"准备导入 {len(chunks)} 条数据,向量维度: {vector_dim}")
# 2. 连接 & 操作
try:
client = get_milvus_client()
collection_name = config.chunks_collection
# 3. 创建集合(如不存在)
if not client.has_collection(collection_name=collection_name):
self.log_step("step_2", f"创建集合: {collection_name}")
self._create_collection(client, collection_name, vector_dim)
# 4. 插入数据
self.log_step("step_3", "执行插入")
self._insert_and_backfill_ids(client, collection_name, chunks)
# 5. 更新 chunks 状态
state["chunks"] = chunks
except MilvusError:
raise
except Exception as e:
raise MilvusError(f"Milvus 操作失败: {e}", node_name=self.name, cause=e)
return state
# ================================================================== #
# 2. 集合创建 #
# ================================================================== #
def _create_collection(self, client, collection_name: str, vector_dim: int):
"""创建 Milvus 集合(schema + 索引)。
Args:
client: MilvusClient 实例。
collection_name: 集合名称。
vector_dim: 稠密向量维度。
"""
# 1. 构建 schema
schema = self._build_schema(client, vector_dim)
# 2. 构建索引
index_params = self._build_index_params(client)
# 3. 创建集合
client.create_collection(
collection_name=collection_name,
schema=schema,
index_params=index_params,
)
self.logger.info(f"集合 {collection_name} 创建成功")
def _build_schema(self, client, vector_dim: int):
"""构建集合 Schema。
Args:
client: MilvusClient 实例。
vector_dim: 稠密向量维度。
Returns:
CollectionSchema 对象。
"""
# 1. 定义 schema
schema = client.create_schema(enable_dynamic_fields=True)
# 2. 添加主键(自增)
schema.add_field(
field_name="chunk_id",
datatype=DataType.INT64,
is_primary=True,
auto_id=True,
)
# 3. 添加标量字段
schema.add_field(field_name="content", datatype=DataType.VARCHAR, max_length=65535)
schema.add_field(field_name="title", datatype=DataType.VARCHAR, max_length=65535)
schema.add_field(field_name="parent_title", datatype=DataType.VARCHAR, max_length=65535)
schema.add_field(field_name="part", datatype=DataType.INT8)
schema.add_field(field_name="file_title", datatype=DataType.VARCHAR, max_length=65535)
schema.add_field(field_name="item_name", datatype=DataType.VARCHAR, max_length=65535)
# 4. 添加向量字段
schema.add_field(field_name="sparse_vector", datatype=DataType.SPARSE_FLOAT_VECTOR)
schema.add_field(field_name="dense_vector", datatype=DataType.FLOAT_VECTOR, dim=vector_dim)
# 5. 返回 schema
return schema
def _build_index_params(self, client):
"""构建索引参数。
Args:
client: MilvusClient 实例。
Returns:
IndexParams 对象。
"""
# 1. 获取索引对象
index_params = client.prepare_index_params()
# 2. 稠密向量添加索引
index_params.add_index(
field_name="dense_vector",
index_name="dense_vector_index",
index_type="AUTOINDEX",
metric_type="IP",
)
# 3. 稀疏向量添加索引
index_params.add_index(
field_name="sparse_vector",
index_name="sparse_inverted_index",
index_type="SPARSE_INVERTED_INDEX",
metric_type="IP",
params={"inverted_index_algo": "DAAT_MAXSCORE"},
)
return index_params
# ================================================================== #
# 3. 数据插入与 chunk_id 回填 #
# ================================================================== #
@staticmethod
def _get_vector_dim(chunks: List[Dict]) -> int:
"""从首条 chunk 中获取稠密向量维度。
Args:
chunks: 切片数据列表。
Returns:
向量维度。
Raises:
MilvusError: 首条 chunk 不包含 dense_vector 时抛出。
"""
dim = len(chunks[0].get("dense_vector", []))
if dim == 0:
raise MilvusError("切片数据不包含 dense_vector", node_name="import_milvus")
return dim
def _insert_and_backfill_ids(
self,
client,
collection_name: str,
chunks: List[Dict],
):
"""插入数据并将 Milvus 返回的主键回填到每条 chunk。
Args:
client: MilvusClient 实例。
collection_name: 目标集合名称。
chunks: 切片数据列表(会被原地修改,添加 chunk_id 字段)。
"""
result = client.insert(collection_name=collection_name, data=chunks)
insert_count = result.get("insert_count", 0)
self.logger.info(f"成功插入 {insert_count} 条数据")
# 回填 chunk_id
inserted_ids = result.get("ids", [])
if inserted_ids and len(inserted_ids) == len(chunks):
for chunk, chunk_id in zip(chunks, inserted_ids):
chunk["chunk_id"] = str(chunk_id)
else:
self.logger.warning(
f"回填 chunk_id 失败: 返回 {len(inserted_ids)} 个 ID,"
f"期望 {len(chunks)} 个"
)5. 测试运行
5.1 测试代码
在 import_milvus.py 文件末尾添加独立测试代码:
🟢 【P2 后面可以查】 测试代码量大但模式固定——读取 JSON → process → 验证 chunk_id 回填。
# ================================================================== #
# 兼容 & 测试 #
# ================================================================== #
node_import_milvus = ImportMilvusNode()
if __name__ == "__main__":
"""
独立测试 ImportMilvusNode
测试流程:
1. 从上一个节点的输出文件读取状态(带向量的切片数据)
2. 执行 Milvus 导入
3. 验证回填的 chunk_id
4. 将结果保存到临时文件
"""
setup_logging()
# ----------------------------------------------------------------
# Step 1: 配置路径
# ----------------------------------------------------------------
temp_dir = r"D:\develop\develop\workspace\pycharm\usage\shopkeeper_brain_v260213\knowledge\processor\import_process\temp"
# 输入:上一个向量化节点处理后的状态
input_path = os.path.join(temp_dir, "chunks_item_name_vector.json")
# 输出:导入 Milvus 后的状态(含 chunk_id)
output_path = os.path.join(temp_dir, "chunks_item_name_vector_ids.json")
# ----------------------------------------------------------------
# Step 2: 读取输入数据
# ----------------------------------------------------------------
print(f"正在读取输入文件: {input_path}")
with open(input_path, "r", encoding="utf-8") as f:
content = json.load(f)
chunks = content.get('chunks', [])
print(f"读取到 {len(chunks)} 个切片")
# ----------------------------------------------------------------
# Step 3: 构建状态并执行处理
# ----------------------------------------------------------------
state = {
"chunks": chunks
}
print("\n开始执行 Milvus 导入...")
result_state = node_import_milvus.process(state)
# ----------------------------------------------------------------
# Step 4: 验证回填结果
# ----------------------------------------------------------------
output_chunks = result_state.get("chunks", [])
print(f"\n处理完成,共 {len(output_chunks)} 个切片")
# 检查 chunk_id 回填情况
chunks_with_id = sum(1 for c in output_chunks if c.get("chunk_id"))
chunks_without_id = len(output_chunks) - chunks_with_id
print(f"\n回填统计:")
print(f" - 成功回填 chunk_id: {chunks_with_id} 个")
print(f" - 未回填 chunk_id: {chunks_without_id} 个")
# 打印前 3 个切片的信息
print("\n前 3 个切片信息:")
for i, chunk in enumerate(output_chunks[:3]):
print(f"\n 切片 {i + 1}:")
print(f" chunk_id: {chunk.get('chunk_id', '无')}")
print(f" title: {chunk.get('title', '')}")
print(f" item_name: {chunk.get('item_name', '')}")
print(f" content: {chunk.get('content', '')[:50]}...")
# ----------------------------------------------------------------
# Step 5: 保存输出文件
# ----------------------------------------------------------------
with open(output_path, "w", encoding="utf-8") as f:
json.dump(result_state, f, ensure_ascii=False, indent=4)
print(f"\n已保存到: {output_path}")5.2 运行测试
# 进入项目目录
cd D:\develop\develop\workspace\pycharm\usage\shopkeeper_brain_v260213
# 执行测试
python -m knowledge.processor.import_process.nodes.import_milvus5.3 预期输出
正在读取输入文件: D:\...\temp\chunks_item_name_vector.json
读取到 15 个切片
开始执行 Milvus 导入...
2024-xx-xx [INFO] import_milvus - [step_1] 准备导入 15 条数据,向量维度: 1024
2024-xx-xx [INFO] import_milvus - [step_2] 创建集合: chunks_test
2024-xx-xx [INFO] import_milvus - 集合 chunks_test 创建成功
2024-xx-xx [INFO] import_milvus - [step_3] 执行插入
2024-xx-xx [INFO] import_milvus - 成功插入 15 条数据
处理完成,共 15 个切片
回填统计:
- 成功回填 chunk_id: 15 个
- 未回填 chunk_id: 0 个
前 3 个切片信息:
切片 1:
chunk_id: 450876234567890
title: 产品介绍
item_name: 福禄克15B+数字万用表
content: 福禄克15B+数字万用表是一款专业级测量仪器...
切片 2:
chunk_id: 450876234567891
title: 基本功能
item_name: 福禄克15B+数字万用表
content: 本产品支持直流电压、交流电压、电阻和电流测量...
切片 3:
chunk_id: 450876234567892
title: 安全须知
item_name: 福禄克15B+数字万用表
content: 使用前请仔细阅读本安全手册,确保输入值不超过...
已保存到: D:\...\temp\chunks_item_name_vector_ids.json5.4 验证 Milvus 数据
可以使用 Attu(Milvus 可视化工具)验证数据导入情况:
- 访问 http://localhost:7000 打开 Attu
- 连接到 Milvus(standalone:19530)
- 查看
chunks_test集合 - 确认数据条数和字段结构
6. 总结
6.1 节点功能概览
| 功能模块 | 说明 |
|---|---|
| 数据校验 | 检查 chunks 是否存在且非空,提取向量维度 |
| 连接管理 | 通过单例模式获取 Milvus 客户端连接 |
| 集合自动创建 | 检测集合存在性,不存在时自动创建 Schema 和索引 |
| Schema 定义 | 定义主键(auto_id)、标量字段和向量字段 |
| 索引配置 | 稠密向量用 AUTOINDEX,稀疏向量用 SPARSE_INVERTED_INDEX |
| 数据插入 | 批量插入切片数据到 Milvus 集合 |
| 主键回填 | 将 Milvus 返回的 chunk_id 写回业务数据 |
6.2 设计要点
幂等性设计
- 使用
has_collection()检测集合是否存在 - 只有不存在时才创建,避免重复创建
- 空数据时直接跳过,不报错
- 使用
Schema 与索引分离构建
_build_schema()负责字段定义_build_index_params()负责索引配置create_collection()一次性创建,保证原子性
双向量索引策略
- 稠密向量:AUTOINDEX(自动选择最优算法)
- 稀疏向量:SPARSE_INVERTED_INDEX(倒排索引 + DAAT_MAXSCORE 剪枝)
- 统一使用 IP 度量(配合归一化向量)
主键回填机制
- 使用
auto_id=True让 Milvus 自动生成主键 - 插入后从返回结果中提取
ids列表 - 使用
zip()逐一回填到对应 chunk - 转为字符串存储,便于跨系统传递
- 使用
企业痛点映射
| 痛点 | 传统方案 | AI Agent 导入方案 | 效率提升 |
|---|---|---|---|
| 向量数据手动导入耗时 | 逐条写入或命令行导入 | 批量 insert 自动创建集合 + 建索引 | 单次导入从 30min 降至 ~5s |
| 集合创建需手动操作 | 手动写 SQL 或 GUI 操作 | has_collection 幂等检测 + 自动创建 | 零手工操作,CI/CD 友好 |
| 主键需手动管理易冲突 | 自建 ID 生成器或 UUID | auto_id 自动生成 + 回填 | 主键冲突概率降为 ~0% |
| Schema 变更影响全局 | 改字段→删集合→重建 | enable_dynamic_fields 灵活 Schema | 字段变更无需重建 |
| 导入错误数据不可追溯 | 不知道哪些数据已导入 | chunk_id 回填到业务数据 | 100% 可追溯每条数据的 Milvus ID |
Remote & Agent 应用场景价值
Remote 场景价值:Milvus 作为远程向量数据库服务运行(通过
MILVUS_URL配置),团队成员无需本地部署。自动创建集合的幂等设计确保 CI/CD Pipeline 多次运行不会产生重复集合。Agent 落地场景:ImportMilvusNode 可封装为"数据入库 Agent"——Agent 接收向量化 chunks → 检查集合存在 → 创建/复用 → 批量插入 → 回填 chunk_id。在 Agent 工作流中,该节点是"处理链"的最后一个环节,输出带 chunk_id 的数据供下游 Agent(如知识图谱构建 Agent)使用。
Git Commit 对应
本节向量数据入库节点对应的提交记录(参考值,以实际版本为准):
<待补充 — 建议在项目仓库中搜索 "import_milvus.py" 相关提交>cd shopkeeper_brain
git log --oneline --all -- knowledge/processor/import_process/nodes/import_milvus.py