资讯动态

企业实战:Milvues向量数据库实践

发布时间:2026/8/19 11:38:01 来源:尧图企业网站定制
目录1 存入 Milvus (node_import_milvus)2 节点作用与实现思路3 导入与配置4 核心辅助函数5 主流程定义6 步骤 1: 检查输入7 步骤 2: 准备集合8 步骤 3: 清理旧数据9 步骤 4: 插入数据10 单元测试11 导入数据节点实现与测试11.1 进行主图调用测试 (main_graph)1 存入 Milvus (node_import_milvus)文件:app/import_process/node_import_milvus.py2 节点作用与实现思路节点作用: 数据加载流程的终点负责将处理好的结构化数据切片内容、元数据、向量持久化存储到向量数据库中构建可供即时查询的索引。实现思路:幂等性设计: 在插入新数据前根据item_name或文件 ID 清理旧数据防止重复导入导致的数据污染。Schema 适配: 严格按照 Milvus 集合的 Schema 定义主键、Dense字段、Sparse字段、JSON元数据字段组织数据确保插入成功率。混合索引构建: 确保存入的数据能够支持 Milvus 的 Hybrid SearchDense Sparse 加权最大化检索效果。3 导入与配置目的: 导入必要的库如pymilvus和项目工具类配置 Milvus 集合名称。关键点:CHUNKS_COLLECTION_NAME: 从环境变量获取集合名。add_running_task: 记录任务执行状态。# 导入Milvus相关依赖frompymilvusimportDataType# 导入自定义模块fromapp.import_process.agent.stateimportImportGraphStatefromapp.clients.milvus_utilsimportget_milvus_clientfromapp.utils.task_utilsimportadd_running_task,add_done_taskfromapp.core.loggerimportlogger,node_log,step_logfromapp.conf.milvus_configimportmilvus_config# 从配置文件读取切片集合名称与配置解耦便于环境切换CHUNKS_COLLECTION_NAMEmilvus_config.chunks_collection4 核心辅助函数功能: 处理幂等性清理删除旧数据和字符串转义。fromapp.utils.escape_milvus_string_utilsimportescape_milvus_string5 主流程定义函数:node_import_milvus逻辑:Step 1: 检查输入 (step_1_check_input)。Step 2: 准备环境 (step_2_prepare_collection)。Step 3: 清理旧数据 (step_3_clean_old_data)。Step 4: 插入数据 (step_4_insert_data)。 1. 检查数据 chunks是否存在 2. 前置准备工作 准备 milvus的集合和字段等 3. 删除旧数据 4. 查询chunks的数据即可 node_log(node_import_milvus)defnode_import_milvus(state:ImportGraphState)-ImportGraphState: 节点: 导入向量库 (node_import_milvus) 为什么叫这个名字: 将处理好的向量数据写入 Milvus 数据库。 # 准备日志和任务列表add_running_task(state[task_id],node_import_milvus)# 1. 检查数据 chunks是否存在chunksstate.get(chunks)ifnotchunks:logger.error(node_import_milvus: chunks数据不存在)raiseValueError(node_import_milvus: chunks数据不存在)# 2. 前置准备工作 创建 Milvus 集合和字段milvus_clientget_milvus_client()step_2_prepare_collection(milvus_client)# 3. 删除旧数据step_3_delete_old_data(milvus_client,state[item_name])# 4. 插入chunks的数据即可with_id_chunksstep_4_insert_collections(milvus_client,chunks)state[chunks]with_id_chunks add_done_task(state[task_id],node_import_milvus)returnstate6 步骤 1: 检查输入功能: 验证chunks是否存在并提取dense_vector维度和item_name。# 1. 检查数据 chunks是否存在chunksstate.get(chunks)ifnotchunks:logger.error(node_import_milvus: chunks数据不存在)raiseValueError(node_import_milvus: chunks数据不存在)7 步骤 2: 准备集合功能: 获取 Milvus 客户端如果集合不存在则创建。step_log(step_2_prepare_collection)defstep_2_prepare_collection(milvus_client): 准备和创建chunks对应的集合 :param milvus_client: :return: # 2. 判断是否存在集合表存在创建集合表ifnotmilvus_client.has_collection(collection_namemilvus_config.chunks_collection):# 创建集合# 3.1. 创建集合对应的列的信息schemamilvus_client.create_schema(auto_idTrue,# 主键自增长enable_dynamic_fieldTrue,# 动态字段)# 3.2. Add fields to schema# pk file_title item_name dense_vector sparse_vectorschema.add_field(field_namechunk_id,datatypeDataType.INT64,is_primaryTrue,auto_idTrue)schema.add_field(field_namefile_title,datatypeDataType.VARCHAR,max_length65535)schema.add_field(field_nameitem_name,datatypeDataType.VARCHAR,max_length65535)schema.add_field(field_namecontent,datatypeDataType.VARCHAR,max_length65535)schema.add_field(field_nametitle,datatypeDataType.VARCHAR,max_length65535)schema.add_field(field_nameparent_title,datatypeDataType.VARCHAR,max_length65535)schema.add_field(field_namepart,datatypeDataType.INT8)schema.add_field(field_namedense_vector,datatypeDataType.FLOAT_VECTOR,dim1024)schema.add_field(field_namesparse_vector,datatypeDataType.SPARSE_FLOAT_VECTOR)# 3.3 查询快配置索引index_paramsmilvus_client.prepare_index_params()index_params.add_index(field_namedense_vector,# 给哪个列创建索引 稠密index_namedense_vector_index,# 索引的名字index_typeHNSW,# 配置查找所用的算法metric_typeCOSINE,# 配置向量匹配和对比的 IP COSINEparams{M:32,# Maximum number of neighbors each node can connect to in the graphefConstruction:300},# or DAAT_WAND or TAAT_NAIVE) 10000 M 16 efConstruction 200 50000 M 32 efConstruction 300 100000 M 64 efConstruction 400 M:图中每个节点在层次结构的每个层级所能拥有的最大边数或连接数。M 越高图的密度就越大搜索结果的召回率和准确率也就越高因为有更多的路径可以探索但同时也会消耗更多内存并由于连接数的增加而减慢插入时间。如上图所示M 5表示 HNSW 图中的每个节点最多与 5 个其他节点直接相连。这就形成了一个中等密度的图结构节点有多条路径到达其他节点。 efConstruction:索引构建过程中考虑的候选节点数量。efConstruction 越高图的质量越好但需要更多时间来构建。 index_params.add_index(field_namesparse_vector,# Name of the vector field to be indexedindex_typeSPARSE_INVERTED_INDEX,# Type of the index to createindex_namesparse_vector_index,# Name of the index to createmetric_typeIP,# Metric type used to measure similarity# 只计算可能得高分的向量跳过大量的 0params{inverted_index_algo:DAAT_MAXSCORE},# Algorithm used for building and querying the index)milvus_client.create_collection(collection_namemilvus_config.chunks_collection,schemaschema,# 字段index_paramsindex_params# 索引)returnmilvus_client8 步骤 3: 清理旧数据功能: 根据item_name删除已存在的切片确保幂等性。step_log(step_3_delete_old_data)defstep_3_delete_old_data(milvus_client,item_name): 删除旧数据 根据item_name删除 :param milvus_client: :param item_name: :return: milvus_client.delete(collection_nameCHUNKS_COLLECTION_NAME,filterfitem_name{item_name})# 调用 load_collection() 会触发 Milvus 重新加载集合数据、刷新索引、清理已标记删除的数据确保删除操作真正生效避免新旧数据混杂导致检索错误。milvus_client.load_collection(collection_nameCHUNKS_COLLECTION_NAME)9 步骤 4: 插入数据功能: 移除临时chunk_id批量插入数据并回填生成的 ID。step_log(step_4_insert_collections)defstep_4_insert_collections(milvus_client,chunks): 插入集合的数据 :param milvus_client :param chunks: :return: chunks - 主键回显 insert_resultmilvus_client.insert(collection_nameCHUNKS_COLLECTION_NAME,datachunks)# 成功插入了几条insert_countinsert_result.get(insert_count,0)logger.info(f完成了数据插入成功插入了{insert_count}条数据)# 获取回显的idsidsinsert_result.get(ids,[])ifidsandlen(ids)len(chunks):forindex,chunkinenumerate(chunks):chunk[chunk_id]ids[index]returnchunks10 单元测试您可以在node_import_milvus.py文件底部直接运行以下测试代码if__name____main__:# --- 单元测试 ---# 目的验证 Milvus 导入节点的完整流程包括连接、创建集合、清理旧数据和插入新数据。importsysimportosfromdotenvimportload_dotenv# 加载环境变量 (自动寻找项目根目录的 .env)current_diros.path.dirname(os.path.abspath(__file__))project_rootos.path.dirname(os.path.dirname(current_dir))load_dotenv(os.path.join(project_root,.env))# 构造测试数据dim1024test_state{task_id:test_milvus_task,item_name:测试项目_Milvus,chunks:[{content:Milvus 测试文本 1,title:测试标题,item_name:测试项目_Milvus,# 必须有 item_name用于幂等清理parent_title:test.pdf,part:1,file_title:test.pdf,dense_vector:[0.1]*dim,# 模拟 Dense Vectorsparse_vector:{1:0.5,10:0.8}# 模拟 Sparse Vector},{content:Milvus 测试文本 2,title:测试标题2,item_name:测试项目_Milvus2,# 必须有 item_name用于幂等清理parent_title:test.pdf2,part:1,file_title:test.pdf2,dense_vector:[0.1]*dim,# 模拟 Dense Vectorsparse_vector:{1:0.5,10:0.8}# 模拟 Sparse Vector}]}print(正在执行 Milvus 导入节点测试...)try:# 检查必要的环境变量ifnotos.getenv(MILVUS_URL):print(❌ 未设置 MILVUS_URL无法连接 Milvus)elifnotos.getenv(CHUNKS_COLLECTION):print(❌ 未设置 CHUNKS_COLLECTION)else:# 执行节点函数result_statenode_import_milvus(test_state)# 验证结果chunksresult_state.get(chunks,[])ifchunksandchunks[0].get(chunk_id):print(f✅ Milvus 导入测试通过生成 ID:{chunks[0][chunk_id]})else:print(❌ 测试失败未能获取 chunk_id)exceptExceptionase:print(f❌ 测试失败:{e})11 导入数据节点实现与测试11.1 进行主图调用测试 (main_graph)主图添加测试代码进行流程完成测试if__name____main__:fromapp.utils.path_utilimportPROJECT_ROOTimportos# 全流程测试验证PDF导入→Milvus入库→KG导入完整链路logger.info( 开始执行知识图谱导入全流程测试 )# 1. 构造测试文件路径复用你项目的doc目录和pdf2md测试文件一致test_pdf_nameos.path.join(doc,万用表RS-12的使用.pdf)test_pdf_pathos.path.join(PROJECT_ROOT,test_pdf_name)# 2. 构造输出目录存放MD/图片等中间文件test_output_diros.path.join(PROJECT_ROOT,output)os.makedirs(test_output_dir,exist_okTrue)# 不存在则创建# 3. 校验测试PDF文件是否存在ifnotos.path.exists(test_pdf_path):logger.error(f全流程测试失败测试PDF文件不存在路径{test_pdf_path})logger.info(请检查文件路径或手动将测试文件放入项目根目录的doc文件夹中)else:# 4. 构造测试状态贴合实际业务入参开启PDF解析开关test_stateImportGraphState({task_id:test_kg_import_workflow_001,# 测试任务IDuser_id:test_user,# 测试用户IDlocal_file_path:test_pdf_path,# 测试PDF文件路径local_dir:test_output_dir,# 中间文件输出目录is_pdf_read_enabled:False,# 开启PDF解析核心开关is_md_read_enabled:False# 关闭MD解析})try:logger.info(f测试任务启动PDF文件路径{test_pdf_path})logger.info(f中间文件输出目录{test_output_dir})logger.info(开始执行全流程节点依次执行entry→pdf2md→md_img→split→item_name→embedding→milvus→kg)# 5. 执行LangGraph全流程流式执行打印节点执行进度final_stateNoneforstepinkb_import_app.stream(test_state,stream_modevalues):# 打印当前执行完成的节点流式输出更直观current_nodelist(step.keys())[-1]ifstepelse未知节点logger.info(f✅ 节点执行完成{current_node})final_statestep# 保存最终状态# 6. 全流程执行完成结果预览和核心指标打印iffinal_state:logger.info(-*80)logger.info( 全流程测试执行成功核心结果预览 )# 提取核心结果指标chunksfinal_state.get(chunks,[])chunk_countlen(chunks)md_contentfinal_state.get(md_content,)[:150]# MD内容前150字符has_embeddingall(dense_vectorincandsparse_vectorincforcinchunks)ifchunkselseFalsehas_chunk_idall(chunk_idincforcinchunks)ifchunkselseFalsekg_idfinal_state.get(kg_id,未生成)# KG导入生成的ID按实际业务字段调整# 打印核心指标logger.info(f PDF转MD内容预览前150字符{md_content}...)logger.info(f 文档切分总切片数{chunk_count})logger.info(f 所有切片是否完成向量化{是ifhas_embeddingelse否})logger.info(f️ 所有切片是否完成Milvus入库含chunk_id{是ifhas_chunk_idelse否})logger.info(f 知识图谱导入ID{kg_id})logger.info(f 最终状态包含的核心键{list(final_state.keys())})logger.info(-*80)exceptExceptionase:logger.exception(f 全流程测试运行失败 )logger.info( 知识图谱导入全流程测试结束 )输出日志—识别完成: HAK 180 烫金机—[Stub]执行节点node_bge_embedding—开始 BGE-M3 向量化处理 (node_bge_embedding)—成功获取第 1-5 项的嵌入。成功获取第 6-6 项的嵌入。— 向量化处理完成共处理 6 条数据 —[Stub] 执行节点node_import_milvus—开始导入 Milvus—从数据中检测到向量维度为: 1024从数据中检测到 item_name 为: HAK 180 烫金机正在连接到 Milvus 并准备集合 ‘kb_chunks’…幂等清理正在删除 Milvus 集合 ‘kb_chunks’ 中 item_nameHAK 180 烫金机 的旧切片…幂等清理完成item_nameHAK 180 烫金机准备了 6 条数据正在执行插入操作…成功插入 6 条数据。[Stub]执行节点node_import_kg 流程执行结束 ✅ 任务ID: full_graph_test_001✅ 文件标题: hak180产品安全手册✅ 识别商品: HAK 180 烫金机✅ 切片数量: 6✅ 向量化结果: 成功 (检测到 6 条带向量的数据)效果体现

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价