ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

WenQu -- 异步处理

WenQu -- 异步处理 前面我们进行了文件处理的工作我们实现了文档的解析分块上传服务理论上功能已经实现比较完整了接着我们会发现程序中存在的另一个问题在我们进行文件上传过程中我们程序的进程被阻塞了这是因为文件的上传解析分块工作并不能马上完成对于较大的文档我们的处理时间要比原有的更长因此这将会导致我们的体验感不佳我们需要想办法进行处理使它能够在处理文档的同时保证我们的其他操作不受影响这就是我们今天要做的事情 – 异步处理。异步处理的核心概念对于异步处理来说最为核心的概念是发起的任务和任务的结果是解耦的。异步处理的三要素回调函数把处理结果的函数作为参数传进去任务完成后自动调用。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
返回列表