diff --git a/rag/svr/task_executor.py b/rag/svr/task_executor.py index 4f3ac2361c..b4d53808fe 100644 --- a/rag/svr/task_executor.py +++ b/rag/svr/task_executor.py @@ -934,8 +934,12 @@ async def run_dataflow(task: dict): time_cost = timer() - start_ts task_time_cost = timer() - task_start_ts + try: + ret = DocumentService.increment_chunk_num(doc_id, task_dataset_id, embedding_token_consumption, len(chunks), task_time_cost) + except Exception: + logging.exception("increment_chunk_num failed for doc %s", doc_id) + ret = None set_progress(task_id, prog=1.0, msg="Indexing done ({:.2f}s). Task done ({:.2f}s)".format(time_cost, task_time_cost)) - ret = DocumentService.increment_chunk_num(doc_id, task_dataset_id, embedding_token_consumption, len(chunks), task_time_cost) get_recording_context().save_func_return_value("DocumentService.increment_chunk_num", ret) logging.info("[Done], chunks({}), token({}), elapsed:{:.2f}".format(len(chunks), embedding_token_consumption, task_time_cost)) get_recording_context().record("dataflow_chunks", chunks) diff --git a/rag/svr/task_executor_refactor/dataflow_service.py b/rag/svr/task_executor_refactor/dataflow_service.py index 447a8f5590..ef4839019b 100644 --- a/rag/svr/task_executor_refactor/dataflow_service.py +++ b/rag/svr/task_executor_refactor/dataflow_service.py @@ -166,13 +166,21 @@ class DataflowService: time_cost = timer() - start_ts task_time_cost = timer() - task_start_ts - self._progress(prog=1.0, msg="Indexing done ({:.2f}s). Task done ({:.2f}s)".format(time_cost, task_time_cost)) - # Update document stats + # Update document stats (chunk counters) BEFORE marking the task + # as done, so the async _sync_progress loop never observes DONE + # with stale chunk_num=0. If the stats update fails, the task + # is still marked DONE — chunk counters are not critical to the + # parse result. if ctx.write_interceptor: ctx.write_interceptor.intercept("DocumentService.increment_chunk_num") else: - DocumentService.increment_chunk_num(doc_id, task_dataset_id, embedding_token_consumption, len(chunks), task_time_cost) + try: + DocumentService.increment_chunk_num(doc_id, task_dataset_id, embedding_token_consumption, len(chunks), task_time_cost) + except Exception: + logging.exception("increment_chunk_num failed for doc %s", doc_id) + + self._progress(prog=1.0, msg="Indexing done ({:.2f}s). Task done ({:.2f}s)".format(time_cost, task_time_cost)) logging.info("[Done], chunks({}), token({}), elapsed:{:.2f}".format(len(chunks), embedding_token_consumption, task_time_cost)) ctx.recording_context.record("dataflow_chunks", chunks)