|
|
|
|
@@ -8,15 +8,16 @@ import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
|
|
|
|
|
import com.volvo.ai.analytic.center.constant.Constant;
|
|
|
|
|
import com.volvo.ai.analytic.center.dto.req.*;
|
|
|
|
|
import com.volvo.ai.analytic.center.dto.resp.*;
|
|
|
|
|
import com.volvo.ai.analytic.center.entity.DataMaskingRule;
|
|
|
|
|
import com.volvo.ai.analytic.center.entity.DiffdefeatApprove;
|
|
|
|
|
import com.volvo.ai.analytic.center.entity.MqMessageRecord;
|
|
|
|
|
import com.volvo.ai.analytic.center.entity.*;
|
|
|
|
|
import com.volvo.ai.analytic.center.enums.*;
|
|
|
|
|
import com.volvo.ai.analytic.center.mapper.AiAnalysisErrorsMapper;
|
|
|
|
|
import com.volvo.ai.analytic.center.mapper.AiAnalysisRequestLogsMapper;
|
|
|
|
|
import com.volvo.ai.analytic.center.mapper.MqMessageRecordMapper;
|
|
|
|
|
import com.volvo.ai.analytic.center.service.DataMaskingRuleService;
|
|
|
|
|
import com.volvo.ai.analytic.center.service.DiFyService;
|
|
|
|
|
import com.volvo.ai.analytic.center.service.DiffdefeatApproveService;
|
|
|
|
|
import com.volvo.ai.analytic.center.service.MqMessageRecordService;
|
|
|
|
|
import com.volvo.ai.analytic.center.utils.AiAnalysisUtils;
|
|
|
|
|
import lombok.extern.slf4j.Slf4j;
|
|
|
|
|
//import org.springframework.amqp.core.AmqpTemplate;
|
|
|
|
|
import org.apache.rocketmq.spring.core.RocketMQTemplate;
|
|
|
|
|
@@ -37,9 +38,6 @@ import java.util.stream.Collectors;
|
|
|
|
|
@Service
|
|
|
|
|
public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMapper, MqMessageRecord> implements MqMessageRecordService {
|
|
|
|
|
|
|
|
|
|
// @Autowired
|
|
|
|
|
// private AmqpTemplate rabbitTemplate;
|
|
|
|
|
|
|
|
|
|
@Autowired
|
|
|
|
|
private DataMaskingRuleService dataMaskingRuleService;
|
|
|
|
|
|
|
|
|
|
@@ -55,6 +53,12 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
|
|
|
|
@Autowired
|
|
|
|
|
private RocketMQTemplate rocketMQTemplate;
|
|
|
|
|
|
|
|
|
|
@Autowired
|
|
|
|
|
private AiAnalysisRequestLogsMapper aiAnalysisRequestLogsMapper;
|
|
|
|
|
|
|
|
|
|
@Autowired
|
|
|
|
|
private AiAnalysisErrorsMapper aiAnalysisErrorsMapper;
|
|
|
|
|
|
|
|
|
|
@Value("${service.dify.user}")
|
|
|
|
|
private String user;
|
|
|
|
|
|
|
|
|
|
@@ -66,38 +70,48 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* 处理Mq消息
|
|
|
|
|
*
|
|
|
|
|
* @param message 消息具体内容
|
|
|
|
|
* @return 返回处理情况
|
|
|
|
|
*/
|
|
|
|
|
@Override
|
|
|
|
|
@Transactional
|
|
|
|
|
public boolean processMessageByMQ(String message) {
|
|
|
|
|
log.info("processMessageByMQ message: {}", message);
|
|
|
|
|
LocalDateTime currTime = LocalDateTime.now();
|
|
|
|
|
MqFormData oldItem = JSONObject.parseObject(message, MqFormData.class);
|
|
|
|
|
log.info("processMessageByMQ 真假战败数据已获取 {}",oldItem.getFormId());
|
|
|
|
|
MqMessageRecord curMQMessageRecord = new MqMessageRecord();
|
|
|
|
|
try {
|
|
|
|
|
curMQMessageRecord.setSinceType(oldItem.getSinceType());
|
|
|
|
|
curMQMessageRecord.setSubSinceType(oldItem.getSubSinceType());
|
|
|
|
|
curMQMessageRecord.setBizNo(oldItem.getFormId());
|
|
|
|
|
curMQMessageRecord.setSubBizNo(oldItem.getFormId());
|
|
|
|
|
curMQMessageRecord.setMessageContent(message);
|
|
|
|
|
curMQMessageRecord.setRetryCount(0);
|
|
|
|
|
curMQMessageRecord.setLastRetryTime(currTime);
|
|
|
|
|
curMQMessageRecord.setTaskStatus(0);
|
|
|
|
|
this.save(curMQMessageRecord);
|
|
|
|
|
log.info("processMessageByMQ 数据入库成功");
|
|
|
|
|
} catch (Exception ex) {
|
|
|
|
|
log.error("processMessageByMQ 数据入库失败", ex);
|
|
|
|
|
log.info("communityProcessMessageByMQ message: {}", message);
|
|
|
|
|
CommunityTargetDTO communityTargetDTO = JSON.parseObject(message, CommunityTargetDTO.class);
|
|
|
|
|
if (communityTargetDTO == null) {
|
|
|
|
|
log.error("communityProcessMessageByMQ message is null");
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
if (Objects.equals(SinceTypeEnum.SINCETYPE2.getCode(),oldItem.getSinceType()) && Objects.equals(SubSinceTypeEnum.SINCETYPE51.getCode(),oldItem.getSubSinceType())) {
|
|
|
|
|
this.processMqSinceType51(oldItem, curMQMessageRecord);
|
|
|
|
|
} else if (Objects.equals(SinceTypeEnum.SINCETYPE2.getCode(),oldItem.getSinceType()) && Objects.equals(SubSinceTypeEnum.SINCETYPE52.getCode(),oldItem.getSubSinceType())) {
|
|
|
|
|
this.processMqSinceType52(oldItem, curMQMessageRecord);
|
|
|
|
|
} else {
|
|
|
|
|
log.info("processMessageByMQ 辨别数据为未知类型sinceType {},subSinceType {}", oldItem.getSinceType(), oldItem.getSubSinceType());
|
|
|
|
|
// 生成ai分析请求id
|
|
|
|
|
String aiAnalysisRequestId = AiAnalysisUtils.getAiAnalysisRequestId(BusinessTypeEnum.COMMUNITYTARGET.getCode());
|
|
|
|
|
|
|
|
|
|
// 保存请求日志
|
|
|
|
|
aiAnalysisRequestLogsMapper.insert(AiAnalysisRequestLogs.builder()
|
|
|
|
|
.aiAnalysisRequestId(aiAnalysisRequestId)
|
|
|
|
|
.businessRequest(message)
|
|
|
|
|
.difyAgentKey(flowId)
|
|
|
|
|
.difyRequest(JSON.toJSONString(communityTargetDTO))
|
|
|
|
|
.aiAnalysisRequestType(BusinessTypeEnum.COMMUNITYTARGET.getCode())
|
|
|
|
|
.build());
|
|
|
|
|
|
|
|
|
|
//构建请求dify对象
|
|
|
|
|
DiFyReq diFyReq = new DiFyReq();
|
|
|
|
|
diFyReq.setUser(user);
|
|
|
|
|
diFyReq.setFlowId(flowId);
|
|
|
|
|
diFyReq.setInputs(communityTargetDTO);
|
|
|
|
|
JSONObject difResult = new JSONObject();
|
|
|
|
|
try {
|
|
|
|
|
//调用dify服务
|
|
|
|
|
difResult = (JSONObject) diFyService.getDiFyObject(diFyReq);
|
|
|
|
|
} catch (Exception e) {
|
|
|
|
|
log.error("请求DiFy异常:{}", e);
|
|
|
|
|
//保存错误日志
|
|
|
|
|
aiAnalysisErrorsMapper.insert(AiAnalysisErrors.builder()
|
|
|
|
|
.aiAnalysisRequestId(aiAnalysisRequestId)
|
|
|
|
|
.difyResponse(difResult.toJSONString())
|
|
|
|
|
.aiAnalysisErrorMessage(e.getMessage())
|
|
|
|
|
.build());
|
|
|
|
|
}
|
|
|
|
|
return true;
|
|
|
|
|
}
|
|
|
|
|
@@ -146,7 +160,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
|
|
|
|
|
|
|
|
|
private void processMqSinceType52(MqFormData oldItem, MqMessageRecord curMQMessageRecord) {
|
|
|
|
|
log.info("processMqSinceType52 辨别数据为分析请求 {}", SubSinceTypeEnum.SINCETYPE52.getCode());
|
|
|
|
|
com.volvo.ai.analytic.center.dto.req.DiffDefeatanApprove curDiffDefeatApprove = JSON.parseObject(JSON.toJSONString(oldItem.getData()), com.volvo.ai.analytic.center.dto.req.DiffDefeatanApprove.class);
|
|
|
|
|
DiffDefeatanApprove curDiffDefeatApprove = JSON.parseObject(JSON.toJSONString(oldItem.getData()), DiffDefeatanApprove.class);
|
|
|
|
|
DiffdefeatApprove approveEntity = diffdefeatApproveService.lambdaQuery().eq(DiffdefeatApprove::getFormId, oldItem.getFormId()).last("limit 1").one();
|
|
|
|
|
if (approveEntity != null) {
|
|
|
|
|
try {
|
|
|
|
|
@@ -416,6 +430,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private String getUserStatus(String oldStr) {
|
|
|
|
|
if (Objects.equals(MessageConvertEnum.CONFIRMED.getCode(), oldStr)) {
|
|
|
|
|
return MessageConvertEnum.CONFIRMED.getMessage();
|
|
|
|
|
|