From 23b20a098ac6393c779f597c6fc792aeae1a2d40 Mon Sep 17 00:00:00 2001 From: Wang Qi Date: Thu, 6 Aug 2026 17:01:53 +0800 Subject: [PATCH] Fix ragflow server hung after parsing a big file (#17936) --- api/ragflow_server.py | 14 ++++++++------ rag/utils/redis_conn.py | 13 ++++++++++++- 2 files changed, 20 insertions(+), 7 deletions(-) diff --git a/api/ragflow_server.py b/api/ragflow_server.py index f9c9d62c08..47409e1348 100644 --- a/api/ragflow_server.py +++ b/api/ragflow_server.py @@ -58,17 +58,19 @@ def update_progress(): redis_lock = RedisDistributedLock("update_progress", lock_value=lock_value, timeout=60) logging.info(f"update_progress lock_value: {lock_value}") while not stop_event.is_set(): + acquired = False try: - if redis_lock.acquire(): + acquired = redis_lock.acquire() + if acquired: DocumentService.update_progress() - redis_lock.release() except Exception: logging.exception("update_progress exception") finally: - try: - redis_lock.release() - except Exception: - logging.exception("update_progress exception") + if acquired: + try: + redis_lock.release() + except Exception: + logging.exception("update_progress exception") stop_event.wait(6) diff --git a/rag/utils/redis_conn.py b/rag/utils/redis_conn.py index 28075c125e..70a275648c 100644 --- a/rag/utils/redis_conn.py +++ b/rag/utils/redis_conn.py @@ -514,7 +514,12 @@ class RedisDB: Do following atomically: Delete a key if its value is equals to the given one, do nothing otherwise. """ - return bool(self.lua_delete_if_equal(keys=[key], args=[expected_value], client=self.REDIS)) + try: + return bool(self.lua_delete_if_equal(keys=[key], args=[expected_value], client=self.REDIS)) + except Exception as e: + logging.warning("RedisDB.delete_if_equal got exception: %s", str(e)) + self.__open__() + return False def delete(self, key) -> bool: try: @@ -537,15 +542,21 @@ class RedisDistributedLock: else: self.lock_value = str(uuid.uuid4()) self.timeout = timeout + self.blocking_timeout = blocking_timeout self.lock = Lock(REDIS_CONN.REDIS, lock_key, timeout=timeout, blocking_timeout=blocking_timeout) + def _refresh_lock(self): + self.lock = Lock(REDIS_CONN.REDIS, self.lock_key, timeout=self.timeout, blocking_timeout=self.blocking_timeout) + def acquire(self): REDIS_CONN.delete_if_equal(self.lock_key, self.lock_value) + self._refresh_lock() return self.lock.acquire(token=self.lock_value) async def spin_acquire(self): REDIS_CONN.delete_if_equal(self.lock_key, self.lock_value) while True: + self._refresh_lock() if self.lock.acquire(token=self.lock_value): break await asyncio.sleep(10)