From d0cd06f9fcb9521d755917a545cc1cdb0492f2b2 Mon Sep 17 00:00:00 2001 From: wanghanlin <1533525126@qq.com> Date: Thu, 3 Sep 2026 14:41:03 +0800 Subject: [PATCH] =?UTF-8?q?fix(document):=20=E4=BF=AE=E5=A4=8D=E6=89=B9?= =?UTF-8?q?=E9=87=8F=E4=B8=8A=E4=BC=A0=E5=90=91=E9=87=8F=E5=8C=96=E5=8F=AA?= =?UTF-8?q?=E6=88=90=E5=8A=9F=E5=B0=8F=E6=96=87=E6=A1=A3=E9=97=AE=E9=A2=98?= =?UTF-8?q?=E5=B9=B6=E5=8A=A0=E5=9B=BA=E5=A4=84=E7=90=86=E9=93=BE=E8=B7=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 移除逐 chunk 串行 AI 关键词提取(MyKeywordEnricher):产出 excerpt_keywords 全库无检索消费点,却将外部调用放大 chunkCount 倍且与在线对话共用模型抢限流 - 向量化改按批入库(默认 50 块/批, knowledge.vector.batch-size)逐批 try-catch 隔离:失败批只记录缺失区间继续,已入库块保留,error_message 聚合"已入库 x/y 块+缺失区间+原因",chunk_count 改为实际入库块数 - 新增文档级失败自动整体重试(≤2 次, 仅 timeout/429/5xx 等瞬时错误):重试前清残留向量防重复;处理期间删除文档即终止写入避免孤儿向量;外层兜底保证状态不悬挂 PROCESSING - Embedding 客户端(DashScope/OpenAI 兼容/豆包多模态)补显式 connect 10s/read 60s 超时,避免慢响应无限挂起占死 documentExecutor 线程 - 文档同步 CLAUDE.md/README.md,新增 knowledge.vector.batch-size 配置说明 --- CLAUDE.md | 4 +- README.md | 9 +- .../config/EmbeddingModelFactory.java | 30 +- .../VolcengineMultimodalEmbeddingModel.java | 16 + .../document/transform/MyKeywordEnricher.java | 30 -- .../service/DocumentProcessingService.java | 452 ++++++++++++++---- .../supportbot/service/DocumentService.java | 8 +- src/main/resources/application.yml | 3 + 8 files changed, 410 insertions(+), 142 deletions(-) delete mode 100644 src/main/java/com/wok/supportbot/document/transform/MyKeywordEnricher.java diff --git a/CLAUDE.md b/CLAUDE.md index ed869ee..af43d11 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -67,7 +67,9 @@ AI 智能客服系统,基于 Spring AI Alibaba + 通义千问 + PGVector,支 - **Open API**: `OpenApiController` 已接入 `ChatPipeline`,补齐角色/RAG/FAQ/MCP/分类隔离能力 ### 文档处理管道 -`DocumentService.uploadDocument()` 统一流程:文档提取 → `MyTokenTextSplitter` 分块 → `MyKeywordEnricher` AI 关键词提取 → `pgVectorVectorStore.add()` 向量化存储。每个分块的 metadata 中注入 `documentId`、`chunkIndex`、`sourceName`、`title` 以关联 `knowledge_document` 表。 +`DocumentService.uploadDocument()` 统一流程:文档提取 → `MyTokenTextSplitter` 分块 → 为每块写 metadata → 按批向量化(默认 50 块/批,配置项 `knowledge.vector.batch-size`)`pgVectorVectorStore.add(batch)` 入库。每个分块的 metadata 注入 `documentId`、`chunkIndex`、`sourceName`、`title`、`categoryId`、`enabled` 关联 `knowledge_document` 表。 + +**向量化加固**(`DocumentProcessingService`):逐批 try-catch 隔离,失败批只记录缺失区间后继续,已入库块保留;`chunk_count` 记实际入库块数,`error_message` 聚合"已入库 x/y 块 + 缺失区间 + 原因";文档级失败自动整体重试至多 2 次(仅瞬时/限流/超时类错误,4xx 不空转),重试前先清残留向量再重建。**无逐块 AI 关键词提取环节**(`MyKeywordEnricher` 已移除,其产出 `excerpt_keywords` 全库无检索消费点)。 ## 关键配置 diff --git a/README.md b/README.md index a5eb006..dc4e2dd 100644 --- a/README.md +++ b/README.md @@ -188,8 +188,7 @@ src/main/java/com/wok/supportbot/ │ │ ├── JsonDocumentLoader.java # JSON 解析(3种模式) │ │ └── SimpleStringDocumentReader.java # 纯文本读取 │ └── transform/ # 文档转换器 -│ ├── MyTokenTextSplitter.java # Token 分块器 -│ └── MyKeywordEnricher.java # AI 关键词提取 +│ └── MyTokenTextSplitter.java # Token 分块器 ├── entity/ # 数据实体类 │ ├── ChatMessage.java # 聊天消息实体 │ ├── KnowledgeDocument.java # 知识文档实体 @@ -234,11 +233,11 @@ src/main/resources/ ↓ [Token 分块] MyTokenTextSplitter (200 token / 100 overlap) ↓ -[关键词提取] MyKeywordEnricher (AI 提取 Top-5 关键词) +[元数据标注] metadata.documentId / chunkIndex / sourceName / title / categoryId / enabled ↓ -[向量化存储] DashScope text-embedding-v2 → PGVector +[分批向量化] 默认 50 块/批循环 add → PGVector(逐批隔离,失败批记录区间后继续,整体失败自动重试 ≤2 次) ↓ -[元数据关联] metadata.documentId / chunkIndex / sourceName / title +[状态更新] READY / FAILED(已入库块保留,chunk_count = 实际入库块数,error_message 聚合缺失区间) ``` ## 🖥️ 前端管理页面 diff --git a/src/main/java/com/wok/supportbot/config/EmbeddingModelFactory.java b/src/main/java/com/wok/supportbot/config/EmbeddingModelFactory.java index b9216b8..69dccf5 100644 --- a/src/main/java/com/wok/supportbot/config/EmbeddingModelFactory.java +++ b/src/main/java/com/wok/supportbot/config/EmbeddingModelFactory.java @@ -13,10 +13,14 @@ import org.springframework.ai.openai.OpenAiEmbeddingModel; import org.springframework.ai.openai.OpenAiEmbeddingOptions; import org.springframework.ai.openai.api.OpenAiApi; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.http.client.JdkClientHttpRequestFactory; import org.springframework.retry.backoff.ExponentialBackOffPolicy; import org.springframework.retry.support.RetryTemplate; import org.springframework.stereotype.Component; +import org.springframework.web.client.RestClient; +import java.net.http.HttpClient; +import java.time.Duration; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; @@ -110,7 +114,8 @@ public class EmbeddingModelFactory { if ("dashscope".equalsIgnoreCase(config.getProvider())) { log.info("创建 DashScope EmbeddingModel: model={}, baseUrl={}", config.getModelName(), baseUrl); DashScopeApi.Builder apiBuilder = DashScopeApi.builder() - .apiKey(config.getApiKey()); + .apiKey(config.getApiKey()) + .restClientBuilder(timeoutRestClientBuilder()); // 团队版/私有化部署使用自定义 baseUrl if (baseUrl != null && !baseUrl.isBlank()) { apiBuilder.baseUrl(baseUrl); @@ -147,6 +152,7 @@ public class EmbeddingModelFactory { .apiKey(config.getApiKey()) .baseUrl(baseUrl) .embeddingsPath(embeddingsPath) + .restClientBuilder(timeoutRestClientBuilder()) .build(); OpenAiEmbeddingOptions options = OpenAiEmbeddingOptions.builder() @@ -223,6 +229,28 @@ public class EmbeddingModelFactory { log.info("EmbeddingModel 缓存已清除"); } + /** 向量化模型连接超时:connect 10s 与聊天路径(ChatModelFactory)保持一致 */ + private static final Duration EMBED_CONNECT_TIMEOUT = Duration.ofSeconds(10); + + /** + * 向量化模型读取超时:embedding 为非流式单次请求,60s 足以容纳大体积批请求, + * 避免慢响应无限挂起占死 documentExecutor 线程 + */ + private static final Duration EMBED_READ_TIMEOUT = Duration.ofSeconds(60); + + /** + * 构建带显式 connect/read 超时的 RestClient(DashScope / OpenAI 兼容 embedding 客户端使用)。 + * 否则第三方默认可能无限等待,长文档批量向量化时一次慢响应即拖垮整个任务 + */ + private static RestClient.Builder timeoutRestClientBuilder() { + HttpClient httpClient = HttpClient.newBuilder() + .connectTimeout(EMBED_CONNECT_TIMEOUT) + .build(); + JdkClientHttpRequestFactory requestFactory = new JdkClientHttpRequestFactory(httpClient); + requestFactory.setReadTimeout(EMBED_READ_TIMEOUT); + return RestClient.builder().requestFactory(requestFactory); + } + // ==================== F1: 连接测试 ==================== /** diff --git a/src/main/java/com/wok/supportbot/config/VolcengineMultimodalEmbeddingModel.java b/src/main/java/com/wok/supportbot/config/VolcengineMultimodalEmbeddingModel.java index 9294823..d68bc87 100644 --- a/src/main/java/com/wok/supportbot/config/VolcengineMultimodalEmbeddingModel.java +++ b/src/main/java/com/wok/supportbot/config/VolcengineMultimodalEmbeddingModel.java @@ -12,9 +12,12 @@ import org.springframework.ai.embedding.EmbeddingResponseMetadata; import org.springframework.ai.chat.metadata.DefaultUsage; import org.springframework.http.HttpHeaders; import org.springframework.http.MediaType; +import org.springframework.http.client.JdkClientHttpRequestFactory; import org.springframework.retry.support.RetryTemplate; import org.springframework.web.client.RestClient; +import java.net.http.HttpClient; +import java.time.Duration; import java.util.ArrayList; import java.util.HashMap; import java.util.List; @@ -39,6 +42,12 @@ public class VolcengineMultimodalEmbeddingModel implements EmbeddingModel { private static final String EMBEDDINGS_PATH = "/embeddings/multimodal"; + /** 连接超时:与聊天路径(ChatModelFactory)保持一致 */ + private static final Duration CONNECT_TIMEOUT = Duration.ofSeconds(10); + + /** 读取超时:单条文本一次请求,60s 足以容纳慢冷启动/大体积文本,避免无限挂起占用文档处理线程 */ + private static final Duration READ_TIMEOUT = Duration.ofSeconds(60); + private final String apiKey; private final String baseUrl; private final String modelName; @@ -54,8 +63,15 @@ public class VolcengineMultimodalEmbeddingModel implements EmbeddingModel { this.modelName = modelName; this.dimensions = dimensions; this.retryTemplate = retryTemplate; + // 显式设置 connect/read 超时:embedding 为非流式单次请求,避免慢响应无限挂起占用文档处理线程 + HttpClient httpClient = HttpClient.newBuilder() + .connectTimeout(CONNECT_TIMEOUT) + .build(); + JdkClientHttpRequestFactory requestFactory = new JdkClientHttpRequestFactory(httpClient); + requestFactory.setReadTimeout(READ_TIMEOUT); this.restClient = RestClient.builder() .baseUrl(baseUrl) + .requestFactory(requestFactory) .defaultHeader(HttpHeaders.AUTHORIZATION, "Bearer " + apiKey) .defaultHeader(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE) .build(); diff --git a/src/main/java/com/wok/supportbot/document/transform/MyKeywordEnricher.java b/src/main/java/com/wok/supportbot/document/transform/MyKeywordEnricher.java deleted file mode 100644 index fcf6aa5..0000000 --- a/src/main/java/com/wok/supportbot/document/transform/MyKeywordEnricher.java +++ /dev/null @@ -1,30 +0,0 @@ -package com.wok.supportbot.document.transform; - -import com.wok.supportbot.config.ChatModelFactory; -import jakarta.annotation.Resource; -import org.springframework.ai.document.Document; -import org.springframework.ai.model.transformer.KeywordMetadataEnricher; -import org.springframework.stereotype.Component; - -import java.util.List; - -/** - * 基于 AI 的文档元信息增强器(为文档补充元信息) - * 通过 ChatModelFactory 获取 ChatModel,支持多提供商动态切换 - */ -@Component -public class MyKeywordEnricher { - - @Resource - private ChatModelFactory chatModelFactory; - - /** - * 使用 AI 提取关键词并添加到元数据 - */ - public List enrichDocuments(List documents) { - KeywordMetadataEnricher enricher = new KeywordMetadataEnricher.Builder(chatModelFactory.getChatModel("CHAT")) - .keywordCount(5) - .build(); - return enricher.apply(documents); - } -} diff --git a/src/main/java/com/wok/supportbot/service/DocumentProcessingService.java b/src/main/java/com/wok/supportbot/service/DocumentProcessingService.java index fcb9a36..89c9504 100644 --- a/src/main/java/com/wok/supportbot/service/DocumentProcessingService.java +++ b/src/main/java/com/wok/supportbot/service/DocumentProcessingService.java @@ -1,38 +1,57 @@ package com.wok.supportbot.service; +import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper; import com.wok.supportbot.dao.KnowledgeDocumentMapper; -import com.wok.supportbot.document.transform.MyKeywordEnricher; import com.wok.supportbot.document.transform.MyTokenTextSplitter; import com.wok.supportbot.entity.KnowledgeDocument; import lombok.extern.slf4j.Slf4j; import org.springframework.ai.document.Document; import org.springframework.ai.vectorstore.VectorStore; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; import org.springframework.jdbc.core.JdbcTemplate; import org.springframework.scheduling.annotation.Async; import org.springframework.stereotype.Service; +import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.function.Supplier; +import java.util.stream.Collectors; /** * 文档异步处理服务 - * 负责文档的分块、关键词提取、向量化等耗时操作,在后台线程中执行 + * 负责文档的分块与向量化(按批隔离、失败可感知、有限自动重试),在后台线程中执行 */ @Service @Slf4j public class DocumentProcessingService { + /** 一批向量化的分块数:控制单批请求体积与"炸点"范围(含 volcengine-vision 逐条串行路径) */ + @Value("${knowledge.vector.batch-size:50}") + private int embedBatchSize = 50; + + /** + * 文档级自动重试次数(不含首次尝试)。 + * 注意:批内 EmbeddingModel 已自带 3 次指数退避重试(EmbeddingModelFactory.createRetryTemplate), + * 文档级重试叠加在最外层;文档级只对"整篇仍有失败批"触发,多数瞬时错误已在内层耗掉,外层命中率低。 + * 若厂商限流严重,可调小本值或调大 RETRY_BASE_DELAY_MS 冷却时间 + */ + private static final int MAX_DOC_LEVEL_RETRIES = 2; + + /** 重试等待基数(ms),指数放大(2s → 4s) */ + private static final long RETRY_BASE_DELAY_MS = 2_000L; + + /** 连续 N 批失败即中止本轮,交给文档级重试,避免打爆全部批次 */ + private static final int MAX_CONSECUTIVE_BATCH_FAILURES = 3; + @Autowired private KnowledgeDocumentMapper documentMapper; @Autowired private MyTokenTextSplitter myTokenTextSplitter; - @Autowired - private MyKeywordEnricher myKeywordEnricher; - @Autowired private VectorStore pgVectorVectorStore; @@ -40,18 +59,18 @@ public class DocumentProcessingService { private JdbcTemplate jdbcTemplate; /** - * 异步处理文档:分块 → 关键词提取 → 向量化 → 更新状态 - * 不加跨方法事务:AI 关键词提取与向量化属于慢速网络调用,期间不占用数据库连接, + * 异步处理文档:分块 → 向量化(分批入库)→ 更新状态 + * 不加跨方法事务:向量化属于慢速网络调用,期间不占用数据库连接, * 避免文件夹批量上传时多个异步任务把连接池占满导致连接超时。 * - * @param docId 文档ID - * @param documents 已解析的原始文档列表 + * @param docId 文档ID + * @param documents 已解析的原始文档列表 * @param sourceName 源文件名 - * @param title 文档标题 + * @param title 文档标题 * @param categoryId 分类ID - * @param tags 标签列表 - * @param chunkSize 分块大小(可选,覆盖全局配置) - * @param overlap 重叠大小(可选,覆盖全局配置) + * @param tags 标签列表 + * @param chunkSize 分块大小(可选,覆盖全局配置) + * @param overlap 重叠大小(可选,覆盖全局配置) */ @Async("documentExecutor") public void processDocumentAsync(Long docId, List documents, String sourceName, @@ -64,55 +83,19 @@ public class DocumentProcessingService { return; } - try { - // 1. 分块处理(使用 per-doc 参数或全局配置) - List splitDocuments = myTokenTextSplitter.splitDocuments(documents, chunkSize, overlap); - - // 2. 为每个分块设置 metadata - for (int i = 0; i < splitDocuments.size(); i++) { - Document d = splitDocuments.get(i); - Map meta = new HashMap<>(d.getMetadata()); - meta.put("documentId", String.valueOf(docId)); - meta.put("chunkIndex", i); - meta.put("sourceName", sourceName); - meta.put("title", title); - if (categoryId != null && categoryId > 0) { - meta.put("categoryId", String.valueOf(categoryId)); - } - if (tags != null && !tags.isEmpty()) { - meta.put("tags", tags); - } - // P1-2.1: 标记启用状态,用于 RAG 检索过滤 - meta.put("enabled", String.valueOf(Boolean.TRUE.equals(doc.getEnabled()))); - splitDocuments.set(i, new Document(d.getId(), d.getText(), meta)); - } - - // 3. 关键词提取 - List enrichedDocuments = myKeywordEnricher.enrichDocuments(splitDocuments); - - // 4. 向量化存储 - pgVectorVectorStore.add(enrichedDocuments); + DocMeta meta = DocMeta.of(docId, sourceName, title, categoryId, tags, + Boolean.TRUE.equals(doc.getEnabled())); - // 5. 更新状态为 READY - doc.setStatus("READY"); - doc.setChunkCount(enrichedDocuments.size()); - documentMapper.updateById(doc); - - log.info("异步处理文档完成: id={}, title={}, chunks={}", docId, title, enrichedDocuments.size()); - - } catch (Exception e) { - doc.setStatus("FAILED"); - doc.setErrorMessage(e.getMessage()); - documentMapper.updateById(doc); - log.error("异步处理文档失败: id={}, title={}", docId, title, e); - } + // 新文档首次处理无需清理(尚无向量);若中途失败触发整体重试,框架内会先清残留向量再重建 + runPipelineWithRetry(doc, meta, false, + () -> myTokenTextSplitter.splitDocuments(documents, chunkSize, overlap)); } /** - * 异步重新处理文档(重新分块 + 向量化) + * 异步重新处理文档(清理旧向量后按批重建向量化) * 使用文档存储的 per-doc 分块参数(extraConfig),如无则使用全局配置 * - * @param docId 文档ID + * @param docId 文档ID * @param documents 解析后的文档列表 */ @Async("documentExecutor") @@ -124,61 +107,265 @@ public class DocumentProcessingService { return; } - try { - // 删除旧向量 - String sql = "SELECT id::text FROM vector_store WHERE metadata->>'documentId' = ?"; - List oldIds = jdbcTemplate.queryForList(sql, String.class, String.valueOf(docId)); - if (!oldIds.isEmpty()) { - pgVectorVectorStore.delete(oldIds); - } + DocMeta meta = DocMeta.fromDocument(doc); - // 从 extraConfig 读取 per-doc 分块参数 + // reprocess 兼容:首轮就清理旧向量后按批重建;分块参数从 extraConfig 读取 + runPipelineWithRetry(doc, meta, true, () -> { Integer chunkSize = null; Integer overlap = null; if (doc.getExtraConfig() != null) { - Object cs = doc.getExtraConfig().get("chunkSize"); - Object ol = doc.getExtraConfig().get("overlap"); - if (cs instanceof Number) chunkSize = ((Number) cs).intValue(); - if (ol instanceof Number) overlap = ((Number) ol).intValue(); - } - - // 重新分块(使用 per-doc 参数或全局配置) - List splitDocuments = myTokenTextSplitter.splitDocuments(documents, chunkSize, overlap); - - for (int i = 0; i < splitDocuments.size(); i++) { - Document d = splitDocuments.get(i); - Map meta = new HashMap<>(d.getMetadata()); - meta.put("documentId", String.valueOf(docId)); - meta.put("chunkIndex", i); - meta.put("sourceName", doc.getSourceName()); - meta.put("title", doc.getTitle()); - if (doc.getCategoryId() != null && doc.getCategoryId() > 0) { - meta.put("categoryId", String.valueOf(doc.getCategoryId())); + if (doc.getExtraConfig().get("chunkSize") instanceof Number cs) chunkSize = cs.intValue(); + if (doc.getExtraConfig().get("overlap") instanceof Number ol) overlap = ol.intValue(); + } + return myTokenTextSplitter.splitDocuments(documents, chunkSize, overlap); + }); + } + + /** + * 文档级处理框架(外层兜底):切分 → 分批向量化 → 失败自动整体重试有限次数。 + * 外层 try-catch 保证任何未预期异常下文档状态都从 PROCESSING 收敛到 FAILED, + * 避免 @Async 异常被 Spring 静默吞掉导致文档永久悬挂(前端轮询永不结束) + */ + private void runPipelineWithRetry(KnowledgeDocument doc, DocMeta meta, + boolean cleanBeforeFirstAttempt, + Supplier> splitter) { + try { + doPipelineWithRetry(doc, meta, cleanBeforeFirstAttempt, splitter); + } catch (Exception e) { + // 兜底落 FAILED;若状态写入也失败(如 DB 抖动),仅记日志,避免再次上抛造成悬挂 + log.error("文档向量化处理出现未捕获异常: id={}", doc.getId(), e); + try { + String reason = e.getMessage() == null + ? e.getClass().getSimpleName() + : truncate(e.getMessage()); + patchStatus(doc.getId(), "FAILED", safeCountVectors(String.valueOf(doc.getId())), + "内部错误: " + reason); + } catch (Exception ex) { + log.error("兜底标记文档处理失败状态也失败: id={}, error={}", doc.getId(), ex.getMessage(), ex); + } + } + } + + private void doPipelineWithRetry(KnowledgeDocument doc, DocMeta meta, + boolean cleanBeforeFirstAttempt, + Supplier> splitter) { + String docIdStr = String.valueOf(doc.getId()); + + for (int attempt = 0; attempt <= MAX_DOC_LEVEL_RETRIES; attempt++) { + // 每次尝试前复查文档仍存在(用户可能在重试等待/sleep 期间删除文档),避免向已删文档写孤儿向量 + if (documentMapper.selectById(doc.getId()) == null) { + log.info("文档在处理期间已被删除,放弃处理: id={}", doc.getId()); + return; + } + + // 重试前:指数等待 + 清理上一轮残留向量,避免 split 重建产生新 Document.id 造成重复向量 + if (attempt > 0) { + if (!sleepBackoff(attempt)) { + patchStatus(doc.getId(), "FAILED", safeCountVectors(docIdStr), "处理被中断"); + return; } - if (doc.getTags() != null && doc.getTags().containsKey("tags")) { - meta.put("tags", doc.getTags().get("tags")); + deleteVectorsByDocumentId(docIdStr); + log.warn("文档向量化整体重试第 {}/{} 次: id={}", attempt, MAX_DOC_LEVEL_RETRIES, doc.getId()); + } else if (cleanBeforeFirstAttempt) { + deleteVectorsByDocumentId(docIdStr); + } + + List chunks; + try { + chunks = splitter.get(); + } catch (Exception e) { + if (attempt < MAX_DOC_LEVEL_RETRIES && isRetryable(e.getMessage())) { + continue; } - // P1-2.1: 标记启用状态 - meta.put("enabled", String.valueOf(Boolean.TRUE.equals(doc.getEnabled()))); - splitDocuments.set(i, new Document(d.getId(), d.getText(), meta)); + markFailed(doc, "文档分块失败: " + truncate(e.getMessage())); + return; } - List enrichedDocuments = myKeywordEnricher.enrichDocuments(splitDocuments); - pgVectorVectorStore.add(enrichedDocuments); + if (chunks == null || chunks.isEmpty()) { + // 分块为空属于内容/参数问题,重试不会自愈 + markFailed(doc, "文档分块结果为空(无可向量化文本或小于最小分块),请检查内容与分块参数"); + return; + } - doc.setStatus("READY"); - doc.setChunkCount(enrichedDocuments.size()); - doc.setErrorMessage(null); - documentMapper.updateById(doc); + List annotated = annotateChunks(chunks, meta); + BatchResult result = vectorizeInBatches(docIdStr, annotated, doc.getId()); - log.info("异步重新处理文档成功: id={}, title={}, chunks={}", docId, doc.getTitle(), enrichedDocuments.size()); + // 处理期间文档被删除:终止且不再写状态(删除是用户明确意图,不覆盖为 FAILED) + if (result.cancelled) { + return; + } + int stored = safeCountVectors(docIdStr); + if (!result.hasFailure()) { + patchStatus(doc.getId(), "READY", stored, null); + log.info("文档向量化成功: id={}, title={}, chunks={}", doc.getId(), doc.getTitle(), stored); + return; + } + + // 有失败批:若仍有余量且错误属瞬时/基础设施类,则整篇重试 + if (attempt < MAX_DOC_LEVEL_RETRIES && isRetryable(result.firstError)) { + log.warn("存在失败批,准备整体重试: id={}, total={}, stored={}, firstError={}", + doc.getId(), result.totalChunks, stored, result.firstError); + continue; + } + + // 已入库块保留在库(chunk_count 记实际入库数),缺失区间聚合进 error_message + markFailed(doc, result.buildSummary(stored)); + return; + } + } + + /** + * 分批向量化入库:按 embedBatchSize 逐批 add(),单批失败仅记录缺失区间后继续; + * 连续失败超过阈值且错误可重试时中止本轮(交给文档级重试),避免打满全部批次 + * + * @return 各批结果聚合(成功批已写入 DB) + */ + private BatchResult vectorizeInBatches(String docIdStr, List annotated, Long docId) { + BatchResult result = new BatchResult(annotated.size()); + int consecutiveFailures = 0; + int batchSize = effectiveBatchSize(); + + for (int start = 0; start < annotated.size(); start += batchSize) { + // 批前复查:文档若在异步处理期间被删除,立即终止后续写入,避免在 vector_store 留下孤儿向量 + if (docId != null && documentMapper.selectById(docId) == null) { + log.info("处理期间文档已被删除,终止剩余批次: docId={}", docIdStr); + result.cancelled = true; + break; + } + + int end = Math.min(start + batchSize, annotated.size()); + List batch = annotated.subList(start, end); + try { + pgVectorVectorStore.add(batch); + consecutiveFailures = 0; + log.info("批次向量化成功: docId={}, 第{}~{}块, 累计进度 {}/{}", + docIdStr, start + 1, end, Math.min(end, result.totalChunks), result.totalChunks); + } catch (Exception e) { + consecutiveFailures++; + // 缺失区间:0-based chunkIndex 闭区间 [start, end-1];错误消息截断防 error_message 过长/回显正文 + result.recordFailure(start, end - 1, truncate(e.getMessage())); + log.error("批次向量化失败: docId={}, 第{}~{}块, error={}", + docIdStr, start + 1, end, truncate(e.getMessage())); + if (consecutiveFailures >= MAX_CONSECUTIVE_BATCH_FAILURES && isRetryable(e.getMessage())) { + result.aborted = true; + // 连续失败中止本轮:把尚未执行的后续批次区间一并记入,避免 error_message 缺失段不完整 + if (end < annotated.size()) { + result.recordFailure(end, annotated.size() - 1, "连续失败已中止本轮,后续批次未执行"); + } + break; + } + } + } + return result; + } + + /** + * 为每个分块写入关联元数据:documentId/chunkIndex/sourceName/title/categoryId/tags/enabled + */ + private List annotateChunks(List splitDocuments, DocMeta meta) { + for (int i = 0; i < splitDocuments.size(); i++) { + Document d = splitDocuments.get(i); + Map m = new HashMap<>(d.getMetadata()); + m.put("documentId", String.valueOf(meta.docId())); + m.put("chunkIndex", i); + m.put("sourceName", meta.sourceName()); + m.put("title", meta.title()); + if (meta.categoryId() != null && meta.categoryId() > 0) { + m.put("categoryId", String.valueOf(meta.categoryId())); + } + if (meta.tagsValue() instanceof List tags && !tags.isEmpty()) { + m.put("tags", tags); + } + // P1-2.1: 标记启用状态,用于 RAG 检索过滤 + m.put("enabled", String.valueOf(meta.enabled())); + splitDocuments.set(i, new Document(d.getId(), d.getText(), m)); + } + return splitDocuments; + } + + // ==================== 私有工具 ==================== + + /** 仅 patch 状态/分块数/错误信息,避免覆盖用户在异步期间的并发编辑(如改标题/分类/toggle enabled) */ + private void patchStatus(Long docId, String status, Integer chunkCount, String errorMessage) { + documentMapper.update(null, new LambdaUpdateWrapper() + .eq(KnowledgeDocument::getId, docId) + .set(KnowledgeDocument::getStatus, status) + .set(KnowledgeDocument::getChunkCount, chunkCount) + .set(KnowledgeDocument::getErrorMessage, errorMessage)); + } + + private void markFailed(KnowledgeDocument doc, String message) { + String docIdStr = String.valueOf(doc.getId()); + patchStatus(doc.getId(), "FAILED", safeCountVectors(docIdStr), message); + log.error("文档向量化失败: id={}, title={}, error={}", doc.getId(), doc.getTitle(), message); + } + + /** 统计实际已入库向量数(失败时容错,DB 抖动不阻断状态收敛) */ + private int safeCountVectors(String docIdStr) { + try { + return countStoredVectors(docIdStr); } catch (Exception e) { - doc.setStatus("FAILED"); - doc.setErrorMessage(e.getMessage()); - documentMapper.updateById(doc); - log.error("异步重新处理文档失败: id={}, title={}", docId, doc.getTitle(), e); + log.warn("统计已入库向量数失败: docId={}, error={}", docIdStr, e.getMessage()); + return 0; + } + } + + /** 截断错误消息,避免 error_message 过长或回显大段文档文本 */ + private static String truncate(String msg) { + if (msg == null) { + return null; + } + return msg.length() <= 500 ? msg : msg.substring(0, 500) + "…(已截断)"; + } + + /** 实际已入库向量条数:以 DB 为准,天然吸收"失败批内偶发的部分写入" */ + private int countStoredVectors(String docIdStr) { + Integer c = jdbcTemplate.queryForObject( + "SELECT count(*) FROM vector_store WHERE metadata->>'documentId' = ?", + Integer.class, docIdStr); + return c == null ? 0 : c; + } + + /** 物理删除该文档在 vector_store 的全部向量(vector_store 无逻辑删除语义) */ + private void deleteVectorsByDocumentId(String docIdStr) { + List ids = jdbcTemplate.queryForList( + "SELECT id::text FROM vector_store WHERE metadata->>'documentId' = ?", + String.class, docIdStr); + if (!ids.isEmpty()) { + pgVectorVectorStore.delete(ids); + log.debug("清理旧向量: documentId={}, count={}", docIdStr, ids.size()); + } + } + + private int effectiveBatchSize() { + return embedBatchSize > 0 ? embedBatchSize : 50; + } + + /** 简单指数退避等待;被中断时复位中断标志并返回 false(放弃处理) */ + private static boolean sleepBackoff(int attempt) { + long ms = RETRY_BASE_DELAY_MS * (1L << (attempt - 1)); + try { + Thread.sleep(ms); + return true; + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return false; + } + } + + /** 仅对疑似瞬时/基础设施类错误做文档级重试,4xx 参数类错误不空转 */ + private static boolean isRetryable(String msg) { + // 消息缺失的未知错误不做整篇重试(避免对 NPE 等非瞬时错误空转 2 次),仅对明确瞬时/限流特征重试 + if (msg == null || msg.isBlank()) { + return false; } + String m = msg.toLowerCase(); + return m.contains("timeout") || m.contains("timed out") || m.contains("connection") + || m.contains("429") || m.contains("too many request") || m.contains("rate limit") + || m.contains(" 500") || m.contains(" 502") || m.contains(" 503") || m.contains(" 504") + || m.contains("socket") || m.contains("i/o error") || m.contains("internal server error") + || m.contains("service unavailable") || m.contains("bad gateway"); } /** @@ -206,4 +393,73 @@ public class DocumentProcessingService { } return null; } + + /** + * 分批向量化结果聚合:成功批已入库;失败批记录缺失 chunkIndex 区间与首个错误 + */ + private static final class BatchResult { + final int totalChunks; + /** 每项为 [start, end](0-based chunkIndex 闭区间),最多保留前 8 段 */ + final List failedRanges = new ArrayList<>(); + String firstError; + boolean aborted; + /** 处理期间文档被删除时置位:调用方应终止且不再写状态 */ + boolean cancelled; + + BatchResult(int totalChunks) { + this.totalChunks = totalChunks; + } + + void recordFailure(int start, int end, String error) { + failedRanges.add(new int[]{start, end}); + if (firstError == null) { + firstError = error; + } + } + + boolean hasFailure() { + return !failedRanges.isEmpty(); + } + + /** 例:已入库 320/400 块,缺失第 321~400 块 向量化失败: <原因> */ + String buildSummary(int storedChunks) { + String ranges = failedRanges.stream().limit(8) + .map(r -> r[0] == r[1] + ? "第 " + (r[0] + 1) + " 块" + : "第 " + (r[0] + 1) + "~" + (r[1] + 1) + " 块") + .collect(Collectors.joining("、")); + if (failedRanges.size() > 8) { + ranges += "…共 " + failedRanges.size() + " 段"; + } + String msg = "已入库 " + storedChunks + "/" + totalChunks + " 块,缺失 " + ranges + " 向量化失败"; + if (aborted) { + msg += "(连续失败已中止本轮)"; + } + if (firstError != null) { + msg += ": " + firstError; + } + return msg; + } + } + + /** + * 分块元数据值对象:统一两个异步入口(新上传/reprocess)的 metadata 标注逻辑 + */ + private record DocMeta(Long docId, String sourceName, String title, Long categoryId, + Object tagsValue, boolean enabled) { + + static DocMeta of(Long docId, String sourceName, String title, Long categoryId, + List tags, boolean enabled) { + return new DocMeta(docId, sourceName, title, categoryId, tags, enabled); + } + + static DocMeta fromDocument(KnowledgeDocument doc) { + Object tagsValue = null; + if (doc.getTags() != null && doc.getTags().containsKey("tags")) { + tagsValue = doc.getTags().get("tags"); + } + return new DocMeta(doc.getId(), doc.getSourceName(), doc.getTitle(), + doc.getCategoryId(), tagsValue, Boolean.TRUE.equals(doc.getEnabled())); + } + } } diff --git a/src/main/java/com/wok/supportbot/service/DocumentService.java b/src/main/java/com/wok/supportbot/service/DocumentService.java index 84a6250..6049c04 100644 --- a/src/main/java/com/wok/supportbot/service/DocumentService.java +++ b/src/main/java/com/wok/supportbot/service/DocumentService.java @@ -7,7 +7,6 @@ import com.wok.supportbot.document.extract.JsonDocumentLoader; import com.wok.supportbot.document.extract.MarkdownDocumentLoader; import com.wok.supportbot.document.extract.SimpleStringDocumentReader; import com.wok.supportbot.document.extract.TikaDocumentReader; -import com.wok.supportbot.document.transform.MyKeywordEnricher; import com.wok.supportbot.document.transform.MyTokenTextSplitter; import com.wok.supportbot.entity.CategoryNode; import com.wok.supportbot.entity.KnowledgeCategory; @@ -56,9 +55,6 @@ public class DocumentService { @Autowired private MyTokenTextSplitter myTokenTextSplitter; - @Autowired - private MyKeywordEnricher myKeywordEnricher; - @Autowired private TikaDocumentReader tikaDocumentReader; @@ -937,9 +933,7 @@ public class DocumentService { meta.put("enabled", String.valueOf(Boolean.TRUE.equals(doc.getEnabled()))); Document newDoc = new Document(vectorId, newContent, meta); - // 关键词提取 - List enriched = myKeywordEnricher.enrichDocuments(List.of(newDoc)); - pgVectorVectorStore.add(enriched); + pgVectorVectorStore.add(List.of(newDoc)); log.info("更新分块: docId={}, chunkIndex={}, vectorId={}", docId, chunkIndex, vectorId); } diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index 2cabaa2..5f01f00 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -58,6 +58,9 @@ knowledge: # 向量维度,需与 Embedding 模型输出维度一致 # 千问 text-embedding-v2: 1024 | 豆包 doubao-embedding-text-240515: 2048 | OpenAI text-embedding-3-small: 1536 dimension: 1024 + # 向量化分批大小:单批分块数(逐批入库,失败批只记录缺失区间不扩散到整篇) + # 默认 50;豆包多模态(vision)模型逐条调用较慢,若大量使用可调小(如 20)减小单批串行耗时 + batch-size: 50 role: # 严格隔离:true=角色未绑定知识库分类时禁止检索任何内容;false=可检索全部知识库 strict-isolation: false