硬编码代码抽到枚举中

This commit is contained in:
zhangfan
2024-12-20 18:25:52 +08:00
parent 046a5c626e
commit 244e6c6263
12 changed files with 240 additions and 126 deletions

View File

@@ -4,5 +4,7 @@ public class Constant {
public static final String rabbitMqFormQueue = "Voc-Defeat-Dcc";
public static final String rabbitToFormQueue = "Voc-DefeatResponse-Dcc";
public static final String CHANNEL_DCC = "Channel_Dcc";
}

View File

@@ -8,7 +8,7 @@ import lombok.Data;
@Data
public class RabbitMqFormData {
private String formId;
private int sinceType;
private int subSinceType;
private Integer sinceType;
private Integer subSinceType;
private Object data;
}

View File

@@ -7,6 +7,7 @@ import lombok.Data;
import lombok.EqualsAndHashCode;
import lombok.experimental.Accessors;
import java.time.LocalDateTime;
import java.util.Date;
/**
* MQ接收数据表实体类
@@ -14,8 +15,7 @@ import java.util.Date;
@Data
@Accessors
@EqualsAndHashCode(callSuper = true)
//@TableName("tt_mq_message_record")
@TableName("tm_custom")
@TableName("tt_mq_message_record")
public class MqMessageRecord extends BaseEntity {
/**
@@ -27,85 +27,73 @@ public class MqMessageRecord extends BaseEntity {
/**
* 场景类型 1质检分析
*/
// @TableField("since_type")
@TableField(select = false)
@TableField("since_type")
private Integer sinceType;
/**
* 1,新车销售2二手车销售3收车销售4 首定保5延保
*/
// @TableField("sub_since_type")
@TableField(select = false)
@TableField("sub_since_type")
private Integer subSinceType;
/**
* sessionId
*/
// @TableField("biz_no")
@TableField(select = false)
@TableField("biz_no")
private String bizNo;
/**
* sourceId
*/
// @TableField("sub_biz_no")
@TableField(select = false)
@TableField("sub_biz_no")
private String subBizNo;
/**
* 消息体小于4000直接存消息体大于存Iobs key
*/
// @TableField("message_content")
@TableField(select = false)
@TableField("message_content")
private String messageContent;
/**
* mq消息体Iobs key
*/
// @TableField("message_iobs_key")
@TableField(select = false)
@TableField("message_iobs_key")
private String messageIobsKey;
/**
* mq message key
*/
// @TableField("message_key")
@TableField(select = false)
@TableField("message_key")
private String messageKey;
/**
* 批处理重试次数,默认3次
*/
// @TableField("retry_count")
@TableField(select = false)
@TableField("retry_count")
private Integer retryCount;
/**
* 最后一次重试时间
*/
// @TableField("last_retry_time")
@TableField(select = false)
private Date lastRetryTime;
@TableField("last_retry_time")
private LocalDateTime lastRetryTime;
/**
* 响应码
*/
// @TableField("resp_code")
@TableField(select = false)
@TableField("resp_code")
private String respCode;
/**
* 错误消息 截取200长度字符
*/
// @TableField("resp_content")
@TableField(select = false)
@TableField("resp_content")
private String respContent;
/**
* 任务状态: 0:待分析,1.分析中, 2.分析完成, 3.分析失败
*/
// @TableField("task_status")
@TableField(select = false)
@TableField("task_status")
private Integer taskStatus;
/**

View File

@@ -0,0 +1,24 @@
package com.volvo.ai.analytic.center.enums;
public enum CategoryEnum {
ENTERPRISE_WECHAT("企微", "企微记录"),
PHONE_VOICE("语音", "语音记录"),
OTHER("未知", "未知记录"),
;
private String code;
private String message;
CategoryEnum(String code, String message) {
this.code = code;
this.message = message;
}
public String getCode() {
return this.code;
}
public String getMessage() {
return this.message;
}
}

View File

@@ -0,0 +1,24 @@
package com.volvo.ai.analytic.center.enums;
public enum HandleStatusEnum {
ANALYSIS_NORMAL(200, "分析完成"),
ANALYSIS_CONTENT_EMPTY(400, "语料信息为空"),
ANALYSIS_CALLING(500, "存在未完成解析的通话信息"),
;
private Integer code;
private String message;
HandleStatusEnum(Integer code, String message) {
this.code = code;
this.message = message;
}
public Integer getCode() {
return this.code;
}
public String getMessage() {
return this.message;
}
}

View File

@@ -0,0 +1,25 @@
package com.volvo.ai.analytic.center.enums;
public enum MessageConvertEnum {
CONFIRMED("确认到店", "客户确认将到店"),
LOOK_ON("继续观望", "客户继续观望中"),
GIVE_UP("放弃购车", "客户放弃购车"),
INVALID("无效通话", "近三天无有效通话记录"),
;
private String code;
private String message;
MessageConvertEnum(String code, String message) {
this.code = code;
this.message = message;
}
public String getCode() {
return this.code;
}
public String getMessage() {
return this.message;
}
}

View File

@@ -1,20 +1,20 @@
package com.volvo.ai.analytic.center.enums;
public enum MqTaskStatusEnum {
Analysis_Wait(0, "待分析"),
Analysis_Underway(1, "分析中"),
Analysis_Completion(2, "分析完成"),
Analysis_Failure(3, "分析失败"),
ANALYSIS_WAIT(0, "待分析"),
ANALYSIS_UNDERWAY(1, "分析中"),
ANALYSIS_COMPLETION(2, "分析完成"),
ANALYSIS_FAILURE(3, "分析失败"),
;
private int code;
private Integer code;
private String name;
MqTaskStatusEnum(int code, String name) {
MqTaskStatusEnum(Integer code, String name) {
this.code = code;
this.name = name;
}
public int getCode() {
public Integer getCode() {
return this.code;
}
}

View File

@@ -0,0 +1,24 @@
package com.volvo.ai.analytic.center.enums;
public enum RoleEnum {
AGENT("AGENT", "客服"),
USER("USER", "客户"),
OTHER("OTHER", "未知"),
;
private String code;
private String message;
RoleEnum(String code, String message) {
this.code = code;
this.message = message;
}
public String getCode() {
return this.code;
}
public String getMessage() {
return this.message;
}
}

View File

@@ -0,0 +1,18 @@
package com.volvo.ai.analytic.center.enums;
public enum SinceTypeEnum {
SINCETYPE2(2, ""),
;
private Integer code;
private String name;
SinceTypeEnum(Integer code, String name) {
this.code = code;
this.name = name;
}
public Integer getCode() {
return this.code;
}
}

View File

@@ -2,19 +2,19 @@ package com.volvo.ai.analytic.center.enums;
public enum SubSinceTypeEnum {
SinceType51(51, "真假战败分析请求"),
SinceType52(52, "真假战败审批结果"),
SinceType53(53, "真假战败分析结果"),
SINCETYPE51(51, "真假战败分析请求"),
SINCETYPE52(52, "真假战败审批结果"),
SINCETYPE53(53, "真假战败分析结果"),
;
private int code;
private Integer code;
private String name;
SubSinceTypeEnum(int code, String name) {
SubSinceTypeEnum(Integer code, String name) {
this.code = code;
this.name = name;
}
public int getCode() {
public Integer getCode() {
return this.code;
}
}

View File

@@ -2,18 +2,24 @@ package com.volvo.ai.analytic.center.rabbitMq;
import com.volvo.ai.analytic.center.constant.Constant;
import com.volvo.ai.analytic.center.service.MqMessageRecordService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
@Slf4j
@Component
public class MessageListener {
@Autowired
private MqMessageRecordService mqMessageRecordService;
//默认情况下,当使用@RabbitListener注解时消息确认模式通常是自动的AcknowledgeMode.AUTO可以在yaml文件中更改
// 消息一旦被消费者接收并处理完成即方法执行完成就会自动发送ack确认给RabbitMQ。
@RabbitListener(queues = Constant.rabbitMqFormQueue)
public void onMessage(String message) {
log.info("Received message: " + message);
boolean result = mqMessageRecordService.processMessageByMQ(message);
}
}

View File

@@ -5,23 +5,21 @@ import com.alibaba.cloud.commons.lang.StringUtils;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.fasterxml.jackson.databind.ObjectMapper;
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.enums.MqTaskStatusEnum;
import com.volvo.ai.analytic.center.enums.SubSinceTypeEnum;
import com.volvo.ai.analytic.center.enums.*;
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 lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.AmqpTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Service;
@@ -37,6 +35,9 @@ import java.util.stream.Collectors;
@Service
public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMapper, MqMessageRecord> implements MqMessageRecordService {
@Autowired
private AmqpTemplate rabbitTemplate;
@Autowired
private DataMaskingRuleService dataMaskingRuleService;
@@ -63,6 +64,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
@Override
public boolean processMessageByMQ(String message) {
log.info("message: {}", message);
LocalDateTime currTime = LocalDateTime.now();
RabbitMqFormData oldItem = JSONObject.parseObject(message, RabbitMqFormData.class);
log.info("真假战败数据已获取 {}",oldItem.getFormId());
MqMessageRecord curMQMessageRecord = new MqMessageRecord();
@@ -73,7 +75,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
curMQMessageRecord.setSubBizNo(oldItem.getFormId());
curMQMessageRecord.setMessageContent(message);
curMQMessageRecord.setRetryCount(0);
curMQMessageRecord.setLastRetryTime(new Date());
curMQMessageRecord.setLastRetryTime(currTime);
curMQMessageRecord.setTaskStatus(0);
this.save(curMQMessageRecord);
log.info("数据入库成功(MQMessageRecord)");
@@ -81,30 +83,30 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
log.error("数据入库失败(MQMessageRecord)", ex);
return false;
}
if (oldItem.getSinceType() == 2 && oldItem.getSubSinceType() == SubSinceTypeEnum.SinceType51.getCode() ) {
if (SinceTypeEnum.SINCETYPE2.getCode().equals(oldItem.getSinceType()) && SubSinceTypeEnum.SINCETYPE51.getCode().equals(oldItem.getSubSinceType()) ) {
log.info("辨别数据为分析请求");
DiffDefeatResult diffDefeatResult = null;
DiffDefeatanAlysis diffDefeatanAlysis = JSON.parseObject(JSON.toJSONString(oldItem.getData()), DiffDefeatanAlysis.class);
DiffDefeatAnalyseOutputResult response = this.processChatRecord(diffDefeatanAlysis);
DiffDefeatAnalyseOutput output = response.getDiffDefeatAnalyseOutput();
log.info("数据分析完成");
if(output.getHandleStatus() == 200){
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.Analysis_Completion.getCode());
if(HandleStatusEnum.ANALYSIS_NORMAL.getCode().equals(output.getHandleStatus())){
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.ANALYSIS_COMPLETION.getCode());
diffDefeatResult = JSONObject.parseObject(output.getResultStr(), DiffDefeatResult.class);
} else if (output.getHandleStatus() == 400) {
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.Analysis_Completion.getCode());
} else if (HandleStatusEnum.ANALYSIS_CONTENT_EMPTY.getCode().equals(output.getHandleStatus())) {
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.ANALYSIS_COMPLETION.getCode());
diffDefeatResult = new DiffDefeatResult();
diffDefeatResult.setUserStatus("无效通话");
diffDefeatResult.setAppointmentResult("因为通话时间太短没有ASR转义文本因此判断为无效通话");
} else {
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.Analysis_Failure.getCode());
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.ANALYSIS_FAILURE.getCode());
}
curMQMessageRecord.setRespCode(output.getHandleStatus()+"");
curMQMessageRecord.setRespCode(output.getHandleStatus() == null ? "" : output.getHandleStatus().toString());
curMQMessageRecord.setRespContent(output.getResultStr());
this.updateById(curMQMessageRecord);
log.info("数据更改分析状态完成");
if (curMQMessageRecord.getTaskStatus() == MqTaskStatusEnum.Analysis_Completion.getCode()){
if (MqTaskStatusEnum.ANALYSIS_COMPLETION.getCode().equals(curMQMessageRecord.getTaskStatus())){
log.info("数据分析结果正常");
DiffdefeatApprove curDiffdefeatApprove = new DiffdefeatApprove();
String userStatus = getUserStatus(diffDefeatResult.getUserStatus());
@@ -119,23 +121,23 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
DiffDefeatanCallbak curDiffDefeatanCallbak = new DiffDefeatanCallbak();
curDiffDefeatanCallbak.setBusinessId(diffDefeatanAlysis.getBusinessId());
curDiffDefeatanCallbak.setResponseCode("200");
curDiffDefeatanCallbak.setResponseCode(HandleStatusEnum.ANALYSIS_NORMAL.getCode().toString());
curDiffDefeatanCallbak.setLabel(userStatus);
curDiffDefeatanCallbak.setDescribe(diffDefeatResult.getAppointmentResult());
RabbitMqToData rabbitMqToData = new RabbitMqToData();
rabbitMqToData.setFormId(oldItem.getFormId());
rabbitMqToData.setSinceType(2);
rabbitMqToData.setSubSinceType(SubSinceTypeEnum.SinceType53.getCode());
rabbitMqToData.setSinceType(SinceTypeEnum.SINCETYPE2.getCode());
rabbitMqToData.setSubSinceType(SubSinceTypeEnum.SINCETYPE53.getCode());
rabbitMqToData.setData(curDiffDefeatanCallbak);
String callbakInput = JSONObject.toJSONString(rabbitMqToData);
//todo 发Mq给LTO
rabbitTemplate.convertAndSend(Constant.rabbitToFormQueue, message);
log.info("发送回调MQ完成: {}", callbakInput);
} else {
log.info("数据分析结果不正常,等待重试");
return false;
}
}
if (oldItem.getSinceType() == 2 && oldItem.getSubSinceType() == SubSinceTypeEnum.SinceType52.getCode()) {
if (SinceTypeEnum.SINCETYPE2.getCode().equals(oldItem.getSinceType()) && SubSinceTypeEnum.SINCETYPE52.getCode().equals(oldItem.getSubSinceType())) {
log.info("辨别数据为审批结果");
com.volvo.ai.analytic.center.dto.req.DiffDefeatanApprove curDiffDefeatanApprove = JSON.parseObject(JSON.toJSONString(oldItem.getData()), com.volvo.ai.analytic.center.dto.req.DiffDefeatanApprove.class);
DiffdefeatApprove approveEntity = diffdefeatApproveService.lambdaQuery().eq(DiffdefeatApprove::getFormId, oldItem.getFormId()).last("limit 1").one();
@@ -146,15 +148,15 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
approveEntity.setApproveOpinion(curDiffDefeatanApprove.getApproveOpinion());
diffdefeatApproveService.updateById(approveEntity);
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.Analysis_Completion.getCode());
curMQMessageRecord.setRespCode("200");
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.ANALYSIS_COMPLETION.getCode());
curMQMessageRecord.setRespCode(HandleStatusEnum.ANALYSIS_NORMAL.getCode().toString());
this.updateById(curMQMessageRecord);
log.info("审批结果更新完成");
} catch (Exception e) {
log.error("审批结果更新失败",e.getMessage());
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.Analysis_Failure.getCode());
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.ANALYSIS_FAILURE.getCode());
curMQMessageRecord.setRespContent(e.getMessage());
curMQMessageRecord.setRespCode("500");
curMQMessageRecord.setRespCode(HandleStatusEnum.ANALYSIS_CALLING.getCode().toString());
this.updateById(curMQMessageRecord);
}
}
@@ -169,7 +171,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
@Override
public DiffDefeatAnalyseOutputResult processChatRecord(DiffDefeatanAlysis input) {
DiffDefeatAnalyseOutput output = new DiffDefeatAnalyseOutput();
String contentStr = "";
StringBuffer contentStr = new StringBuffer();
//获取脱敏配置信息
List<DataMaskingRule> maskingRuleItems = dataMaskingRuleService.getDataMaskingRuleListByApplicationChannel(Constant.CHANNEL_DCC);
@@ -178,12 +180,12 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
List<String> sourceIds = Optional.ofNullable(input.getCallList()).map(list -> list.stream().map(CallItem::getSourceId).collect(Collectors.toList())).orElse(Collections.emptyList());
if (!CollectionUtils.isEmpty(sourceIds)) {
//查询未完成解析的通话信息
String sourceId = sourceIds.stream().map(code -> "'"+String.valueOf(code)+"'").collect(Collectors.joining(","));
String sourceId = sourceIds.stream().map(code -> "'"+code+"'").collect(Collectors.joining(","));
int asrCount = clickhouseJdbcTemplate.queryForObject("select count(1) from asr_speechdetail where source_id in ("+ sourceId+") and file_status='InProgress'",Integer.class);
if (asrCount > 0) {
output.setHandleStatus(500);
output.setResultStr("存在未完成解析的通话信息");
return new DiffDefeatAnalyseOutputResult(output, contentStr);
output.setHandleStatus(HandleStatusEnum.ANALYSIS_CALLING.getCode());
output.setResultStr(HandleStatusEnum.ANALYSIS_CALLING.getMessage());
return new DiffDefeatAnalyseOutputResult(output, contentStr.toString());
}
//查询到的通话数据
List<Map<String, Object>> hishistoryList = clickhouseJdbcTemplate.queryForList("select id,msg_json,source_id from asr_hishistory where dialect_text=0 and source_id in ("+ sourceId+") order by id");
@@ -194,7 +196,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
.findFirst().map(entity -> entity.get("msg_json") == null ? "" : entity.get("msg_json").toString());
DiffDefeatCorpuItem curDiffDefeatCorpuItem = new DiffDefeatCorpuItem();
curDiffDefeatCorpuItem.setSourceId(item.getSourceId());
curDiffDefeatCorpuItem.setCategory("语音");
curDiffDefeatCorpuItem.setCategory(CategoryEnum.PHONE_VOICE.getCode());
curDiffDefeatCorpuItem.setCorpuText(result.isPresent() ? result.get() : "");
curDiffDefeatCorpuItem.setHappenTime(item.getAudioTime());
curDiffDefeatCorpuItems.add(curDiffDefeatCorpuItem);
@@ -239,47 +241,47 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
this.addSliceData(curDiffDefeatCorpuItems, slice2StopTime, slice3StopTime, userInfoList);
}
} catch (ParseException e) {
e.printStackTrace();
log.error("企微数据处理异常",e);
}
}
// 对集合进行排序
curDiffDefeatCorpuItems.sort(Comparator.comparing(DiffDefeatCorpuItem::getHappenTime));
if (!CollectionUtils.isEmpty(curDiffDefeatCorpuItems)){
for (DiffDefeatCorpuItem item : curDiffDefeatCorpuItems) {
if ("企微".equals(item.getCategory())) {
contentStr += "企微记录\n";
contentStr += item.getCorpuText();
} else if ("语音".equals(item.getCategory()) && StringUtils.isNotEmpty(item.getCorpuText())) {
if (CategoryEnum.ENTERPRISE_WECHAT.getCode().equals(item.getCategory())) {
contentStr.append(CategoryEnum.ENTERPRISE_WECHAT.getMessage()).append("\n");
contentStr.append(item.getCorpuText());
} else if (CategoryEnum.PHONE_VOICE.getCode().equals(item.getCategory()) && StringUtils.isNotEmpty(item.getCorpuText())) {
// 反序列化CorpuText为KafkaJson对象
KafkaJson oldItem = JSON.parseObject(item.getCorpuText(), KafkaJson.class);
// 反序列化display为CollectTranscriberJobResponse对象
CollectTranscriberJobResponse resp = JSON.parseObject(oldItem.getDisplay(), CollectTranscriberJobResponse.class);
// 检查状态并处理Segments
if ("FINISHED".equals(resp.getStatus()) && !CollectionUtils.isEmpty(resp.getSegments())) {
String vocStr = "";
StringBuffer vocStr = new StringBuffer();
for (Segment segment : resp.getSegments()) {
if ("AGENT".equals(segment.getResult().getAnalysisInfo().getRole())) {
vocStr += "客服" + segment.getResult().getText() + "\n";
} else if ("USER".equals(segment.getResult().getAnalysisInfo().getRole())) {
vocStr += "客户" + segment.getResult().getText() + "\n";
if (RoleEnum.AGENT.getCode().equals(segment.getResult().getAnalysisInfo().getRole())) {
vocStr.append(RoleEnum.AGENT.getMessage()+"" + segment.getResult().getText() + "\n");
} else if (RoleEnum.USER.getCode().equals(segment.getResult().getAnalysisInfo().getRole())) {
vocStr.append(RoleEnum.USER.getMessage()+"" + segment.getResult().getText() + "\n");
} else {
vocStr += "未知" + segment.getResult().getText() + "\n";
vocStr.append(RoleEnum.OTHER.getMessage()+"" + segment.getResult().getText() + "\n");
}
}
contentStr += vocStr;
contentStr.append(vocStr);
}
}
}
}
if (StringUtils.isEmpty(contentStr)) {
output.setHandleStatus(400);
output.setResultStr("语料信息为空");
return new DiffDefeatAnalyseOutputResult(output, contentStr);
if (StringUtils.isEmpty(contentStr.toString())) {
output.setHandleStatus(HandleStatusEnum.ANALYSIS_CONTENT_EMPTY.getCode());
output.setResultStr(HandleStatusEnum.ANALYSIS_CONTENT_EMPTY.getMessage());
return new DiffDefeatAnalyseOutputResult(output, contentStr.toString());
}
log.info("开始脱敏 {}", contentStr);
RunMaskingRuleInput runMaskingRuleInput = new RunMaskingRuleInput();
runMaskingRuleInput.setOpinionId(input.getBusinessId());
runMaskingRuleInput.setOldStr(contentStr);
runMaskingRuleInput.setOldStr(contentStr.toString());
runMaskingRuleInput.setDataMaskingRules(maskingRuleItems);
String summaryText = dataMaskingRuleService.runMaskingRule(runMaskingRuleInput);
log.info("脱敏结果 {}", summaryText);
@@ -292,20 +294,20 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
diFyReq.setFlowId(flowId);
diFyReq.setInputs(record);
JSONObject difyResult = (JSONObject) diFyService.getDiFyObject(diFyReq);
output.setHandleStatus(200);
output.setHandleStatus(HandleStatusEnum.ANALYSIS_NORMAL.getCode());
output.setResultStr(difyResult.getString("result"));
log.info("DiFy平台处理结果:{}", output.getResultStr());
return new DiffDefeatAnalyseOutputResult(output, contentStr);
return new DiffDefeatAnalyseOutputResult(output, contentStr.toString());
}
@Override
public void processMessageByTask() {
log.info("开始处理任务");
Date currTime = new Date();
Date startTime = DateUtil.offsetDay(currTime, -3);
Date stopTime = DateUtil.offsetMinute(currTime, -5);
LocalDateTime currTime = LocalDateTime.now();
LocalDateTime startTime = currTime.plusDays(-3);
LocalDateTime stopTime = currTime.plusMinutes(-5);
List<MqMessageRecord> mqMessageRecords = this.lambdaQuery()
.eq(MqMessageRecord::getTaskStatus, MqTaskStatusEnum.Analysis_Failure.getCode())
.eq(MqMessageRecord::getTaskStatus, MqTaskStatusEnum.ANALYSIS_FAILURE.getCode())
.eq(MqMessageRecord::getSinceType, 2)
.le(MqMessageRecord::getRetryCount, 3)
.ge(MqMessageRecord::getCreateTime, startTime)
@@ -317,25 +319,25 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
for (MqMessageRecord curMQMessageRecord : mqMessageRecords){
log.info("辨别数据为分析请求");
RabbitMqFormData oldItem = JSONObject.parseObject(curMQMessageRecord.getMessageContent(), RabbitMqFormData.class);
if (curMQMessageRecord.getSinceType() == 2 && curMQMessageRecord.getSubSinceType() == SubSinceTypeEnum.SinceType51.getCode() ) {
if (curMQMessageRecord.getSinceType() == 2 && curMQMessageRecord.getSubSinceType() == SubSinceTypeEnum.SINCETYPE51.getCode() ) {
try {
DiffDefeatResult diffDefeatResult = null;
DiffDefeatanAlysis diffDefeatanAlysis = JSON.parseObject(JSON.toJSONString(oldItem.getData()), DiffDefeatanAlysis.class);
DiffDefeatAnalyseOutputResult response = this.processChatRecord(diffDefeatanAlysis);
DiffDefeatAnalyseOutput output = response.getDiffDefeatAnalyseOutput();
log.info("数据分析完成");
if(output.getHandleStatus() == 200){
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.Analysis_Completion.getCode());
if(HandleStatusEnum.ANALYSIS_NORMAL.getCode().equals(output.getHandleStatus())){
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.ANALYSIS_COMPLETION.getCode());
diffDefeatResult = JSONObject.parseObject(output.getResultStr(), DiffDefeatResult.class);
} else if (output.getHandleStatus() == 400) {
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.Analysis_Completion.getCode());
} else if (HandleStatusEnum.ANALYSIS_CONTENT_EMPTY.getCode().equals(output.getHandleStatus())) {
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.ANALYSIS_COMPLETION.getCode());
diffDefeatResult = new DiffDefeatResult();
diffDefeatResult.setUserStatus("无效通话");
diffDefeatResult.setAppointmentResult("因为通话时间太短没有ASR转义文本因此判断为无效通话");
} else {
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.Analysis_Failure.getCode());
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.ANALYSIS_FAILURE.getCode());
}
if (curMQMessageRecord.getTaskStatus() == MqTaskStatusEnum.Analysis_Completion.getCode()){
if (curMQMessageRecord.getTaskStatus() == MqTaskStatusEnum.ANALYSIS_COMPLETION.getCode()){
String userStatus = getUserStatus(diffDefeatResult.getUserStatus());
DiffdefeatApprove curDiffdefeatApprove = diffdefeatApproveService.lambdaQuery()
.eq(DiffdefeatApprove::getFormId, oldItem.getFormId())
@@ -361,36 +363,36 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
DiffDefeatanCallbak curDiffDefeatanCallbak = new DiffDefeatanCallbak();
curDiffDefeatanCallbak.setBusinessId(diffDefeatanAlysis.getBusinessId());
curDiffDefeatanCallbak.setResponseCode("200");
curDiffDefeatanCallbak.setResponseCode(HandleStatusEnum.ANALYSIS_NORMAL.getCode().toString());
curDiffDefeatanCallbak.setLabel(userStatus);
curDiffDefeatanCallbak.setDescribe(diffDefeatResult.getAppointmentResult());
RabbitMqToData rabbitMqToData = new RabbitMqToData();
rabbitMqToData.setFormId(oldItem.getFormId());
rabbitMqToData.setSinceType(2);
rabbitMqToData.setSubSinceType(SubSinceTypeEnum.SinceType53.getCode());
rabbitMqToData.setSubSinceType(SubSinceTypeEnum.SINCETYPE53.getCode());
rabbitMqToData.setData(curDiffDefeatanCallbak);
String callbakInput = JSONObject.toJSONString(rabbitMqToData);
//todo 发Mq给LTO
rabbitTemplate.convertAndSend(Constant.rabbitToFormQueue, callbakInput);
log.info("发送回调MQ完成: {}", callbakInput);
}
curMQMessageRecord.setRespCode(output.getHandleStatus()+"");
curMQMessageRecord.setRespContent(output.getResultStr());
curMQMessageRecord.setLastRetryTime(new Date());
curMQMessageRecord.setLastRetryTime(currTime);
curMQMessageRecord.setRetryCount(curMQMessageRecord.getRetryCount()+1);
this.updateById(curMQMessageRecord);
log.info("数据更改分析状态完成");
} catch (Exception e) {
log.error("数据解析异常: {}", e.getMessage());
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.Analysis_Failure.getCode());
curMQMessageRecord.setRespCode("500");
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.ANALYSIS_FAILURE.getCode());
curMQMessageRecord.setRespCode(HandleStatusEnum.ANALYSIS_CALLING.getCode().toString());
curMQMessageRecord.setRespContent(e.getMessage());
curMQMessageRecord.setLastRetryTime(new Date());
curMQMessageRecord.setLastRetryTime(currTime);
curMQMessageRecord.setRetryCount(curMQMessageRecord.getRetryCount()+1);
this.updateById(curMQMessageRecord);
}
}
if (curMQMessageRecord.getSinceType() == 2 && curMQMessageRecord.getSubSinceType() == SubSinceTypeEnum.SinceType52.getCode()) {
if (curMQMessageRecord.getSinceType() == 2 && curMQMessageRecord.getSubSinceType() == SubSinceTypeEnum.SINCETYPE52.getCode()) {
try {
DiffdefeatApprove approveEntity = diffdefeatApproveService.lambdaQuery()
.eq(DiffdefeatApprove::getFormId, oldItem.getFormId())
@@ -404,19 +406,19 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
approveEntity.setApproveOpinion(curDiffDefeatanApprove.getApproveOpinion());
diffdefeatApproveService.updateById(approveEntity);
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.Analysis_Completion.getCode());
curMQMessageRecord.setRespCode("200");
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.ANALYSIS_COMPLETION.getCode());
curMQMessageRecord.setRespCode(HandleStatusEnum.ANALYSIS_NORMAL.getCode().toString());
curMQMessageRecord.setRespContent("");
curMQMessageRecord.setLastRetryTime(new Date());
curMQMessageRecord.setLastRetryTime(currTime);
curMQMessageRecord.setRetryCount(curMQMessageRecord.getRetryCount()+1);
this.updateById(curMQMessageRecord);
}
} catch (Exception e) {
log.error("数据解析异常: {}", e.getMessage());
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.Analysis_Failure.getCode());
curMQMessageRecord.setRespCode("500");
curMQMessageRecord.setTaskStatus(MqTaskStatusEnum.ANALYSIS_FAILURE.getCode());
curMQMessageRecord.setRespCode(HandleStatusEnum.ANALYSIS_CALLING.getCode().toString());
curMQMessageRecord.setRespContent(e.getMessage());
curMQMessageRecord.setLastRetryTime(new Date());
curMQMessageRecord.setLastRetryTime(currTime);
curMQMessageRecord.setRetryCount(curMQMessageRecord.getRetryCount()+1);
this.updateById(curMQMessageRecord);
}
@@ -427,17 +429,18 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
private String getUserStatus(String oldStr)
{
if (oldStr == "确认到店"){
return "客户确认将到店";
}else if (oldStr == "继续观望"){
return "客户继续观望中";
}else if (oldStr == "放弃购车"){
return "客户放弃购车";
}else if (oldStr == "无效通话"){
return "近三天无有效通话记录";
}
if (MessageConvertEnum.CONFIRMED.getCode().equals(oldStr)){
return MessageConvertEnum.CONFIRMED.getMessage();
}else if (MessageConvertEnum.LOOK_ON.getCode().equals(oldStr)){
return MessageConvertEnum.LOOK_ON.getMessage();
}else if (MessageConvertEnum.GIVE_UP.getCode().equals(oldStr)){
return MessageConvertEnum.GIVE_UP.getMessage();
}else if (MessageConvertEnum.INVALID.getCode().equals(oldStr)){
return MessageConvertEnum.INVALID.getMessage();
}else {
return oldStr;
}
}
public void addSliceData(List<DiffDefeatCorpuItem> curDiffDefeatCorpuItems, Date sliceStartTime, Date sliceStopTime, List<SessionItem> sessionItems) {
// 使用Java 8的Stream API来过滤sessionItems
@@ -455,9 +458,9 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
ContentText contentText = JSON.parseObject(sessionItem.getContent(), ContentText.class);
// 根据fromuserrole添加内容到字符串
if (sessionItem.getFromuserrole() == 1) {
str.append("客服").append(contentText.getContent()).append("\n");
str.append(RoleEnum.AGENT.getMessage()).append("").append(contentText.getContent()).append("\n");
} else {
str.append("客户").append(contentText.getContent()).append("\n");
str.append(RoleEnum.USER.getMessage()).append("").append(contentText.getContent()).append("\n");
}
} catch (Exception e) {
log.error("反序列化异常: ", e);
@@ -470,7 +473,7 @@ public class MqMessageRecordServiceImpl extends ServiceImpl<MqMessageRecordMappe
SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
DiffDefeatCorpuItem curDiffDefeatCorpuItem = new DiffDefeatCorpuItem();
curDiffDefeatCorpuItem.setSourceId(slices.get(0).getMsgid());
curDiffDefeatCorpuItem.setCategory("企微");
curDiffDefeatCorpuItem.setCategory(CategoryEnum.ENTERPRISE_WECHAT.getCode());
curDiffDefeatCorpuItem.setHappenTime(sdf.format(slices.get(0).getMsgtime()));
curDiffDefeatCorpuItem.setCorpuText(str.toString());
curDiffDefeatCorpuItems.add(curDiffDefeatCorpuItem);