feat(knowledge): 补 parse 轮询推进 worker,打通摄入 parsing→searchable 终态

处理链 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) <noreply@anthropic.com>
This commit is contained in:
lili 2026-06-23 21:14:07 -07:00
parent ec4a02ded2
commit 635045cc30
4 changed files with 327 additions and 0 deletions

View File

@ -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
*
* <p>处理链在 {@code startParseDocuments} 触发 RAGFlow 解析后置 task=parsing/parse_status=processing 即返回
* RAGFlow 解析+embedding+索引是异步的 worker 周期轮询 RAGFlow 文档状态把完成的任务推进到
* completedversion processingStatus=searchable使文档可进检索失败的标 failed未完成的保持 parsing</p>
*
* <p>跨租户运行先以 ignore-tenant 捞取 parsing 任务再逐个切到任务自身租户上下文调用 RAGFlow 与写库
* RAGFlow 临时不可达/轮询异常一律保持 parsing 等下轮重试绝不误判失败默认关闭需显式开启</p>
*/
@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<MuseKnowledgeProcessingTaskDO> 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();
}
}

View File

@ -214,6 +214,46 @@ public class MuseKnowledgeProcessingTaskService {
updateVersionStatus(documentVersionId, "passed", "failed", "failed", failureName(result));
}
/**
* parse 轮询 worker 确认 RAGFlow 解析+索引完成后推进终态
* task completedversion 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;
}

View File

@ -42,4 +42,17 @@ public interface MuseKnowledgeProcessingTaskMapper extends BaseMapperX<MuseKnowl
.in(MuseKnowledgeProcessingTaskDO::getStatus, statuses));
}
/**
* 捞取处于 parsing 且已落 RAGFlow dataset/document id 的任务 parse 轮询 worker 推进
* 跨租户查询调用方需用 TenantUtils.executeIgnore 包裹worker 不绑定单一租户
*/
default List<MuseKnowledgeProcessingTaskDO> selectParsingForPoll(int limit) {
return selectList(new LambdaQueryWrapperX<MuseKnowledgeProcessingTaskDO>()
.eq(MuseKnowledgeProcessingTaskDO::getStatus, "parsing")
.isNotNull(MuseKnowledgeProcessingTaskDO::getRagflowDatasetId)
.isNotNull(MuseKnowledgeProcessingTaskDO::getRagflowDocumentId)
.orderByAsc(MuseKnowledgeProcessingTaskDO::getUpdateTime)
.last("LIMIT " + limit));
}
}

View File

@ -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);
}
}