From 6ce4503b229495757e975a3531f728a490d871e2 Mon Sep 17 00:00:00 2001 From: rainv123 <2148537152@qq.com> Date: Mon, 20 Apr 2026 14:49:01 +0800 Subject: [PATCH] =?UTF-8?q?fix:=E5=AE=9E=E7=8E=B0=E7=9F=A5=E8=AF=86?= =?UTF-8?q?=E5=BA=93=E6=96=87=E6=A1=A3=E4=B8=8ERAGFlow=E5=8F=8C=E5=90=91?= =?UTF-8?q?=E5=90=8C=E6=AD=A5=EF=BC=8C=E6=94=AF=E6=8C=81=E8=BF=9C=E7=AB=AF?= =?UTF-8?q?=E4=B8=8A=E4=BC=A0/=E5=88=A0=E9=99=A4=E8=87=AA=E5=8A=A8?= =?UTF-8?q?=E5=AF=B9=E8=B4=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../java/xiaozhi/common/redis/RedisKeys.java | 7 + .../service/KnowledgeFilesService.java | 9 + .../impl/KnowledgeFilesServiceImpl.java | 207 +++++++++++++++++- 3 files changed, 221 insertions(+), 2 deletions(-) diff --git a/main/manager-api/src/main/java/xiaozhi/common/redis/RedisKeys.java b/main/manager-api/src/main/java/xiaozhi/common/redis/RedisKeys.java index 6b143930..ec2efc6b 100644 --- a/main/manager-api/src/main/java/xiaozhi/common/redis/RedisKeys.java +++ b/main/manager-api/src/main/java/xiaozhi/common/redis/RedisKeys.java @@ -187,4 +187,11 @@ public class RedisKeys { public static String getOtaUploadCountKey(Long username) { return "ota:upload:count:" + username; } + + /** + * 文档全量同步冷却标记key + */ + public static String getDocumentSyncKey(String datasetId) { + return "knowledge:doc:sync:" + datasetId; + } } diff --git a/main/manager-api/src/main/java/xiaozhi/modules/knowledge/service/KnowledgeFilesService.java b/main/manager-api/src/main/java/xiaozhi/modules/knowledge/service/KnowledgeFilesService.java index fa420e38..f46b7417 100644 --- a/main/manager-api/src/main/java/xiaozhi/modules/knowledge/service/KnowledgeFilesService.java +++ b/main/manager-api/src/main/java/xiaozhi/modules/knowledge/service/KnowledgeFilesService.java @@ -120,4 +120,13 @@ public interface KnowledgeFilesService { * 同步所有处于 RUNNING 状态的文档 (供定时任务调用) */ void syncRunningDocuments(); + + /** + * 从RAGFlow全量同步文档到本地影子表 + * 拉取远端所有文档,与本地影子表对比,插入缺失的记录 + * + * @param datasetId 数据集ID + * @return 新同步的文档数量 + */ + int syncDocumentsFromRAG(String datasetId); } \ No newline at end of file diff --git a/main/manager-api/src/main/java/xiaozhi/modules/knowledge/service/impl/KnowledgeFilesServiceImpl.java b/main/manager-api/src/main/java/xiaozhi/modules/knowledge/service/impl/KnowledgeFilesServiceImpl.java index 3bf09b60..49a8820f 100644 --- a/main/manager-api/src/main/java/xiaozhi/modules/knowledge/service/impl/KnowledgeFilesServiceImpl.java +++ b/main/manager-api/src/main/java/xiaozhi/modules/knowledge/service/impl/KnowledgeFilesServiceImpl.java @@ -4,6 +4,8 @@ import java.util.ArrayList; import java.util.Date; import java.util.List; import java.util.Map; +import java.util.Set; +import java.util.stream.Collectors; import org.apache.commons.lang3.StringUtils; import org.springframework.beans.BeanUtils; @@ -76,6 +78,27 @@ public class KnowledgeFilesServiceImpl extends BaseServiceImpl().eq("dataset_id", datasetId)); + localEmpty = localCount == null || localCount == 0; + } + if (!cooldownActive || localEmpty) { + try { + self.syncDocumentsFromRAG(datasetId); + // 60秒冷却,避免每次翻页都触发远端调用(本地为空时缩短为10秒) + int cooldownSec = localEmpty ? 10 : 60; + redisUtils.set(syncKey, "1", cooldownSec); + } catch (Exception e) { + log.warn("从RAGFlow全量同步文档失败(不影响本地查询): datasetId={}, error={}", datasetId, e.getMessage()); + } + } + // 1. 获取本地影子表数据 (MyBatis-Plus 分页) Page pageParams = new Page<>(page, limit); QueryWrapper queryWrapper = new QueryWrapper<>(); @@ -512,7 +535,7 @@ public class KnowledgeFilesServiceImpl extends BaseServiceImpl ragConfig = knowledgeBaseService.getRAGConfigByDatasetId(datasetId); + KnowledgeBaseAdapter adapter = KnowledgeBaseAdapterFactory.getAdapter(extractAdapterType(ragConfig), ragConfig); + + // 2. 分页拉取远端所有文档 + List allRemoteDocs = new ArrayList<>(); + int pageNum = 1; + int pageSize = 100; + long totalRemote = Long.MAX_VALUE; + + while ((long) (pageNum - 1) * pageSize < totalRemote) { + DocumentDTO.ListReq req = DocumentDTO.ListReq.builder() + .page(pageNum) + .pageSize(pageSize) + .build(); + PageData remotePage = adapter.getDocumentList(datasetId, req); + if (remotePage == null || remotePage.getList() == null || remotePage.getList().isEmpty()) { + break; + } + allRemoteDocs.addAll(remotePage.getList()); + totalRemote = remotePage.getTotal(); + pageNum++; + } + + // 3. 获取本地已有文档 + List localDocs = documentDao.selectList( + new QueryWrapper().eq("dataset_id", datasetId)); + Set localDocIds = localDocs.stream() + .map(DocumentEntity::getDocumentId) + .collect(Collectors.toSet()); + + // 4. 远端文档ID集合 + Set remoteDocIds = allRemoteDocs.stream() + .map(KnowledgeFilesDTO::getDocumentId) + .filter(id -> id != null) + .collect(Collectors.toSet()); + + // 5. 补充: 插入远端存在但本地缺失的文档 + List newDocs = allRemoteDocs.stream() + .filter(doc -> doc.getDocumentId() != null && !localDocIds.contains(doc.getDocumentId())) + .collect(Collectors.toList()); + + int syncCount = 0; + if (!newDocs.isEmpty()) { + for (KnowledgeFilesDTO doc : newDocs) { + try { + self.saveDocumentShadow(datasetId, doc, doc.getName(), doc.getChunkMethod(), doc.getParserConfig()); + // 同步远端已有的 token/chunk 统计 + Long tokenCount = doc.getTokenCount() != null ? doc.getTokenCount() : 0L; + long chunkCount = doc.getChunkCount() != null ? doc.getChunkCount().longValue() : 0L; + if (tokenCount > 0 || chunkCount > 0) { + knowledgeBaseService.updateStatistics(datasetId, 0, chunkCount, tokenCount); + } + syncCount++; + } catch (Exception e) { + log.warn("同步单个文档影子记录失败: docId={}, error={}", doc.getDocumentId(), e.getMessage()); + } + } + log.info("从RAGFlow新增同步 {} 个文档影子记录, datasetId={}", syncCount, datasetId); + } + + // 6. 清理: 删除远端已不存在但本地仍保留的影子记录 + List deletedDocs = localDocs.stream() + .filter(entity -> !remoteDocIds.contains(entity.getDocumentId())) + .collect(Collectors.toList()); + + if (!deletedDocs.isEmpty()) { + List deletedDocIds = new ArrayList<>(); + long totalChunkDelta = 0; + long totalTokenDelta = 0; + + for (DocumentEntity entity : deletedDocs) { + deletedDocIds.add(entity.getDocumentId()); + totalChunkDelta += entity.getChunkCount() != null ? entity.getChunkCount() : 0L; + totalTokenDelta += entity.getTokenCount() != null ? entity.getTokenCount() : 0L; + } + try { + self.deleteDocumentShadows(deletedDocIds, datasetId, totalChunkDelta, totalTokenDelta); + log.info("清理远端已删除的影子记录: {} 个, datasetId={}", deletedDocs.size(), datasetId); + } catch (Exception e) { + log.warn("清理远端已删除的影子记录失败: datasetId={}, error={}", datasetId, e.getMessage()); + } + } + + // 7. 更新: 远端和本地都存在但状态不一致的文档(处理RAGFlow复用documentId重传场景) + // 当文档在RAGFlow被删除后重新上传,RAGFlow可能复用同一个documentId, + // 导致本地影子记录仍保留旧的CANCEL/FAIL状态而不会被步骤5/6处理 + Map remoteDocMap = allRemoteDocs.stream() + .filter(doc -> doc.getDocumentId() != null) + .collect(Collectors.toMap(KnowledgeFilesDTO::getDocumentId, doc -> doc, (a, b) -> b)); + + Map localDocMap = localDocs.stream() + .collect(Collectors.toMap(DocumentEntity::getDocumentId, e -> e, (a, b) -> b)); + + int updateCount = 0; + for (Map.Entry entry : remoteDocMap.entrySet()) { + String docId = entry.getKey(); + DocumentEntity local = localDocMap.get(docId); + if (local == null) { + continue; // 不在本地,由步骤5处理 + } + KnowledgeFilesDTO remote = entry.getValue(); + + // 判断是否需要更新:run状态或status不同,说明远端有变化 + boolean runChanged = remote.getRun() != null && !remote.getRun().equals(local.getRun()); + boolean statusChanged = remote.getStatus() != null && !remote.getStatus().equals(local.getStatus()); + boolean nameChanged = remote.getName() != null && !remote.getName().equals(local.getName()); + + if (runChanged || statusChanged || nameChanged) { + log.info("影子更新:检测到远端文档状态变化, docId={}, 本地run={}, 远端run={}, 本地status={}, 远端status={}", + docId, local.getRun(), remote.getRun(), local.getStatus(), remote.getStatus()); + + UpdateWrapper updateWrapper = new UpdateWrapper() + .set("run", remote.getRun()) + .set("status", remote.getStatus() != null ? remote.getStatus() : local.getStatus()) + .set("progress", remote.getProgress()) + .set("chunk_count", remote.getChunkCount()) + .set("token_count", remote.getTokenCount()) + .set("error", remote.getError()) + .set("process_duration", remote.getProcessDuration()) + .set("updated_at", new Date()) + .set("last_sync_at", new Date()) + .eq("document_id", docId) + .eq("dataset_id", datasetId); + + if (remote.getName() != null) { + updateWrapper.set("name", remote.getName()); + } + if (remote.getThumbnail() != null) { + updateWrapper.set("thumbnail", remote.getThumbnail()); + } + if (remote.getMetaFields() != null) { + try { + updateWrapper.set("meta_fields", objectMapper.writeValueAsString(remote.getMetaFields())); + } catch (Exception e) { + log.warn("同步更新元数据序列化失败: docId={}, error={}", docId, e.getMessage()); + } + } + + documentDao.update(null, updateWrapper); + + // 如果状态从终态(CANCEL/FAIL)变为活跃态(UNSTART/RUNNING),需要同步token统计 + boolean wasTerminal = "CANCEL".equals(local.getRun()) || "FAIL".equals(local.getRun()); + boolean isActive = "RUNNING".equals(remote.getRun()) || "UNSTART".equals(remote.getRun()); + Long remoteTokenCount = remote.getTokenCount() != null ? remote.getTokenCount() : 0L; + Long localTokenCount = local.getTokenCount() != null ? local.getTokenCount() : 0L; + long remoteChunkCount = remote.getChunkCount() != null ? remote.getChunkCount().longValue() : 0L; + long localChunkCount = local.getChunkCount() != null ? local.getChunkCount().longValue() : 0L; + if (wasTerminal && isActive) { + long tokenDelta = remoteTokenCount - localTokenCount; + long chunkDelta = remoteChunkCount - localChunkCount; + if (tokenDelta != 0 || chunkDelta != 0) { + knowledgeBaseService.updateStatistics(datasetId, 0, chunkDelta, tokenDelta); + log.info("影子更新: 修正知识库统计, docId={}, chunkDelta={}, tokenDelta={}", docId, chunkDelta, tokenDelta); + } + } + + updateCount++; + } + } + + if (syncCount == 0 && deletedDocs.isEmpty() && updateCount == 0) { + log.info("本地影子表已与RAGFlow完全同步, datasetId={}", datasetId); + } else { + log.info("同步完成: 新增={}, 清理={}, 更新={}, datasetId={}", syncCount, deletedDocs.size(), updateCount, datasetId); + } + + return syncCount; + } + @Override public void syncRunningDocuments() { // 1. 查询所有 RUNNING 状态的文档 @@ -755,7 +958,7 @@ public class KnowledgeFilesServiceImpl extends BaseServiceImpl> groupedDocs = runningDocs.stream() - .collect(java.util.stream.Collectors.groupingBy(DocumentEntity::getDatasetId)); + .collect(Collectors.groupingBy(DocumentEntity::getDatasetId)); groupedDocs.forEach((datasetId, docs) -> { KnowledgeBaseAdapter adapter = null;