去掉校验限制,增加error状态更新
This commit is contained in:
@@ -4,8 +4,6 @@ package com.volvo.ai.analytic.center.mq;
|
|||||||
import com.fasterxml.jackson.core.JsonProcessingException;
|
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||||
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
|
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
|
||||||
import com.volvo.ai.analytic.center.dto.corpus.DisplayDTO;
|
|
||||||
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
|
|
||||||
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
|
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
|
||||||
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
|
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
@@ -57,16 +55,8 @@ public class CorpusDccMqConsumer implements RocketMQListener<MessageExt> {
|
|||||||
String message = new String(messageExt.getBody());
|
String message = new String(messageExt.getBody());
|
||||||
log.info("dcc_mq message: " + message);
|
log.info("dcc_mq message: " + message);
|
||||||
AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class);
|
AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class);
|
||||||
AiAnalysisRequestLogs aiAnalysisRequestLogs = null;
|
|
||||||
if(checkDccRepeat.equals("true")){
|
|
||||||
aiAnalysisRequestLogs = aiAnalysisRequestLogsService.queryAiAnalysisRequestLogsByBusinessReponse(aicorpusTelephone.getSourceId());
|
|
||||||
log.info("corpusDccMqProducer 消费校验是否已经发送: {}", aiAnalysisRequestLogs);
|
|
||||||
}
|
|
||||||
|
|
||||||
if(null == aiAnalysisRequestLogs){
|
|
||||||
tmTelephoneCorpusService.runTelephoneCorpusDify(aicorpusTelephone);
|
tmTelephoneCorpusService.runTelephoneCorpusDify(aicorpusTelephone);
|
||||||
log.info("corpusDccMqProducer mq 处理完成: {}", aicorpusTelephone.getSourceId());
|
log.info("corpusDccMqProducer mq 处理完成: {}", aicorpusTelephone.getSourceId());
|
||||||
}
|
|
||||||
log.info("dcc_mq 处理完成,耗时:{}", System.currentTimeMillis() - startTime);
|
log.info("dcc_mq 处理完成,耗时:{}", System.currentTimeMillis() - startTime);
|
||||||
} catch (JsonProcessingException e) {
|
} catch (JsonProcessingException e) {
|
||||||
log.info(" dcc mq 处理失败:{}", e.getMessage());
|
log.info(" dcc mq 处理失败:{}", e.getMessage());
|
||||||
|
|||||||
@@ -24,7 +24,6 @@ import org.springframework.web.bind.annotation.RestController;
|
|||||||
|
|
||||||
import javax.annotation.Resource;
|
import javax.annotation.Resource;
|
||||||
import java.time.LocalDateTime;
|
import java.time.LocalDateTime;
|
||||||
import java.time.format.DateTimeFormatter;
|
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.concurrent.ExecutorService;
|
import java.util.concurrent.ExecutorService;
|
||||||
import java.util.concurrent.Executors;
|
import java.util.concurrent.Executors;
|
||||||
@@ -51,9 +50,6 @@ public class CorpusProcessKafkaProducer {
|
|||||||
@Value("${rocketmq.producer.corpus.dcctopic}")
|
@Value("${rocketmq.producer.corpus.dcctopic}")
|
||||||
private String dccMqTipic;
|
private String dccMqTipic;
|
||||||
|
|
||||||
@Value("${dify.corpus.kafkaTimeLimit}")
|
|
||||||
private String kafkaTimeLimit;
|
|
||||||
|
|
||||||
@Resource
|
@Resource
|
||||||
private RocketMQTemplate rocketMqTemplate;
|
private RocketMQTemplate rocketMqTemplate;
|
||||||
|
|
||||||
@@ -94,15 +90,6 @@ public class CorpusProcessKafkaProducer {
|
|||||||
try {
|
try {
|
||||||
String message = (String) record.value();
|
String message = (String) record.value();
|
||||||
AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class);
|
AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class);
|
||||||
String transcribeTimeStr = aicorpusTelephone.getTranscribeTime();
|
|
||||||
if (transcribeTimeStr != null) {
|
|
||||||
DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
|
|
||||||
LocalDateTime transcribeTime = LocalDateTime.parse(transcribeTimeStr, formatter);
|
|
||||||
LocalDateTime kafkaStartTime = LocalDateTime.parse(kafkaTimeLimit, formatter);
|
|
||||||
if (transcribeTime.isBefore(kafkaStartTime)) {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
log.info("aicorpusTelephone categoryCode:{}, display: {}", aicorpusTelephone.getCategoryCode(), aicorpusTelephone.getDisplay());
|
log.info("aicorpusTelephone categoryCode:{}, display: {}", aicorpusTelephone.getCategoryCode(), aicorpusTelephone.getDisplay());
|
||||||
DisplayDTO display = objectMapper.readValue(aicorpusTelephone.getDisplay(), DisplayDTO.class);
|
DisplayDTO display = objectMapper.readValue(aicorpusTelephone.getDisplay(), DisplayDTO.class);
|
||||||
log.info("aicorpusTelephone display getSegments: {}", display.getSegments());
|
log.info("aicorpusTelephone display getSegments: {}", display.getSegments());
|
||||||
|
|||||||
@@ -1,7 +1,9 @@
|
|||||||
package com.volvo.ai.analytic.center.service.impl;
|
package com.volvo.ai.analytic.center.service.impl;
|
||||||
|
|
||||||
|
import com.alibaba.fastjson.JSON;
|
||||||
import com.alibaba.fastjson.JSONArray;
|
import com.alibaba.fastjson.JSONArray;
|
||||||
import com.alibaba.fastjson.JSONObject;
|
import com.alibaba.fastjson.JSONObject;
|
||||||
|
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
|
||||||
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
|
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
|
||||||
import com.volvo.ai.analytic.center.constant.Constant;
|
import com.volvo.ai.analytic.center.constant.Constant;
|
||||||
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
|
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
|
||||||
@@ -11,6 +13,7 @@ import com.volvo.ai.analytic.center.dto.req.DiFyReq;
|
|||||||
import com.volvo.ai.analytic.center.dto.req.RunMaskingRuleInput;
|
import com.volvo.ai.analytic.center.dto.req.RunMaskingRuleInput;
|
||||||
import com.volvo.ai.analytic.center.dto.resp.CarModelRespDTO;
|
import com.volvo.ai.analytic.center.dto.resp.CarModelRespDTO;
|
||||||
import com.volvo.ai.analytic.center.dto.resp.ResultDTO;
|
import com.volvo.ai.analytic.center.dto.resp.ResultDTO;
|
||||||
|
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
|
||||||
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
|
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
|
||||||
import com.volvo.ai.analytic.center.entity.DataMaskingRule;
|
import com.volvo.ai.analytic.center.entity.DataMaskingRule;
|
||||||
import com.volvo.ai.analytic.center.entity.TmTelephoneCorpus;
|
import com.volvo.ai.analytic.center.entity.TmTelephoneCorpus;
|
||||||
@@ -19,10 +22,7 @@ import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
|
|||||||
import com.volvo.ai.analytic.center.enums.CategoryEnum;
|
import com.volvo.ai.analytic.center.enums.CategoryEnum;
|
||||||
import com.volvo.ai.analytic.center.feign.RemoteCarModelClient;
|
import com.volvo.ai.analytic.center.feign.RemoteCarModelClient;
|
||||||
import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper;
|
import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper;
|
||||||
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
|
import com.volvo.ai.analytic.center.service.*;
|
||||||
import com.volvo.ai.analytic.center.service.DataMaskingRuleService;
|
|
||||||
import com.volvo.ai.analytic.center.service.DiFyService;
|
|
||||||
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
|
|
||||||
import com.volvo.ai.analytic.center.utils.ConstantStr;
|
import com.volvo.ai.analytic.center.utils.ConstantStr;
|
||||||
import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil;
|
import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
@@ -77,6 +77,9 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
|
|||||||
@Autowired
|
@Autowired
|
||||||
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
|
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
private AiAnalysisErrorsService aiAnalysisErrorsService;
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@Transactional
|
@Transactional
|
||||||
public void saveTelephoneCorpus(TmTelephoneCorpus tmTelephoneCorpus) {
|
public void saveTelephoneCorpus(TmTelephoneCorpus tmTelephoneCorpus) {
|
||||||
@@ -165,6 +168,11 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
|
|||||||
log.info("send mq {}",ltoMap);
|
log.info("send mq {}",ltoMap);
|
||||||
sendMq( CategoryEnum.PHONE_VOICE.getCode(), JSONObject.toJSONString(ltoMap));
|
sendMq( CategoryEnum.PHONE_VOICE.getCode(), JSONObject.toJSONString(ltoMap));
|
||||||
try {
|
try {
|
||||||
|
AiAnalysisErrors aiAnalysisErrors = new AiAnalysisErrors();
|
||||||
|
aiAnalysisErrors.setAiAnalysisRequestId(aiAnalysisRequestId);
|
||||||
|
aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1);
|
||||||
|
aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("1");
|
||||||
|
aiAnalysisErrorsService.saveAiAnalysisErrors(aiAnalysisErrors);
|
||||||
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder().aiAnalysisRequestId(execDifyFlow.getString("aiAnalysisRequestId")).businessResponse(JSONObject.toJSONString(ltoMap)).build());
|
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder().aiAnalysisRequestId(execDifyFlow.getString("aiAnalysisRequestId")).businessResponse(JSONObject.toJSONString(ltoMap)).build());
|
||||||
|
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
@@ -215,4 +223,5 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
}
|
}
|
||||||
Reference in New Issue
Block a user