
前面我們進(jìn)行了文件處理的工作我們實(shí)現(xiàn)了文檔的解析分塊上傳服務(wù)理論上功能已經(jīng)實(shí)現(xiàn)比較完整了接著我們會(huì)發(fā)現(xiàn)程序中存在的另一個(gè)問題在我們進(jìn)行文件上傳過程中我們程序的進(jìn)程被阻塞了這是因?yàn)槲募纳蟼鹘馕龇謮K工作并不能馬上完成對(duì)于較大的文檔我們的處理時(shí)間要比原有的更長因此這將會(huì)導(dǎo)致我們的體驗(yàn)感不佳我們需要想辦法進(jìn)行處理使它能夠在處理文檔的同時(shí)保證我們的其他操作不受影響這就是我們今天要做的事情 – 異步處理。異步處理的核心概念對(duì)于異步處理來說最為核心的概念是發(fā)起的任務(wù)和任務(wù)的結(jié)果是解耦的。異步處理的三要素回調(diào)函數(shù)把處理結(jié)果的函數(shù)作為參數(shù)傳進(jìn)去任務(wù)完成后自動(dòng)調(diào)用。Promise / Future返回一個(gè)“承諾對(duì)象”可以稍后通過.then()或await來獲取結(jié)果。事件驅(qū)動(dòng)就像瀏覽器有一個(gè)循環(huán)不斷檢查“有沒有任務(wù)做完了有的話就執(zhí)行對(duì)應(yīng)的回調(diào)?!痹?WenQu 中是怎么實(shí)現(xiàn)的了解完了基本的概念我們來看看在 WenQu 中我是怎么實(shí)現(xiàn)的。首先我們需要搭建異步處理的基本設(shè)施AsyncConfig這個(gè)配置類可以使我們打開 spring 的異步處理功能ConfigurationEnableAsync// 打開 Spring 的異步開關(guān)publicclassAsyncConfig{Bean(docTaskExecutor)publicExecutordocTaskExecutor(){ThreadPoolTaskExecutorexecutornewThreadPoolTaskExecutor();executor.setCorePoolSize(2);// 常駐線程executor.setMaxPoolSize(4);// 最大線程executor.setQueueCapacity(200);// 排隊(duì)等待的任務(wù)數(shù)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);}這段代碼就是輪詢的方法當(dāng)我們進(jìn)行文件上傳啟動(dòng)這個(gè)異步處理時(shí)輪詢開啟每隔 5s 就進(jìn)行查詢看看處理的狀態(tài)。接著我們繼續(xù)我們先整體看看/** * 異步處理實(shí)現(xiàn)類 */Slf4jServiceRequiredArgsConstructorpublicclassDocumentProcessServiceImplimplementsDocumentProcessService{privatefinalDocumentMapperdocumentMapper;privatefinalDocChunkMapperdocChunkMapper;privatefinalTextExtractorFactoryextractorFactory;privatefinalKnowledgeBaseMapperknowledgeBaseMapper;privatefinalTextChunkersentenceChunker;OverrideAsync(docTaskExecutor)publicvoidprocess(LongdocumentId){// 查文檔DocumentdocdocumentMapper.selectEntityById(documentId);if(docnull){log.warn(文檔不存在{},documentId);return;}// 狀態(tài)設(shè)置為解析中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(解析結(jié)果為空);}documentMapper.updateStatus(documentId,DocumentStatus.CHUNKING,null);// 查知識(shí)庫的 chunk_size 和 overlapKnowledgeBasekbknowledgeBaseMapper.selectConfigByIdAndUserId(doc.getKbId(),doc.getUserId());if(kbnull){thrownewIllegalStateException(知識(shí)庫不存在或無權(quán)訪問);}// 切分intchunkSizekb.getChunkSize()null?400:kb.getChunkSize();intoverlapkb.getOverlap()null?80:kb.getOverlap();ListStringchunkssentenceChunker.chunk(text,chunkSize,overlap);// 組裝成實(shí)體并批量入庫seq 從 0 開始編號(hào)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);// 回填塊數(shù)置成功documentMapper.updateChunkCount(documentId,docChunks.size());documentMapper.updateStatus(documentId,DocumentStatus.READY,null);log.info([{}] 處理完成共 {} 塊,doc.getName(),docChunks.size());}catch(Exceptione){// 任何一步失敗 → 記 FAILED 錯(cuò)誤信息// error_msg 列是 varchar(1000)異常棧太長會(huì)截?cái)鄨?bào)錯(cuò)只保留簡短摘要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)來啟動(dòng)異步處理spring 讀取到了這個(gè)注解后就會(huì)攔截這個(gè)方法進(jìn)行后續(xù)的異步請(qǐng)求處理。這部分是異步處理的內(nèi)部方法可以注意到在這個(gè)方法內(nèi)我們是進(jìn)行了很多的數(shù)據(jù)庫操作的為什么我們沒有使用事務(wù)去保證數(shù)據(jù)一致性呢這是因?yàn)槲覀冊(cè)谶@里面的操作對(duì)數(shù)據(jù)庫操作時(shí)其實(shí)本質(zhì)上也進(jìn)行了一致性校驗(yàn)要是數(shù)據(jù)出現(xiàn)不一致問題程序就會(huì)進(jìn)行報(bào)錯(cuò)而事務(wù)這是我們刻意設(shè)計(jì)的因?yàn)楫惒教幚碛械牟僮餍枰喈?dāng)長時(shí)間使用事務(wù)會(huì)阻塞其他操作。我們?cè)賮砜纯慈刂械钠渌麄z/** * 文檔功能實(shí)現(xiàn)類 */ServiceSlf4jRequiredArgsConstructorpublicclassDocumentServiceImplimplementsDocumentService{privatefinalKnowledgeBaseServiceknowledgeBaseService;privatefinalDocumentMapperdocumentMapper;privatefinalTextExtractorFactoryextractorFactory;privatefinalDocumentProcessServicedocumentProcessService;/** * 文件上傳根目錄 */Value(${wenqu.upload-dir:./uploads})privateStringuploadDir;/** * 文件上傳文件落盤MySQL 只存文件路徑 */OverrideTransactionalpublicDocumentVOupload(LongkbId,MultipartFilefile,LonguserId)throwsIOException{// 查庫是否存在并校驗(yàn)身份knowledgeBaseService.getKnowledgeBase(kbId,userId);// 校驗(yàn)文件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}/{時(shí)間戳}_{原文件名}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)// 設(shè)置成中間狀態(tài).build();documentMapper.insert(doc);// 更新文件路徑和狀態(tài)doc.setFilePath(target.toString());doc.setStatus(DocumentStatus.UPLOADED);documentMapper.updateFilePath(doc);// 觸發(fā)后臺(tái)異步處理解析 → 切分 → 入庫// 必須在事務(wù)提交后觸發(fā)否則異步線程查不到剛插入的文檔TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){OverridepublicvoidafterCommit(){documentProcessService.process(doc.getId());}});// 返回VO對(duì)象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){// 校驗(yàn)知識(shí)庫存在且屬于當(dāng)前用戶knowledgeBaseService.validateKnowledgeBase(kbId,userId);// 開啟分頁查詢PageHelper.startPage(page,pageSize);// 調(diào)mapper層查詢PageDocumentVOpagesdocumentMapper.query(kbId,page,pageSize);Longtotalpages.getTotal();ListDocumentVOrecordspages.getResult();returnnewPageResult(records,total,page,pageSize);}/** * 文檔詳情先查文檔再校驗(yàn)所屬知識(shí)庫歸屬 */OverridepublicDocumentVOgetDocument(Longid,LonguserId){DocumentVOdocdocumentMapper.selectById(id);if(docnull){thrownewBusinessException(ResultCode.DOCUMENT_NOT_FOUND);}// 校驗(yàn)所屬知識(shí)庫存在且屬于當(dāng)前用戶knowledgeBaseService.validateKnowledgeBase(doc.getKbId(),userId);returndoc;}}這是文檔處理的詳細(xì)代碼我們重點(diǎn)來看看下面這段要非常注意的事情是我們?cè)谶M(jìn)行異步處理時(shí)必須在事務(wù)提交后觸發(fā)否則異步線程會(huì)找不到剛插入的文檔。// 觸發(fā)后臺(tái)異步處理解析 → 切分 → 入庫// 必須在事務(wù)提交后觸發(fā)否則異步線程查不到剛插入的文檔TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){OverridepublicvoidafterCommit(){documentProcessService.process(doc.getId());}});這里就是我們所說的函數(shù)回調(diào)但是我們也發(fā)現(xiàn)并沒有使用到.then()等方法這是因?yàn)樵谖覀冞@個(gè)處理中不需要使用到我們就簡單使用注解Async去解決了。至此我們的異步處理就解決了我們采用的是較為輕量型的方案對(duì)于本項(xiàng)目來說個(gè)人學(xué)習(xí)已經(jīng)夠用當(dāng)然我們不會(huì)止步于此為了更高的并發(fā)更安全的線程策略以及處理因?yàn)榉?wù)器宕機(jī)導(dǎo)致任務(wù)被截?cái)喽a(chǎn)生的“僵尸任務(wù)”問題我們后續(xù)將重構(gòu)這部分代碼引入更為規(guī)范的方法 – 消息隊(duì)列。當(dāng)然這并不是我們現(xiàn)階段要做的事情了。后面我們首先要做的就是向量化。我是 _AgAiN請(qǐng)見證我的學(xué)習(xí)之路。項(xiàng)目鏈接WenQu