diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java index 14b1f2a..240c631 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java @@ -3,6 +3,7 @@ package com.volvo.ai.analytic.center.job; import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSONObject; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.volvo.ai.analytic.center.dto.PageDto; import com.volvo.ai.analytic.center.dto.corpus.CorpusReportDTO; import com.volvo.ai.analytic.center.dto.req.DiFyReq; import com.volvo.ai.analytic.center.entity.AiAnalysisErrors; @@ -53,7 +54,8 @@ public class CorpusFailJob { @Resource private RocketMQTemplate rocketMqTemplate; - + @Value("${batch.size}") + public int pageSize = 100; @Autowired private TmTelephoneCorpusService tmTelephoneCorpusService; @@ -65,69 +67,78 @@ public class CorpusFailJob { public void corpusFailTask() { log.info("语料解析失败重试处理"); - List aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(BusinessTypeEnum.SMART_ASSISTANT.getCode()); - if(CollectionUtils.isNotEmpty(aiAnalysisErrorsListlist)) { - log.info("语料解析失败重试处理 size:{}", aiAnalysisErrorsListlist.size()); - aiAnalysisErrorsListlist.stream().forEach(aiAnalysisErrors -> { + Integer total = aiAnalysisErrorsService.queryCountAnalysisErrorList(BusinessTypeEnum.SMART_ASSISTANT.getCode()); + log.info("语料解析失败重试处理数据量:{}", total); + int totalPages = PageDto.getTotalPages(total, pageSize); + log.info("语料解析失败重试处理数据量:{},总页数:{}", total, totalPages); + for (int i = 1; i <= totalPages; i++) { + int offset = (i - 1) * pageSize; + + List aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(BusinessTypeEnum.SMART_ASSISTANT.getCode(),offset, pageSize); + if(CollectionUtils.isNotEmpty(aiAnalysisErrorsListlist)) { + log.info("语料解析失败重试处理 size:{}", aiAnalysisErrorsListlist.size()); + aiAnalysisErrorsListlist.stream().forEach(aiAnalysisErrors -> { - try { - LambdaQueryWrapper queryWrapper = new LambdaQueryWrapper<>(); - queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisErrors.getAiAnalysisRequestId()); - AiAnalysisRequestLogs oldAiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(queryWrapper); + try { + LambdaQueryWrapper queryWrapper = new LambdaQueryWrapper<>(); + queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisErrors.getAiAnalysisRequestId()); + AiAnalysisRequestLogs oldAiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(queryWrapper); - if (null != oldAiAnalysisRequestLogs) { - DiFyReq diFyReq = JSONObject.parseObject(oldAiAnalysisRequestLogs.getDifyRequest(), DiFyReq.class); - CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class); + if (null != oldAiAnalysisRequestLogs) { + DiFyReq diFyReq = JSONObject.parseObject(oldAiAnalysisRequestLogs.getDifyRequest(), DiFyReq.class); + CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class); + + JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyReq); + log.info(" corpusFailTask runDify execDifyFlow {}", execDifyFlow); + if (null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")) { + String text = execDifyFlow.getJSONObject("outputs").getString("text"); + String resultStrOne = FlowResultSplitUtil.flowOutputTextSplit(text, "任务1", "任务2"); + String resultStrTwo = FlowResultSplitUtil.flowOutputTextSplit(text, "任务2", null); + if (StringUtils.isBlank(resultStrOne) || StringUtils.isBlank(resultStrTwo)) { + log.info(" corpusFailTask 企微语料解析为空,text:{}", text); + return; + } + Map ltoMap = new HashMap<>(); + ltoMap.put("analysisRecordId", oldAiAnalysisRequestLogs.getAiAnalysisRequestId()); + ltoMap.put("analysisScene", corpusReportDTO.getAnalysisScene() + ""); + String tag = ""; + if (null != corpusReportDTO && corpusReportDTO.getAnalysisScene() == 1) { //企微 + ltoMap.put("unionId", corpusReportDTO.getUnionId()); + ltoMap.put("consultantId", corpusReportDTO.getUserId()); + tag = CategoryEnum.ENTERPRISE_WECHAT.getCode(); + } else { + ltoMap.put("recordId", corpusReportDTO.getUserId()); + tag = CategoryEnum.PHONE_VOICE.getCode(); + } + + ltoMap.put("communicateDate", corpusReportDTO.getCorpusTime()); + ltoMap.put("analysisResult", resultStrOne.replace("#", "")); + ltoMap.put("analysisDetail", resultStrTwo); + // 发送MQ + log.info("corpusFailTask send mq {}", ltoMap); + tmTelephoneCorpusService.sendMq(tag, JSONObject.toJSONString(ltoMap)); + // 保存报告 + oldAiAnalysisRequestLogs.setBusinessResponse(JSONObject.toJSONString(ltoMap)); - JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyReq); - log.info(" corpusFailTask runDify execDifyFlow {}",execDifyFlow); - if(null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")){ - String text = execDifyFlow.getJSONObject("outputs").getString("text"); - String resultStrOne = FlowResultSplitUtil.flowOutputTextSplit(text, "任务1", "任务2"); - String resultStrTwo = FlowResultSplitUtil.flowOutputTextSplit(text, "任务2", null); - if (StringUtils.isBlank(resultStrOne) || StringUtils.isBlank(resultStrTwo)) { - log.info(" corpusFailTask 企微语料解析为空,text:{}", text); - return; } - Map ltoMap = new HashMap<>(); - ltoMap.put("analysisRecordId", oldAiAnalysisRequestLogs.getAiAnalysisRequestId()); - ltoMap.put("analysisScene", corpusReportDTO.getAnalysisScene() + ""); - String tag = ""; - if (null != corpusReportDTO && corpusReportDTO.getAnalysisScene() == 1) { //企微 - ltoMap.put("unionId", corpusReportDTO.getUnionId()); - ltoMap.put("consultantId", corpusReportDTO.getUserId()); - tag = CategoryEnum.ENTERPRISE_WECHAT.getCode(); - } else { - ltoMap.put("recordId", corpusReportDTO.getUserId()); - tag = CategoryEnum.PHONE_VOICE.getCode(); - } - - ltoMap.put("communicateDate", corpusReportDTO.getCorpusTime()); - ltoMap.put("analysisResult", resultStrOne.replace("#", "")); - ltoMap.put("analysisDetail", resultStrTwo); - // 发送MQ - log.info("corpusFailTask send mq {}", ltoMap); - tmTelephoneCorpusService.sendMq(tag, JSONObject.toJSONString(ltoMap)); - // 保存报告 - oldAiAnalysisRequestLogs.setBusinessResponse(JSONObject.toJSONString(ltoMap)); - + oldAiAnalysisRequestLogs.setDifyResponse(JSON.toJSONString(execDifyFlow)); + updateDiFyRequest(oldAiAnalysisRequestLogs, oldAiAnalysisRequestLogs.getAiAnalysisRequestId()); + aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1); + aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("1"); + updateAiAnalysisErrors(aiAnalysisErrors, oldAiAnalysisRequestLogs.getAiAnalysisRequestId()); } - oldAiAnalysisRequestLogs.setDifyResponse(JSON.toJSONString(execDifyFlow)); - updateDiFyRequest(oldAiAnalysisRequestLogs,oldAiAnalysisRequestLogs.getAiAnalysisRequestId()); + } catch (Exception e) { + log.info("语料解析失败补偿异常:{}", e); + aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("0"); aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1); - aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("1"); - updateAiAnalysisErrors(aiAnalysisErrors,oldAiAnalysisRequestLogs.getAiAnalysisRequestId()); + updateAiAnalysisErrors(aiAnalysisErrors, aiAnalysisErrors.getAiAnalysisRequestId()); } - } catch (Exception e) { - log.info("语料解析失败补偿异常:{}",e); - aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("0"); - aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount()+ 1); - updateAiAnalysisErrors(aiAnalysisErrors,aiAnalysisErrors.getAiAnalysisRequestId()); - } - }); + }); + } + } + - } } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisErrorsMapper.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisErrorsMapper.java index d196aee..c4c684f 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisErrorsMapper.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisErrorsMapper.java @@ -10,7 +10,7 @@ import java.util.List; @Mapper public interface AiAnalysisErrorsMapper extends BaseMapper { - List queryAnalysisErrorList(@Param("businessType") String businessType); - + int queryCountAnalysisErrorList(@Param("businessType") String businessType); + List queryAnalysisErrorList( @Param("businessType") String businessType, @Param("offset") int offset, @Param("pageSize") int pageSize); } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java index 4e14e08..6e7d896 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java @@ -17,6 +17,8 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.cloud.context.config.annotation.RefreshScope; import org.springframework.kafka.annotation.KafkaListener; +import org.springframework.kafka.annotation.PartitionOffset; +import org.springframework.kafka.annotation.TopicPartition; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Component; import org.springframework.web.bind.annotation.PostMapping; @@ -54,7 +56,27 @@ public class CorpusProcessKafkaProducer { private RocketMQTemplate rocketMqTemplate; @PostMapping("corpusProcessKafkaConsumer") - @KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") + @KafkaListener( + topicPartitions = @TopicPartition( + topic = "${spring.kafka.topic}", + partitions = {"0", "1","2", "3","4", "5","6", "7","8", "9", "10", "11"}, + partitionOffsets = { + @PartitionOffset(partition = "0", initialOffset = "2792520"), + @PartitionOffset(partition = "1", initialOffset = "2596153"), + @PartitionOffset(partition = "2", initialOffset = "2536889"), + @PartitionOffset(partition = "3", initialOffset = "2782173"), + @PartitionOffset(partition = "4", initialOffset = "2616677"), + @PartitionOffset(partition = "5", initialOffset = "2529585"), + @PartitionOffset(partition = "6", initialOffset = "2761950"), + @PartitionOffset(partition = "7", initialOffset = "2581957"), + @PartitionOffset(partition = "8", initialOffset = "2530225"), + @PartitionOffset(partition = "9", initialOffset = "2803179"), + @PartitionOffset(partition = "10", initialOffset = "2599129"), + @PartitionOffset(partition = "11", initialOffset = "2546277") + } + ), + groupId = "${spring.kafka.group}" + ) public void listen(List recordMessages) { long startTime = System.currentTimeMillis(); try { @@ -97,11 +119,6 @@ public class CorpusProcessKafkaProducer { } }, 10000); -// CompletableFuture.runAsync(() -> { -// tmTelephoneCorpusService.runTelephoneCorpusDify(aicorpusTelephone); -// }, runTelephoneExecutor); -// // 关闭线程池 -// runTelephoneExecutor.shutdown(); log.info(" dify处理耗时:{}", System.currentTimeMillis() - startTimeDify); } } catch (Exception e) { diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiAnalysisErrorsService.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiAnalysisErrorsService.java index 862d0dd..fdb949d 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiAnalysisErrorsService.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiAnalysisErrorsService.java @@ -2,12 +2,13 @@ package com.volvo.ai.analytic.center.service; import com.baomidou.mybatisplus.extension.service.IService; import com.volvo.ai.analytic.center.entity.AiAnalysisErrors; +import org.apache.ibatis.annotations.Param; import java.util.List; public interface AiAnalysisErrorsService extends IService { boolean saveAiAnalysisErrors(AiAnalysisErrors entity); - - List queryAnalysisErrorList(String businessType); + int queryCountAnalysisErrorList(String businessType); + List queryAnalysisErrorList( String businessType, int offset,int pageSize); } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisErrorsServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisErrorsServiceImpl.java index 3e9265f..b109b06 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisErrorsServiceImpl.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisErrorsServiceImpl.java @@ -30,9 +30,14 @@ public class AiAnalysisErrorsServiceImpl extends ServiceImpl queryAnalysisErrorList(String businessType) { + @Override + public int queryCountAnalysisErrorList(String businessType) { + return aiAnalysisErrorsMapper.queryCountAnalysisErrorList(businessType); + } + + public List queryAnalysisErrorList( String businessType, int offset,int pageSize) { //捞取异常表中属于社区的异常数据 - return aiAnalysisErrorsMapper.queryAnalysisErrorList(businessType); + return aiAnalysisErrorsMapper.queryAnalysisErrorList(businessType, offset, pageSize); } } diff --git a/ai-analytic-center-biz/src/main/resources/mapper/AiAnalysisErrorsMapper.xml b/ai-analytic-center-biz/src/main/resources/mapper/AiAnalysisErrorsMapper.xml index ae5f953..eaca7d6 100644 --- a/ai-analytic-center-biz/src/main/resources/mapper/AiAnalysisErrorsMapper.xml +++ b/ai-analytic-center-biz/src/main/resources/mapper/AiAnalysisErrorsMapper.xml @@ -1,6 +1,23 @@ + \ No newline at end of file