From 635045cc30da1cb13f00d61da66dcdc2e1ea32a0 Mon Sep 17 00:00:00 2001 From: lili Date: Tue, 23 Jun 2026 21:14:07 -0700 Subject: [PATCH] =?UTF-8?q?feat(knowledge):=20=E8=A1=A5=20parse=20?= =?UTF-8?q?=E8=BD=AE=E8=AF=A2=E6=8E=A8=E8=BF=9B=20worker=EF=BC=8C=E6=89=93?= =?UTF-8?q?=E9=80=9A=E6=91=84=E5=85=A5=20parsing=E2=86=92searchable=20?= =?UTF-8?q?=E7=BB=88=E6=80=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 处理链 startParseDocuments 触发 RAGFlow 解析后置 parsing/processing 即返回,RAGFlow 解析+embedding+索引为异步;此前无 worker 轮询其完成态(pollDocumentStatuses 仅 IT 手动调用),文档永久卡 parsing 无法 searchable。 新增 MuseKnowledgeParseStatusPollWorker(@Scheduled,默认关闭开关 muse.knowledge.parse-poll-worker.enabled): - 跨租户捞 parsing 任务(selectParsingForPoll + executeIgnore),逐个切回任务租户轮询 RAGFlow 文档状态 - museStatus=searchable(RAGFlow run=done)→ markRagflowParseCompleted(task completed + version processingStatus=searchable → 文档 isSearchable=true,可进检索) - failed → markRagflowParseFailed;processing → markRagflowParsePolling 刷新进度并保持 parsing - RAGFlow 临时不可达/轮询异常一律保持 parsing 等下轮重试,绝不误判失败 ProcessingTaskService 补 markRagflowParseCompleted/markRagflowParsePolling;Mapper 补 selectParsingForPoll。 验证: worker 接线 7 + 回归(DocumentService 17 + Facade 3 + ScanService 9 + P-B RetrievalApiImpl 11) 全绿。端到端启用开关后真验。 Co-Authored-By: Claude Opus 4.8 (1M context) --- .../MuseKnowledgeParseStatusPollWorker.java | 133 +++++++++++++++++ .../MuseKnowledgeProcessingTaskService.java | 40 +++++ .../MuseKnowledgeProcessingTaskMapper.java | 13 ++ ...useKnowledgeParseStatusPollWorkerTest.java | 141 ++++++++++++++++++ 4 files changed, 327 insertions(+) create mode 100644 muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeParseStatusPollWorker.java create mode 100644 muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/test/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeParseStatusPollWorkerTest.java diff --git a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeParseStatusPollWorker.java b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeParseStatusPollWorker.java new file mode 100644 index 00000000..4be19148 --- /dev/null +++ b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeParseStatusPollWorker.java @@ -0,0 +1,133 @@ +package cn.iocoder.muse.module.knowledge.application.muse; + +import cn.iocoder.muse.framework.tenant.core.util.TenantUtils; +import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient; +import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeProcessingTaskDO; +import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeProcessingTaskMapper; +import jakarta.annotation.Resource; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; +import org.springframework.util.StringUtils; + +import java.util.List; + +/** + * Knowledge 文档 parse 轮询推进 worker。 + * + *

处理链在 {@code startParseDocuments} 触发 RAGFlow 解析后置 task=parsing/parse_status=processing 即返回, + * RAGFlow 解析+embedding+索引是异步的。本 worker 周期轮询 RAGFlow 文档状态,把完成的任务推进到 + * completed(version processingStatus=searchable,使文档可进检索),失败的标 failed,未完成的保持 parsing。

+ * + *

跨租户运行:先以 ignore-tenant 捞取 parsing 任务,再逐个切到任务自身租户上下文调用 RAGFlow 与写库。 + * RAGFlow 临时不可达/轮询异常一律保持 parsing 等下轮重试,绝不误判失败。默认关闭,需显式开启。

+ */ +@Slf4j +@Component +public class MuseKnowledgeParseStatusPollWorker { + + private static final int BATCH_LIMIT = 20; + private static final String MUSE_STATUS_SEARCHABLE = "searchable"; + private static final String MUSE_STATUS_FAILED = "failed"; + + @Resource + private MuseKnowledgeProcessingTaskMapper processingTaskMapper; + @Resource + private RagFlowKnowledgeRuntimeClient ragFlowClient; + @Resource + private MuseKnowledgeProcessingTaskService processingTaskService; + + @Value("${muse.knowledge.parse-poll-worker.enabled:false}") + private boolean enabled; + + public MuseKnowledgeParseStatusPollWorker() { + } + + // 测试构造:直接注入依赖与开关,避免起 Spring 上下文。 + MuseKnowledgeParseStatusPollWorker(MuseKnowledgeProcessingTaskMapper processingTaskMapper, + RagFlowKnowledgeRuntimeClient ragFlowClient, + MuseKnowledgeProcessingTaskService processingTaskService, + boolean enabled) { + this.processingTaskMapper = processingTaskMapper; + this.ragFlowClient = ragFlowClient; + this.processingTaskService = processingTaskService; + this.enabled = enabled; + } + + @Scheduled(initialDelayString = "${muse.knowledge.parse-poll-worker.initial-delay-ms:5000}", + fixedDelayString = "${muse.knowledge.parse-poll-worker.fixed-delay-ms:5000}") + public void pollScheduled() { + pollOnce(); + } + + /** + * 轮询一批 parsing 任务,返回本轮推进到终态(completed/failed)的任务数。 + */ + public int pollOnce() { + if (!enabled) { + return 0; + } + List tasks = TenantUtils.executeIgnore( + () -> processingTaskMapper.selectParsingForPoll(BATCH_LIMIT)); + if (tasks == null || tasks.isEmpty()) { + return 0; + } + int advanced = 0; + for (MuseKnowledgeProcessingTaskDO task : tasks) { + // 缺租户/外部 id 的任务无法安全轮询,跳过(避免空指针与跨租户写串)。 + if (task.getTenantId() == null || !StringUtils.hasText(task.getRagflowDatasetId()) + || !StringUtils.hasText(task.getRagflowDocumentId())) { + continue; + } + advanced += TenantUtils.execute(task.getTenantId(), () -> pollOne(task)); + } + return advanced; + } + + private int pollOne(MuseKnowledgeProcessingTaskDO task) { + RagFlowKnowledgeRuntimeClient.RuntimeResult result; + try { + result = ragFlowClient.pollDocumentStatuses(new RagFlowKnowledgeRuntimeClient.PollDocumentStatusesCommand( + task.getTenantId(), task.getOwnerUserId(), task.getKbId(), task.getRagflowDatasetId(), + List.of(task.getRagflowDocumentId()), task.getTaskId(), + pollCorrelationId(task), pollRequestHash(task), 1)); + } catch (RuntimeException ex) { + // RAGFlow 临时不可达:保持 parsing,下轮重试,不误判失败。 + log.warn("Knowledge parse 轮询异常,taskId={}, errorType={}", task.getTaskId(), ex.getClass().getSimpleName()); + return 0; + } + if (result.status() == RagFlowKnowledgeRuntimeClient.Status.FAILED) { + log.warn("Knowledge parse 轮询 RAGFlow 返回失败,taskId={}, failureClass={}", + task.getTaskId(), result.failureClass()); + return 0; // 轮询本身失败(非文档失败)→ 保持 parsing 重试 + } + if (result.documentStatuses() == null || result.documentStatuses().isEmpty()) { + return 0; + } + RagFlowKnowledgeRuntimeClient.DocumentStatus status = result.documentStatuses().getFirst(); + String museStatus = status.museStatus(); + if (MUSE_STATUS_SEARCHABLE.equals(museStatus)) { + processingTaskService.markRagflowParseCompleted(task.getDocumentVersionId(), task.getTaskId(), + status.ragflowDocumentId(), status.progress(), result); + return 1; + } + if (MUSE_STATUS_FAILED.equals(museStatus)) { + processingTaskService.markRagflowParseFailed(task.getDocumentVersionId(), task.getTaskId(), + status.ragflowDocumentId(), result); + return 1; + } + // 仍在 RAGFlow 解析中:刷新进度痕迹,保持 parsing。 + processingTaskService.markRagflowParsePolling(task.getTaskId(), status.ragflowRun(), status.progress()); + return 0; + } + + private String pollCorrelationId(MuseKnowledgeProcessingTaskDO task) { + return StringUtils.hasText(task.getCorrelationId()) ? task.getCorrelationId() : "poll:" + task.getTaskId(); + } + + private String pollRequestHash(MuseKnowledgeProcessingTaskDO task) { + return "poll:" + task.getTaskId(); + } + +} diff --git a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeProcessingTaskService.java b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeProcessingTaskService.java index 00f48001..0ad25e7a 100644 --- a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeProcessingTaskService.java +++ b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeProcessingTaskService.java @@ -214,6 +214,46 @@ public class MuseKnowledgeProcessingTaskService { updateVersionStatus(documentVersionId, "passed", "failed", "failed", failureName(result)); } + /** + * parse 轮询 worker 确认 RAGFlow 解析+索引完成后推进终态: + * task → completed,version processingStatus → searchable(使文档 isSearchable=true、可进检索)。 + */ + @Transactional(propagation = Propagation.REQUIRES_NEW, rollbackFor = Exception.class) + public void markRagflowParseCompleted(Long documentVersionId, String taskId, String ragflowDocumentId, + Integer progress, RagFlowKnowledgeRuntimeClient.RuntimeResult result) { + MuseKnowledgeProcessingTaskDO task = processingTaskMapper.selectByTaskId(taskId); + if (task != null) { + task.setStatus("completed"); + task.setParseStatus("completed"); + task.setIndexStatus("completed"); + task.setRagflowDocumentId(ragflowDocumentId); + task.setRagflowRunStatus("done"); + task.setParseProgress(progress == null ? 100 : progress); + task.setRetryable(false); + task.setProgressMessage("RAGFlow parse + index completed"); + task.setFinishedAt(LocalDateTime.now()); + task.setResultSummary(JsonUtils.toJsonString(Map.of("ragflow", safeSummary(result)))); + processingTaskMapper.updateById(task); + } + // processingStatus=searchable 是 DocumentService.setSearchable 的唯一真值来源,必须落到版本上。 + updateVersionStatus(documentVersionId, "passed", "completed", "searchable", null); + } + + /** + * parse 轮询过程态:仅刷新 RAGFlow run/progress 痕迹,保持 parsing 不变(下轮继续轮询)。 + */ + @Transactional(propagation = Propagation.REQUIRES_NEW, rollbackFor = Exception.class) + public void markRagflowParsePolling(String taskId, String ragflowRun, Integer progress) { + MuseKnowledgeProcessingTaskDO task = processingTaskMapper.selectByTaskId(taskId); + if (task != null && "parsing".equals(task.getStatus())) { + task.setRagflowRunStatus(ragflowRun); + if (progress != null) { + task.setParseProgress(progress); + } + processingTaskMapper.updateById(task); + } + } + public String taskId(Long documentVersionId) { return "knowledge-processing-" + documentVersionId; } diff --git a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/dal/mysql/muse/MuseKnowledgeProcessingTaskMapper.java b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/dal/mysql/muse/MuseKnowledgeProcessingTaskMapper.java index 187a8ead..a1b8669b 100644 --- a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/dal/mysql/muse/MuseKnowledgeProcessingTaskMapper.java +++ b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/main/java/cn/iocoder/muse/module/knowledge/dal/mysql/muse/MuseKnowledgeProcessingTaskMapper.java @@ -42,4 +42,17 @@ public interface MuseKnowledgeProcessingTaskMapper extends BaseMapperX selectParsingForPoll(int limit) { + return selectList(new LambdaQueryWrapperX() + .eq(MuseKnowledgeProcessingTaskDO::getStatus, "parsing") + .isNotNull(MuseKnowledgeProcessingTaskDO::getRagflowDatasetId) + .isNotNull(MuseKnowledgeProcessingTaskDO::getRagflowDocumentId) + .orderByAsc(MuseKnowledgeProcessingTaskDO::getUpdateTime) + .last("LIMIT " + limit)); + } + } diff --git a/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/test/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeParseStatusPollWorkerTest.java b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/test/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeParseStatusPollWorkerTest.java new file mode 100644 index 00000000..da185ea0 --- /dev/null +++ b/muse-cloud/muse-module-knowledge/muse-module-knowledge-server/src/test/java/cn/iocoder/muse/module/knowledge/application/muse/MuseKnowledgeParseStatusPollWorkerTest.java @@ -0,0 +1,141 @@ +package cn.iocoder.muse.module.knowledge.application.muse; + +import cn.iocoder.muse.framework.test.core.ut.BaseMockitoUnitTest; +import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient; +import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.DocumentStatus; +import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.Operation; +import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.RuntimeResult; +import cn.iocoder.muse.module.knowledge.application.muse.facade.RagFlowKnowledgeRuntimeClient.Status; +import cn.iocoder.muse.module.knowledge.dal.dataobject.muse.MuseKnowledgeProcessingTaskDO; +import cn.iocoder.muse.module.knowledge.dal.mysql.muse.MuseKnowledgeProcessingTaskMapper; +import org.junit.jupiter.api.Test; +import org.mockito.Mock; + +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +/** + * parse 轮询 worker 接线测试:验各 museStatus 分支推进 + 失败/异常保持 parsing 不误判。 + */ +class MuseKnowledgeParseStatusPollWorkerTest extends BaseMockitoUnitTest { + + @Mock + private MuseKnowledgeProcessingTaskMapper processingTaskMapper; + @Mock + private RagFlowKnowledgeRuntimeClient ragFlowClient; + @Mock + private MuseKnowledgeProcessingTaskService processingTaskService; + + private MuseKnowledgeParseStatusPollWorker worker(boolean enabled) { + return new MuseKnowledgeParseStatusPollWorker(processingTaskMapper, ragFlowClient, processingTaskService, enabled); + } + + private MuseKnowledgeProcessingTaskDO parsingTask() { + MuseKnowledgeProcessingTaskDO t = new MuseKnowledgeProcessingTaskDO(); + t.setTaskId("knowledge-processing-4"); + t.setTenantId(1L); + t.setKbId(2L); + t.setOwnerUserId(1L); + t.setDocumentVersionId(4L); + t.setRagflowDatasetId("ds-1"); + t.setRagflowDocumentId("doc-1"); + t.setCorrelationId("corr-1"); + t.setStatus("parsing"); + return t; + } + + private RuntimeResult poll(Status status, DocumentStatus doc) { + return new RuntimeResult(null, Operation.POLL_DOCUMENT_STATUSES, "corr", "hash", 1, 5L, + status, null, "knowledge-processing-4", "{}", + doc == null ? List.of() : List.of(doc), Map.of()); + } + + @Test + void should_returnZeroAndNoCalls_whenDisabled() { + int advanced = worker(false).pollOnce(); + assertEquals(0, advanced); + verifyNoInteractions(processingTaskMapper, ragFlowClient, processingTaskService); + } + + @Test + void should_markCompleted_whenRagflowDone() { + when(processingTaskMapper.selectParsingForPoll(anyInt())).thenReturn(List.of(parsingTask())); + when(ragFlowClient.pollDocumentStatuses(any())).thenReturn( + poll(Status.SUCCEEDED, new DocumentStatus("doc-1", "searchable", "DONE", 100, "done"))); + + int advanced = worker(true).pollOnce(); + + assertEquals(1, advanced); + verify(processingTaskService).markRagflowParseCompleted(eq(4L), eq("knowledge-processing-4"), + eq("doc-1"), eq(100), any()); + } + + @Test + void should_markFailed_whenRagflowFailed() { + when(processingTaskMapper.selectParsingForPoll(anyInt())).thenReturn(List.of(parsingTask())); + when(ragFlowClient.pollDocumentStatuses(any())).thenReturn( + poll(Status.SUCCEEDED, new DocumentStatus("doc-1", "failed", "FAILED", null, "err"))); + + int advanced = worker(true).pollOnce(); + + assertEquals(1, advanced); + verify(processingTaskService).markRagflowParseFailed(eq(4L), eq("knowledge-processing-4"), eq("doc-1"), any()); + } + + @Test + void should_keepParsing_whenStillProcessing() { + when(processingTaskMapper.selectParsingForPoll(anyInt())).thenReturn(List.of(parsingTask())); + when(ragFlowClient.pollDocumentStatuses(any())).thenReturn( + poll(Status.SUCCEEDED, new DocumentStatus("doc-1", "processing", "RUNNING", 50, "running"))); + + int advanced = worker(true).pollOnce(); + + assertEquals(0, advanced); + verify(processingTaskService).markRagflowParsePolling("knowledge-processing-4", "RUNNING", 50); + verify(processingTaskService, never()).markRagflowParseCompleted(any(), any(), any(), any(), any()); + } + + @Test + void should_keepParsing_whenPollReturnsFailedStatus() { + when(processingTaskMapper.selectParsingForPoll(anyInt())).thenReturn(List.of(parsingTask())); + when(ragFlowClient.pollDocumentStatuses(any())).thenReturn(poll(Status.FAILED, null)); + + int advanced = worker(true).pollOnce(); + + assertEquals(0, advanced); + verify(processingTaskService, never()).markRagflowParseCompleted(any(), any(), any(), any(), any()); + verify(processingTaskService, never()).markRagflowParseFailed(any(), any(), any(), any()); + } + + @Test + void should_keepParsing_whenClientThrows() { + when(processingTaskMapper.selectParsingForPoll(anyInt())).thenReturn(List.of(parsingTask())); + when(ragFlowClient.pollDocumentStatuses(any())).thenThrow(new RuntimeException("boom")); + + int advanced = worker(true).pollOnce(); + + assertEquals(0, advanced); + verify(processingTaskService, never()).markRagflowParseCompleted(any(), any(), any(), any(), any()); + } + + @Test + void should_skip_whenMissingRagflowIds() { + MuseKnowledgeProcessingTaskDO t = parsingTask(); + t.setRagflowDocumentId(null); + when(processingTaskMapper.selectParsingForPoll(anyInt())).thenReturn(List.of(t)); + + int advanced = worker(true).pollOnce(); + + assertEquals(0, advanced); + verifyNoInteractions(ragFlowClient); + } +}