diff --git a/app/utils/knowledge_utils.py b/app/utils/knowledge_utils.py index 335ffc0..2e202a8 100644 --- a/app/utils/knowledge_utils.py +++ b/app/utils/knowledge_utils.py @@ -1,11 +1,15 @@ import asyncio +import traceback from typing import List from fastapi import UploadFile +from llama_index.core import SimpleDirectoryReader, Document +from app.config.env import env from app.model.FileModel import FileSaveService from app.model.KnowledgeDoc import KnowledgeDocModel, KnowledgeDocService from app.utils.db_utils import AsyncSessionDep, async_session +from app.utils.milvus_utils import milvus_service class KnowledgeService: @@ -19,6 +23,7 @@ class KnowledgeService: async with async_session() as session: return await FileSaveService.saveFile(session=session, file=file, filename=file.filename, file_record={}) + # 将 file_dict_list 保存为 doc_cls_list async def save_knowledge_doc_list(self, session: AsyncSessionDep, file_dict_list: List[dict], kb_code: str): # 先插入文档对象 @@ -31,11 +36,60 @@ class KnowledgeService: return insert_cls_list - async def process_file_dict(self, file_dict: dict, kb_code: str): - print("process start", file_dict['name']) - await asyncio.sleep(3) - print("process end", file_dict['name']) - return {} + # 将doc_cls对应的文件创建索引保存到Milvus中 + async def process_doc_cls(self, doc_cls: KnowledgeDocModel): + try: + # 找到文件的保存路径 + file_path = doc_cls.path.replace(env.file_public_path, env.file_save_path) + # 加载文件内容 + documents = await self.read_document_async(file_path=file_path) + # 处理文档块的元信息 + documents = [ + Document( + doc_id=doc_cls.id, + text=doc.text, + metadata={ + **doc.metadata, + **doc_cls.model_dump(), + } + ) + for doc in documents + ] + # 为文档创建索引 + try: + await milvus_service.async_create_index_from_documents(documents) + except Exception as e: + raise Exception(f"为文档创建索引失败:{str(e)}") + # 更新文档状态 + async with async_session() as session: + await KnowledgeDocService.item_update(session, row_dict={"id": doc_cls.id, "status": "success"}) + return {} + except Exception as e: + print(e) + traceback.print_exc() + # 更新文档状态 + async with async_session() as session: + await KnowledgeDocService.item_update( + session, + row_dict={ + "id": doc_cls.id, + "status": "fail", + "error": str(e) + } + ) + return {} + + # 异步读取文件 + async def read_document_async(self, file_path: str) -> List[Document]: + try: + def _read_document(): + reader = SimpleDirectoryReader(input_files=[file_path]) + return reader.load_data() + + documents = await asyncio.to_thread(_read_document) + return documents + except Exception as e: + raise Exception(f"读取文档失败:{str(e)}") knowledge_service = KnowledgeService()