update:增加文档解析状态

This commit is contained in:
rainv123
2025-11-06 18:20:01 +08:00
parent eaf4f4b08e
commit 90980652ce
8 changed files with 460 additions and 134 deletions
@@ -57,11 +57,19 @@ public class KnowledgeFilesController {
public Result<PageData<KnowledgeFilesDTO>> getPageList(
@PathVariable("dataset_id") String datasetId,
@RequestParam(required = false) String name,
@RequestParam(required = false) Integer status,
@RequestParam(required = false, defaultValue = "1") Integer page,
@RequestParam(required = false, defaultValue = "10") Integer page_size) {
// 验证知识库权限
validateKnowledgeBasePermission(datasetId);
// 如果指定了状态参数,使用状态查询接口
if (status != null) {
PageData<KnowledgeFilesDTO> pageData = knowledgeFilesService.getPageListByStatus(datasetId, status, page, page_size);
return new Result<PageData<KnowledgeFilesDTO>>().ok(pageData);
}
// 否则使用通用查询接口
KnowledgeFilesDTO knowledgeFilesDTO = new KnowledgeFilesDTO();
knowledgeFilesDTO.setDatasetId(datasetId);
knowledgeFilesDTO.setName(name);
@@ -69,6 +77,21 @@ public class KnowledgeFilesController {
return new Result<PageData<KnowledgeFilesDTO>>().ok(pageData);
}
@GetMapping("/documents/status/{status}")
@Operation(summary = "根据状态分页查询文档列表")
@RequiresPermissions("sys:role:normal")
public Result<PageData<KnowledgeFilesDTO>> getPageListByStatus(
@PathVariable("dataset_id") String datasetId,
@PathVariable("status") Integer status,
@RequestParam(required = false, defaultValue = "1") Integer page,
@RequestParam(required = false, defaultValue = "10") Integer page_size) {
// 验证知识库权限
validateKnowledgeBasePermission(datasetId);
PageData<KnowledgeFilesDTO> pageData = knowledgeFilesService.getPageListByStatus(datasetId, status, page, page_size);
return new Result<PageData<KnowledgeFilesDTO>>().ok(pageData);
}
@PostMapping("/documents")
@Operation(summary = "上传文档到知识库")
@RequiresPermissions("sys:role:normal")
@@ -186,7 +209,6 @@ public class KnowledgeFilesController {
return new Result<Map<String, Object>>().error("召回测试失败: " + e.getMessage());
}
}
/**
* 解析JSON字符串为Map对象
*/
@@ -47,6 +47,9 @@ public class KnowledgeFilesDTO implements Serializable {
@Schema(description = "状态")
private Integer status;
@Schema(description = "文档解析状态")
private String run;
@Schema(description = "创建者")
private Long creator;
@@ -58,4 +61,36 @@ public class KnowledgeFilesDTO implements Serializable {
@Schema(description = "更新时间")
private Date updatedAt;
// 文档解析状态常量定义
private static final Integer STATUS_UNSTART = 0;
private static final Integer STATUS_RUNNING = 1;
private static final Integer STATUS_CANCEL = 2;
private static final Integer STATUS_DONE = 3;
private static final Integer STATUS_FAIL = 4;
/**
* 获取文档解析状态码(基于run字段转换)
*/
public Integer getParseStatusCode() {
if (run == null) {
return STATUS_UNSTART;
}
// 根据run字段的值直接映射到对应的状态码
switch (run.toUpperCase()) {
case "RUNNING":
return STATUS_RUNNING;
case "CANCEL":
return STATUS_CANCEL;
case "DONE":
return STATUS_DONE;
case "FAIL":
return STATUS_FAIL;
case "UNSTART":
default:
return STATUS_UNSTART;
}
}
}
@@ -47,6 +47,17 @@ public interface KnowledgeFilesService {
Map<String, Object> metaFields, String chunkMethod,
Map<String, Object> parserConfig);
/**
* 根据状态分页查询文档列表
*
* @param datasetId 知识库ID
* @param status 文档解析状态(0-未开始,1-进行中,2-已取消,3-已完成,4-失败)
* @param page 页码
* @param limit 每页数量
* @return 分页数据
*/
PageData<KnowledgeFilesDTO> getPageListByStatus(String datasetId, Integer status, Integer page, Integer limit);
/**
* 根据文档ID和知识库ID删除文档
*
@@ -9,6 +9,7 @@ import java.util.Date;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
import org.apache.commons.lang3.StringUtils;
import org.springframework.core.io.AbstractResource;
@@ -36,6 +37,7 @@ import xiaozhi.common.constant.Constant;
import xiaozhi.common.exception.ErrorCode;
import xiaozhi.common.exception.RenException;
import xiaozhi.common.page.PageData;
import xiaozhi.common.redis.RedisUtils;
import xiaozhi.modules.knowledge.dto.KnowledgeFilesDTO;
import xiaozhi.modules.knowledge.service.KnowledgeFilesService;
import xiaozhi.modules.model.dao.ModelConfigDao;
@@ -49,6 +51,7 @@ public class KnowledgeFilesServiceImpl implements KnowledgeFilesService {
private final ModelConfigService modelConfigService;
private final ModelConfigDao modelConfigDao;
private final RedisUtils redisUtils;
private RestTemplate restTemplate = new RestTemplate();
private ObjectMapper objectMapper = new ObjectMapper();
@@ -197,6 +200,8 @@ public class KnowledgeFilesServiceImpl implements KnowledgeFilesService {
for (Map<String, Object> docMap : documentList) {
KnowledgeFilesDTO dto = convertRAGDocumentToDTO(docMap);
if (dto != null) {
// 在文档列表获取时也进行状态同步检查
syncDocumentStatusWithRAGFlow(dto);
documents.add(dto);
}
}
@@ -239,6 +244,88 @@ public class KnowledgeFilesServiceImpl implements KnowledgeFilesService {
}
}
/**
* 同步文档状态与RAGFlow实际状态
* 优化状态同步逻辑,确保解析中状态能够正常显示
* 只有当文档有切片且解析时间超过30秒时,才更新为完成状态
*/
private void syncDocumentStatusWithRAGFlow(KnowledgeFilesDTO dto) {
if (dto == null || StringUtils.isBlank(dto.getDocumentId())) {
return;
}
String documentId = dto.getDocumentId();
Integer currentStatus = dto.getStatus();
// 只有当状态明确为处理中(1)时,才进行状态同步检查
// 避免在状态不确定或已完成的文档上重复检查
if (currentStatus != null && currentStatus == 1) {
try {
long currentTime = System.currentTimeMillis();
// 调用RAGFlow API获取文档切片信息
Map<String, Object> ragConfig = getDefaultRAGConfig();
String baseUrl = (String) ragConfig.get("base_url");
String apiKey = (String) ragConfig.get("api_key");
String url = baseUrl + "/api/v1/documents/" + documentId + "/chunks";
HttpHeaders headers = new HttpHeaders();
headers.setContentType(MediaType.APPLICATION_JSON);
headers.set("Authorization", "Bearer " + apiKey);
HttpEntity<String> requestEntity = new HttpEntity<>(headers);
log.debug("检查文档切片状态,documentId: {}", documentId);
ResponseEntity<String> response = restTemplate.exchange(url, HttpMethod.GET, requestEntity,
String.class);
if (response.getStatusCode().is2xxSuccessful()) {
String responseBody = response.getBody();
Map<String, Object> responseMap = objectMapper.readValue(responseBody, Map.class);
Integer code = (Integer) responseMap.get("code");
if (code != null && code == 0) {
Object dataObj = responseMap.get("data");
if (dataObj instanceof Map) {
Map<String, Object> dataMap = (Map<String, Object>) dataObj;
// 检查是否有切片数据
Object chunksObj = getValueFromMultipleKeys(dataMap, "chunks", "items", "list", "data");
if (chunksObj instanceof List) {
List<?> chunks = (List<?>) chunksObj;
// 如果有切片且数量大于0,说明解析已完成
if (!chunks.isEmpty()) {
// 检查文档创建时间,确保解析过程有足够的时间显示
Date createdAt = dto.getCreatedAt();
long parseDuration = currentTime
- (createdAt != null ? createdAt.getTime() : currentTime);
// 只有当解析时间超过30秒时,才更新为完成状态
// 这样可以确保解析中状态有足够的时间显示
if (parseDuration > 30000) {
log.info("状态同步:文档已有切片且解析时间超过30秒,更新为完成状态,documentId: {}, 切片数量: {}, 解析时长: {}ms",
documentId, chunks.size(), parseDuration);
// 更新状态为完成(3)
dto.setStatus(3);
} else {
log.debug("文档已有切片但解析时间不足30秒,保持解析中状态,documentId: {}, 解析时长: {}ms",
documentId, parseDuration);
}
}
}
}
}
}
} catch (Exception e) {
log.debug("检查文档切片状态失败,documentId: {}, 错误: {}", documentId, e.getMessage());
// 忽略检查失败,保持原状态
}
}
}
/**
* 将RAGFlow文档数据转换为KnowledgeFilesDTO
*/
@@ -313,6 +400,21 @@ public class KnowledgeFilesServiceImpl implements KnowledgeFilesService {
}
}
// 设置文档解析状态信息 - 直接使用RAGFlow最新状态
String documentId = dto.getDocumentId();
if (StringUtils.isNotBlank(documentId)) {
// 获取RAGFlow的最新状态
Object runObj = getValueFromMultipleKeys(docMap, "run", "status", "parse_status");
Integer ragFlowStatus = null;
if (runObj != null) {
dto.setRun(runObj.toString());
ragFlowStatus = dto.getParseStatusCode();
log.debug("获取RAGFlow最新状态,documentId: {}, run: {}, status: {}",
documentId, runObj, ragFlowStatus);
}
}
return dto;
} catch (Exception e) {
@@ -433,6 +535,115 @@ public class KnowledgeFilesServiceImpl implements KnowledgeFilesService {
}
}
@Override
public PageData<KnowledgeFilesDTO> getPageListByStatus(String datasetId, Integer status, Integer page,
Integer limit) {
if (StringUtils.isBlank(datasetId)) {
throw new RenException(ErrorCode.PARAMS_GET_ERROR, "datasetId不能为空");
}
log.info("=== 开始根据状态查询文档列表 ===");
log.info("datasetId: {}, status: {}, page: {}, limit: {}", datasetId, status, page, limit);
try {
// 获取RAG配置
Map<String, Object> ragConfig = getDefaultRAGConfig();
String baseUrl = (String) ragConfig.get("base_url");
String apiKey = (String) ragConfig.get("api_key");
// 构建请求URL - 获取文档列表
StringBuilder urlBuilder = new StringBuilder();
urlBuilder.append(baseUrl).append("/api/v1/datasets/").append(datasetId).append("/documents");
// 构建查询参数
List<String> params = new ArrayList<>();
if (page != null && page > 0) {
params.add("page=" + page);
}
if (limit != null && limit > 0) {
params.add("page_size=" + limit);
}
if (!params.isEmpty()) {
urlBuilder.append("?").append(String.join("&", params));
}
String url = urlBuilder.toString();
log.debug("请求URL: {}", url);
// 设置请求头
HttpHeaders headers = new HttpHeaders();
headers.setContentType(MediaType.APPLICATION_JSON);
headers.set("Authorization", "Bearer " + apiKey);
HttpEntity<String> requestEntity = new HttpEntity<>(headers);
// 发送GET请求
log.info("发送GET请求到RAGFlow API获取文档列表...");
ResponseEntity<String> response = restTemplate.exchange(url, HttpMethod.GET, requestEntity, String.class);
log.info("RAGFlow API响应状态码: {}", response.getStatusCode());
if (!response.getStatusCode().is2xxSuccessful()) {
log.error("RAGFlow API调用失败,状态码: {}, 响应内容: {}", response.getStatusCode(), response.getBody());
throw new RenException(ErrorCode.RAG_API_ERROR);
}
String responseBody = response.getBody();
log.debug("RAGFlow API响应内容: {}", responseBody);
// 解析响应
Map<String, Object> responseMap = objectMapper.readValue(responseBody, Map.class);
Integer code = (Integer) responseMap.get("code");
if (code != null && code == 0) {
Object dataObj = responseMap.get("data");
// 解析文档列表并过滤状态
PageData<KnowledgeFilesDTO> pageData = parseDocumentListResponse(dataObj, page, limit);
if (status != null) {
// 根据状态过滤文档列表
List<KnowledgeFilesDTO> filteredDocuments = pageData.getList().stream()
.filter(doc -> status.equals(doc.getStatus()))
.collect(Collectors.toList());
// 更新分页数据
pageData.setList(filteredDocuments);
pageData.setTotal(filteredDocuments.size());
}
log.info("根据状态查询文档列表成功,datasetId: {}, 状态: {}, 文档数量: {}",
datasetId, status, pageData.getList().size());
return pageData;
} else {
log.error("RAGFlow API调用失败,响应码: {}, 响应内容: {}", code, responseBody);
throw new RenException(ErrorCode.RAG_API_ERROR, "RAGFlow API调用失败,响应码: " + code);
}
} catch (IOException e) {
log.error("解析RAGFlow API响应失败: {}", e.getMessage(), e);
throw new RenException(ErrorCode.RAG_API_ERROR,
"解析RAGFlow响应失败: " + e.getMessage());
} catch (HttpClientErrorException e) {
log.error("RAGFlow API调用失败 - HTTP错误: {}, 状态码: {}, 响应内容: {}",
e.getMessage(), e.getStatusCode(), e.getResponseBodyAsString());
throw new RenException(ErrorCode.RAG_API_ERROR, "获取RAGFlow文档列表失败: " + e.getMessage());
} catch (HttpServerErrorException e) {
log.error("RAGFlow API调用失败 - 服务器错误: {}, 状态码: {}, 响应内容: {}",
e.getMessage(), e.getStatusCode(), e.getResponseBodyAsString());
throw new RenException(ErrorCode.RAG_API_ERROR, "获取RAGFlow文档列表失败: " + e.getMessage());
} catch (ResourceAccessException e) {
log.error("RAGFlow API调用失败 - 网络连接错误: {}", e.getMessage(), e);
throw new RenException(ErrorCode.RAG_API_ERROR, "获取RAGFlow文档列表失败: 网络连接错误 - " + e.getMessage());
} catch (Exception e) {
log.error("根据状态查询文档列表失败: {}", e.getMessage(), e);
throw new RenException(ErrorCode.RAG_API_ERROR, "查询文档列表失败: " + e.getMessage());
} finally {
log.info("=== 根据状态查询文档列表操作结束 ===");
}
}
@Override
public KnowledgeFilesDTO uploadDocument(String datasetId, MultipartFile file, String name,
Map<String, Object> metaFields, String chunkMethod,