Skip to content

导入向量数据库节点

本文档详细介绍知识库导入流程中的导入向量数据库节点(ImportMilvusNode),该节点负责将向量化后的切片数据批量导入 Milvus 向量数据库,包括集合创建、Schema 定义、索引构建和数据插入等完整流程。


学习理念:向量数据入库是导入流程的"终点"——前面的切片、向量化、商品名识别等所有工作,最终都要写入 Milvus 才能被检索到。ImportMilvusNode 的核心是幂等性设计(集合不存在才创建)+ Schema + 索引分离构建 + 主键自动回填

海外对标:ImportMilvusNode 的"Schema 定义 → 索引构建 → 批量插入 → 主键回填"流程对标 Pinecone 的 create_index + upsert API,以及 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自动 IDSchema 配置,让 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 本章目标

通过本章学习,你将掌握:

  1. Milvus 集合管理:学会创建集合、定义 Schema、配置索引
  2. Schema 设计:理解标量字段与向量字段的定义方式
  3. 索引类型选择:掌握 AUTOINDEX 和 SPARSE_INVERTED_INDEX 的适用场景
  4. 数据插入与主键回填:学会批量插入数据并获取自动生成的主键
  5. 幂等性设计:理解"集合不存在才创建"的设计模式

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: 图状态字典,包含向量化后的切片数据

处理逻辑

  1. 从 state 中获取 "chunks" 字段
  2. 如果 chunks 为空,记录警告并直接返回(幂等设计)
  3. 从第一条切片中提取 dense_vector 的维度
  4. 如果向量维度为 0,抛出 MilvusError 异常

输出

  • 验证通过的切片列表 chunks
  • 向量维度 vector_dim

代码片段

🟡 【P1 看注释就行】 Step 1 校验代码——chunks 为空则跳过(幂等设计),_get_vector_dim 从首条 chunk 提取稠密向量维度。

python
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 dim
Step2:连接 Milvus

获取 Milvus 客户端连接实例。

处理逻辑

  1. 调用 get_milvus_client() 获取客户端单例
  2. 客户端根据环境变量 MILVUS_URL 配置连接地址
  3. 首次调用时建立连接,后续调用复用

环境变量配置

代码片段

🟡 【P1 看注释就行】 Step 2 连接代码——get_milvus_client() 获取单例客户端。

python
try:
    client = get_milvus_client()
    collection_name = config.chunks_collection  # 从配置获取集合名
Step3:检查并创建集合

检查目标集合是否存在,不存在则创建。

处理逻辑

  1. 调用 client.has_collection() 检查集合是否存在
  2. 如果不存在,调用 _create_collection() 创建集合
  3. 创建过程包括:构建 Schema → 构建索引参数 → 创建集合

幂等性保证

  • 只有在集合不存在时才创建
  • 重复执行不会报错或创建重复集合

代码片段

🟡 【P1 看注释就行】 Step 3 幂等创建——has_collection 检查 → 不存在才创建。避免重复运行报错。

python
# 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

定义集合的字段结构。

字段定义

  1. chunk_id: INT64 主键,自动生成
  2. content: VARCHAR,存储切片文本内容
  3. title: VARCHAR,存储切片标题
  4. parent_title: VARCHAR,存储父级标题
  5. part: INT8,存储切片在同标题下的序号
  6. file_title: VARCHAR,存储源文件标题
  7. item_name: VARCHAR,存储商品名称
  8. sparse_vector: SPARSE_FLOAT_VECTOR,存储稀疏向量
  9. 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 允许动态添加字段。

python
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 schema
Step5:构建索引参数

为稠密向量和稀疏向量分别配置索引。

索引配置

  1. 稠密向量索引

    • 类型:AUTOINDEX(自动选择最优算法)
    • 度量:IP(内积)
  2. 稀疏向量索引

    • 类型:SPARSE_INVERTED_INDEX(倒排索引)
    • 度量:IP(内积)
    • 算法:DAAT_MAXSCORE(动态剪枝,高效查询)

代码片段

🔥 【P0 必须要学】 Step 5 索引配置——稠密向量用 AUTOINDEX(自动选择 HNSW/IVF 等),稀疏向量用 SPARSE_INVERTED_INDEX + DAAT_MAXSCORE 剪枝。两者都用 IP 度量(配合 L2 归一化后 IP = 余弦相似度)。

python
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_params
Step6:创建集合

组合 Schema 和索引参数,创建集合。

处理逻辑

  1. 调用 _build_schema() 构建 Schema
  2. 调用 _build_index_params() 构建索引参数
  3. 调用 client.create_collection() 创建集合

代码片段

🟡 【P1 看注释就行】 Step 6 创建集合——组合 schema + index_params 一次性创建。注意 _build_schema + _build_index_params 的职责分离。

