
前面我們進行了文件處理的工作我們實現了文檔的解析分塊上傳服務理論上功能已經實現比較完整了接著我們會發現程序中存在的另一個問題在我們進行文件上傳過程中我們程序的進程被阻塞了這是因為文件的上傳解析分塊工作并不能馬上完成對于較大的文檔我們的處理時間要比原有的更長因此這將會導致我們的體驗感不佳我們需要想辦法進行處理使它能夠在處理文檔的同時保證我們的其他操作不受影響這就是我們今天要做的事情 – 異步處理。異步處理的核心概念對于異步處理來說最為核心的概念是發起的任務和任務的結果是解耦的。異步處理的三要素回調函數把處理結果的函數作為參數傳進去任務完成后自動調用。Promise / Future返回一個“承諾對象”可以稍后通過.then()或await來獲取結果。事件驅動就像瀏覽器有一個循環不斷檢查“有沒有任務做完了有的話就執行對應的回調。”在 WenQu 中是怎么實現的了解完了基本的概念我們來看看在 WenQu 中我是怎么實現的。首先我們需要搭建異步處理的基本設施AsyncConfig這個配置類可以使我們打開 spring 的異步處理功能ConfigurationEnableAsync// 打開 Spring 的異步開關publicclassAsyncConfig{Bean(docTaskExecutor)publicExecutordocTaskExecutor(){ThreadPoolTaskExecutorexecutornewThreadPoolTaskExecutor();executor.setCorePoolSize(2);// 常駐線程executor.setMaxPoolSize(4);// 最大線程executor.setQueueCapacity(200);// 排隊等待的任務數executor.setThreadNamePrefix(doc-parse-);executor.initialize();returnexecutor;}}接著我們來瞧瞧三要素我們先從前端看起functionpollDoc(kbId,docId){if(state.docPollTimers[docId])return;state.docPollTimers[docId]setInterval(async(){try{constdawaitapi(GET,/documents/${docId});if(!d)return;if(state.currentKb?.id!kbId)return;constidxstate.docs.findIndex(xx.iddocId);if(idx0)state.docs[idx]d;if([READY,FAILED].includes(d.status)){clearInterval(state.docPollTimers[docId]);deletestate.docPollTimers[docId];renderDocs();if(state.currentKb)loadKbs();}}catch(err){if(errinstanceofApiError(err.code40400||err.code40401)){clearInterval(state.docPollTimers[docId]);deletestate.docPollTimers[docId];if(state.currentKb?.idkbId)renderDocs();}elseif(errinstanceofApiError(err.code40100||err.code40101)){clearInterval(state.docPollTimers[docId]);deletestate.docPollTimers[docId];}}},5000);}這段代碼就是輪詢的方法當我們進行文件上傳啟動這個異步處理時輪詢開啟每隔 5s 就進行查詢看看處理的狀態。接著我們繼續我們先整體看看/** * 異步處理實現類 */Slf4jServiceRequiredArgsConstructorpublicclassDocumentProcessServiceImplimplementsDocumentProcessService{privatefinalDocumentMapperdocumentMapper;privatefinalDocChunkMapperdocChunkMapper;privatefinalTextExtractorFactoryextractorFactory;privatefinalKnowledgeBaseMapperknowledgeBaseMapper;privatefinalTextChunkersentenceChunker;OverrideAsync(docTaskExecutor)publicvoidprocess(LongdocumentId){// 查文檔DocumentdocdocumentMapper.selectEntityById(documentId);if(docnull){log.warn(文檔不存在{},documentId);return;}// 狀態設置為解析中documentMapper.updateStatus(documentId,DocumentStatus.PARSING,null);log.info([{}]開始處理,doc.getName());try{// 解析StringtextextractorFactory.get(doc.getType()).extract(Paths.get(doc.getFilePath()));if(textnull||text.isBlank()){thrownewIllegalStateException(解析結果為空);}documentMapper.updateStatus(documentId,DocumentStatus.CHUNKING,null);// 查知識庫的 chunk_size 和 overlapKnowledgeBasekbknowledgeBaseMapper.selectConfigByIdAndUserId(doc.getKbId(),doc.getUserId());if(kbnull){thrownewIllegalStateException(知識庫不存在或無權訪問);}// 切分intchunkSizekb.getChunkSize()null?400:kb.getChunkSize();intoverlapkb.getOverlap()null?80:kb.getOverlap();ListStringchunkssentenceChunker.chunk(text,chunkSize,overlap);// 組裝成實體并批量入庫seq 從 0 開始編號ListDocChunkdocChunksIntStream.range(0,chunks.size()).mapToObj(i-DocChunk.builder().documentId(documentId).kbId(doc.getKbId()).seq(i).content(chunks.get(i)).build()).collect(Collectors.toList());docChunkMapper.batchInsert(docChunks);// 回填塊數置成功documentMapper.updateChunkCount(documentId,docChunks.size());documentMapper.updateStatus(documentId,DocumentStatus.READY,null);log.info([{}] 處理完成共 {} 塊,doc.getName(),docChunks.size());}catch(Exceptione){// 任何一步失敗 → 記 FAILED 錯誤信息// error_msg 列是 varchar(1000)異常棧太長會截斷報錯只保留簡短摘要Stringmsge.getMessage()null?e.getClass().getSimpleName():e.getMessage();if(msg.length()900){msgmsg.substring(0,900);}log.error([{}] 處理失敗,doc.getName(),e);documentMapper.updateStatus(documentId,DocumentStatus.FAILED,msg);}}}在具體方法里我們使用注解Async(docTaskExecutor)來啟動異步處理spring 讀取到了這個注解后就會攔截這個方法進行后續的異步請求處理。這部分是異步處理的內部方法可以注意到在這個方法內我們是進行了很多的數據庫操作的為什么我們沒有使用事務去保證數據一致性呢這是因為我們在這里面的操作對數據庫操作時其實本質上也進行了一致性校驗要是數據出現不一致問題程序就會進行報錯而事務這是我們刻意設計的因為異步處理有的操作需要相當長時間使用事務會阻塞其他操作。我們再來看看三要素中的其他倆/** * 文檔功能實現類 */ServiceSlf4jRequiredArgsConstructorpublicclassDocumentServiceImplimplementsDocumentService{privatefinalKnowledgeBaseServiceknowledgeBaseService;privatefinalDocumentMapperdocumentMapper;privatefinalTextExtractorFactoryextractorFactory;privatefinalDocumentProcessServicedocumentProcessService;/** * 文件上傳根目錄 */Value(${wenqu.upload-dir:./uploads})privateStringuploadDir;/** * 文件上傳文件落盤MySQL 只存文件路徑 */OverrideTransactionalpublicDocumentVOupload(LongkbId,MultipartFilefile,LonguserId)throwsIOException{// 查庫是否存在并校驗身份knowledgeBaseService.getKnowledgeBase(kbId,userId);// 校驗文件if(filenull||file.isEmpty()){thrownewBusinessException(ResultCode.FILE_EMPTY);}StringnameObjects.requireNonNull(file.getOriginalFilename());Stringextname.contains(.)?name.substring(name.lastIndexOf(.)1).toLowerCase():;if(!Set.of(txt,md,doc,docx,pdf,xls,xlsx,ppt,pptx,html,csv,epub).contains(ext)){thrownewBusinessException(ResultCode.UNSUPPORTED_FILE_TYPE);}if(file.getSize()20L*1024*1024){thrownewBusinessException(ResultCode.FILE_TOO_LARGE);}// 先落盤./uploads/{userId}/{kbId}/{時間戳}_{原文件名}PathdirPaths.get(uploadDir,String.valueOf(userId),String.valueOf(kbId));Files.createDirectories(dir);Pathtargetdir.resolve(System.currentTimeMillis()_name).toAbsolutePath();try{file.transferTo(target);// 寫庫DocumentdocDocument.builder().kbId(kbId).userId(userId).name(name).type(ext).size(file.getSize()).status(DocumentStatus.UPLOADING)// 設置成中間狀態.build();documentMapper.insert(doc);// 更新文件路徑和狀態doc.setFilePath(target.toString());doc.setStatus(DocumentStatus.UPLOADED);documentMapper.updateFilePath(doc);// 觸發后臺異步處理解析 → 切分 → 入庫// 必須在事務提交后觸發否則異步線程查不到剛插入的文檔TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){OverridepublicvoidafterCommit(){documentProcessService.process(doc.getId());}});// 返回VO對象returnDocumentVO.builder().id(doc.getId()).kbId(kbId).name(name).type(ext).size(doc.getSize()).status(DocumentStatus.UPLOADED).chunkCount(0L).createdAt(System.currentTimeMillis()).build();}catch(Exceptione){// 文件已落盤但 DB 寫入失敗 → 刪除文件避免孤兒文件Files.deleteIfExists(target);throwe;}}/** * 文章列表查詢 */OverridepublicPageResultDocumentVOpageQuery(LongkbId,LonguserId,intpage,intpageSize){// 校驗知識庫存在且屬于當前用戶knowledgeBaseService.validateKnowledgeBase(kbId,userId);// 開啟分頁查詢PageHelper.startPage(page,pageSize);// 調mapper層查詢PageDocumentVOpagesdocumentMapper.query(kbId,page,pageSize);Longtotalpages.getTotal();ListDocumentVOrecordspages.getResult();returnnewPageResult(records,total,page,pageSize);}/** * 文檔詳情先查文檔再校驗所屬知識庫歸屬 */OverridepublicDocumentVOgetDocument(Longid,LonguserId){DocumentVOdocdocumentMapper.selectById(id);if(docnull){thrownewBusinessException(ResultCode.DOCUMENT_NOT_FOUND);}// 校驗所屬知識庫存在且屬于當前用戶knowledgeBaseService.validateKnowledgeBase(doc.getKbId(),userId);returndoc;}}這是文檔處理的詳細代碼我們重點來看看下面這段要非常注意的事情是我們在進行異步處理時必須在事務提交后觸發否則異步線程會找不到剛插入的文檔。// 觸發后臺異步處理解析 → 切分 → 入庫// 必須在事務提交后觸發否則異步線程查不到剛插入的文檔TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){OverridepublicvoidafterCommit(){documentProcessService.process(doc.getId());}});這里就是我們所說的函數回調但是我們也發現并沒有使用到.then()等方法這是因為在我們這個處理中不需要使用到我們就簡單使用注解Async去解決了。至此我們的異步處理就解決了我們采用的是較為輕量型的方案對于本項目來說個人學習已經夠用當然我們不會止步于此為了更高的并發更安全的線程策略以及處理因為服務器宕機導致任務被截斷而產生的“僵尸任務”問題我們后續將重構這部分代碼引入更為規范的方法 – 消息隊列。當然這并不是我們現階段要做的事情了。后面我們首先要做的就是向量化。我是 _AgAiN請見證我的學習之路。項目鏈接WenQu