From 583a17b37d13c5ff4715ec5b301da1117701430f Mon Sep 17 00:00:00 2001 From: Wenjie Zhang Date: Mon, 4 Aug 2025 20:15:39 +0800 Subject: [PATCH] =?UTF-8?q?feat(knowledge):=20=E6=B7=BB=E5=8A=A0=E6=96=87?= =?UTF-8?q?=E4=BB=B6=E5=A4=84=E7=90=86=E9=98=9F=E5=88=97=E7=AE=A1=E7=90=86?= =?UTF-8?q?=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 添加文件处理队列管理功能,包括添加/移除文件到处理队列的方法 添加状态检查方法修复异常的processing状态 确保文件处理状态的一致性 --- src/knowledge/chroma_kb.py | 3 ++ src/knowledge/knowledge_base.py | 88 +++++++++++++++++++++++++++++++++ src/knowledge/lightrag_kb.py | 3 ++ src/knowledge/milvus_kb.py | 9 ++++ 4 files changed, 103 insertions(+) diff --git a/src/knowledge/chroma_kb.py b/src/knowledge/chroma_kb.py index bbecf6c0..fd91d2f0 100644 --- a/src/knowledge/chroma_kb.py +++ b/src/knowledge/chroma_kb.py @@ -186,6 +186,7 @@ class ChromaKB(KnowledgeBase): self.files_meta[file_id] = file_record self._save_metadata() + self._add_to_processing_queue(file_id) try: # 根据内容类型处理内容 if content_type == "file": @@ -222,6 +223,8 @@ class ChromaKB(KnowledgeBase): self.files_meta[file_id]["status"] = "failed" self._save_metadata() file_record['status'] = "failed" + finally: + self._remove_from_processing_queue(file_id) processed_items_info.append(file_record) diff --git a/src/knowledge/knowledge_base.py b/src/knowledge/knowledge_base.py index d0580b3c..458d55d8 100644 --- a/src/knowledge/knowledge_base.py +++ b/src/knowledge/knowledge_base.py @@ -28,6 +28,10 @@ class KBOperationError(KnowledgeBaseException): class KnowledgeBase(ABC): """知识库抽象基类,定义统一接口""" + # 类级别的处理队列,跟踪所有正在处理的文件 + _processing_files = set() + _processing_lock = None + def __init__(self, work_dir: str): """ 初始化知识库 @@ -35,9 +39,16 @@ class KnowledgeBase(ABC): Args: work_dir: 工作目录 """ + import threading + self.work_dir = work_dir self.databases_meta: dict[str, dict] = {} self.files_meta: dict[str, dict] = {} + + # 初始化类级别的锁 + if KnowledgeBase._processing_lock is None: + KnowledgeBase._processing_lock = threading.Lock() + os.makedirs(work_dir, exist_ok=True) # 自动加载元数据 @@ -208,6 +219,9 @@ class KnowledgeBase(ABC): meta = self.databases_meta[db_id].copy() meta["db_id"] = db_id + # 检查并修复异常的processing状态 + self._check_and_fix_processing_status(db_id) + # 获取文件信息 db_files = {} for file_id, file_info in self.files_meta.items(): @@ -238,6 +252,9 @@ class KnowledgeBase(ABC): """ databases = [] for db_id, meta in self.databases_meta.items(): + # 检查并修复异常的processing状态 + self._check_and_fix_processing_status(db_id) + db_dict = meta.copy() db_dict["db_id"] = db_id @@ -264,6 +281,77 @@ class KnowledgeBase(ABC): return {"databases": databases} + @classmethod + def _add_to_processing_queue(cls, file_id: str) -> None: + """ + 将文件添加到处理队列 + + Args: + file_id: 文件ID + """ + with cls._processing_lock: + cls._processing_files.add(file_id) + logger.debug(f"Added file {file_id} to processing queue") + + @classmethod + def _remove_from_processing_queue(cls, file_id: str) -> None: + """ + 从处理队列中移除文件 + + Args: + file_id: 文件ID + """ + with cls._processing_lock: + cls._processing_files.discard(file_id) + logger.debug(f"Removed file {file_id} from processing queue") + + @classmethod + def _is_file_in_processing_queue(cls, file_id: str) -> bool: + """ + 检查文件是否在处理队列中 + + Args: + file_id: 文件ID + + Returns: + bool: 文件是否在处理队列中 + """ + with cls._processing_lock: + return file_id in cls._processing_files + + + + def _check_and_fix_processing_status(self, db_id: str) -> None: + """ + 检查并修复异常的processing状态 + 如果文件状态为processing但实际不在处理队列中,则修改为error状态 + + Args: + db_id: 数据库ID + """ + try: + status_changed = False + + # 检查该数据库下所有processing状态的文件 + for file_id, file_info in self.files_meta.items(): + if (file_info.get("database_id") == db_id and + file_info.get("status") == "processing"): + + # 检查文件是否真的在处理队列中 + if not self._is_file_in_processing_queue(file_id): + logger.warning(f"File {file_id} has processing status but is not in processing queue, marking as error") + self.files_meta[file_id]["status"] = "error" + self.files_meta[file_id]["error"] = "Processing interrupted - file not found in processing queue" + status_changed = True + + # 如果有状态变更,保存元数据 + if status_changed: + self._save_metadata() + logger.info(f"Fixed processing status for database {db_id}") + + except Exception as e: + logger.error(f"Error checking processing status for database {db_id}: {e}") + @abstractmethod async def delete_file(self, db_id: str, file_id: str) -> None: """ diff --git a/src/knowledge/lightrag_kb.py b/src/knowledge/lightrag_kb.py index 8ff6df73..db6ea6e3 100644 --- a/src/knowledge/lightrag_kb.py +++ b/src/knowledge/lightrag_kb.py @@ -164,6 +164,7 @@ class LightRagKB(KnowledgeBase): self.files_meta[file_id] = file_record self._save_metadata() + self._add_to_processing_queue(file_id) try: # 根据内容类型处理内容 if content_type == "file": @@ -195,6 +196,8 @@ class LightRagKB(KnowledgeBase): self._save_metadata() file_record['status'] = "failed" file_record['error'] = error_msg + finally: + self._remove_from_processing_queue(file_id) processed_items_info.append(file_record) diff --git a/src/knowledge/milvus_kb.py b/src/knowledge/milvus_kb.py index d1c5c1ea..9b0401ae 100644 --- a/src/knowledge/milvus_kb.py +++ b/src/knowledge/milvus_kb.py @@ -256,6 +256,9 @@ class MilvusKB(KnowledgeBase): self._save_metadata() file_record["file_id"] = file_id + + # 添加到处理队列 + self._add_to_processing_queue(file_id) try: if content_type == "file": @@ -292,6 +295,8 @@ class MilvusKB(KnowledgeBase): self.files_meta[file_id]["status"] = "done" self._save_metadata() file_record['status'] = "done" + # 从处理队列中移除 + self._remove_from_processing_queue(file_id) except Exception as e: logger.error(f"处理{content_type} {item} 失败: {e}, {traceback.format_exc()}") @@ -299,6 +304,10 @@ class MilvusKB(KnowledgeBase): self.files_meta[file_id]["status"] = "failed" self._save_metadata() file_record['status'] = "failed" + # 从处理队列中移除 + self._remove_from_processing_queue(file_id) + finally: + pass processed_items_info.append(file_record)