python
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: 包含向量的切片数据列表

处理逻辑

  1. 调用 client.insert() 批量插入数据
  2. Milvus 会自动为���条数据生成 chunk_id
  3. 插入结果包含 insert_countids

注意事项

  • 插入时不需要提供 chunk_id 字段(auto_id=True)
  • Milvus 会自动处理数据类型转换

代码片段

🟡 【P1 看注释就行】 Step 7 插入数据——client.insert(collection_name, data=chunks) 批量插入。auto_id=True 时不需要传 chunk_id。

python
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 返回的主键回填到原始切片数据中。

处理逻辑

  1. 从插入结果中获取 ids 列表
  2. 验证返回的 ID 数量与切片数量一致
  3. 使用 zip() 将 ID 逐一回填到对应的 chunk
  4. 将 ID 转换为字符串存储(便于后续 JSON 序列化和跨系统传递)

代码片段

🔥 【P0 必须要学】 Step 8 主键回填是关键设计result.get("ids", []) 获取 Milvus 自动生成的 ID,zip(chunks, inserted_ids) 逐一回填,转 str(chunk_id) 便于 JSON 序列化。回填失败则无法关联知识图谱

python
# 回填 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:更新状态并返回

将处理结果写回状态并返回。

处理逻辑

  1. 将已回填 chunk_id 的 chunks 写回 state
  2. 返回更新后的 state

代码片段

🟢 【P2 后面可以查】 Step 9 更新 state——state["chunks"] = chunks

python
# 5. 更新 chunks 状态
state["chunks"] = chunks

return state

4.4 代码实现

🔥 【P0 必须要学】 ImportMilvusNode 是导入流程的"终点"——所有处理完的切片最终落库到 Milvus。重点理解:主流程(校验→集合创建→插入→回填)+ 幂等设计(has_collection)+ Schema/索引分离构建。_insert_and_backfill_idszip 回填是数据关联的关键。

python
"""
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 回填。

python
# ================================================================== #
#                        兼容 & 测试                                   #
# ================================================================== #

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 运行测试

bash
# 进入项目目录
cd D:\develop\develop\workspace\pycharm\usage\shopkeeper_brain_v260213

# 执行测试
python -m knowledge.processor.import_process.nodes.import_milvus

5.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.json

5.4 验证 Milvus 数据

可以使用 Attu(Milvus 可视化工具)验证数据导入情况:

  1. 访问 http://localhost:7000 打开 Attu
  2. 连接到 Milvus(standalone:19530)
  3. 查看 chunks_test 集合
  4. 确认数据条数和字段结构

6. 总结

6.1 节点功能概览

功能模块说明
数据校验检查 chunks 是否存在且非空,提取向量维度
连接管理通过单例模式获取 Milvus 客户端连接
集合自动创建检测集合存在性,不存在时自动创建 Schema 和索引
Schema 定义定义主键(auto_id)、标量字段和向量字段
索引配置稠密向量用 AUTOINDEX,稀疏向量用 SPARSE_INVERTED_INDEX
数据插入批量插入切片数据到 Milvus 集合
主键回填将 Milvus 返回的 chunk_id 写回业务数据

6.2 设计要点

  1. 幂等性设计

    • 使用 has_collection() 检测集合是否存在
    • 只有不存在时才创建,避免重复创建
    • 空数据时直接跳过,不报错
  2. Schema 与索引分离构建

    • _build_schema() 负责字段定义
    • _build_index_params() 负责索引配置
    • create_collection() 一次性创建,保证原子性
  3. 双向量索引策略

    • 稠密向量:AUTOINDEX(自动选择最优算法)
    • 稀疏向量:SPARSE_INVERTED_INDEX(倒排索引 + DAAT_MAXSCORE 剪枝)
    • 统一使用 IP 度量(配合归一化向量)
  4. 主键回填机制

    • 使用 auto_id=True 让 Milvus 自动生成主键
    • 插入后从返回结果中提取 ids 列表
    • 使用 zip() 逐一回填到对应 chunk
    • 转为字符串存储,便于跨系统传递

企业痛点映射

痛点传统方案AI Agent 导入方案效率提升
向量数据手动导入耗时逐条写入或命令行导入批量 insert 自动创建集合 + 建索引单次导入从 30min 降至 ~5s
集合创建需手动操作手动写 SQL 或 GUI 操作has_collection 幂等检测 + 自动创建零手工操作,CI/CD 友好
主键需手动管理易冲突自建 ID 生成器或 UUIDauto_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" 相关提交>
bash
cd shopkeeper_brain
git log --oneline --all -- knowledge/processor/import_process/nodes/import_milvus.py

OPC 超级个体实战指南