点赞点踩

This commit is contained in:
zren25
2025-03-11 17:05:51 +08:00
parent a6485821e9
commit d75eb448b4
11 changed files with 338 additions and 72 deletions

View File

@@ -0,0 +1,18 @@
package com.volvo.ai.analytic.center.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.volvo.ai.analytic.center.entity.TmCorpusReport;
import com.volvo.ai.analytic.center.entity.TmTelephoneCorpus;
import org.apache.ibatis.annotations.Mapper;
import org.springframework.stereotype.Repository;
/**
* @description 电话语料表-同步表
* @author BEJSON
* @date 2025-03-04
*/
@Mapper
public interface TmCorpusReportMapper extends BaseMapper<TmCorpusReport> {
}

View File

@@ -1,10 +1,17 @@
package com.volvo.ai.analytic.center.mq;
import com.alibaba.fastjson.JSONObject;
import com.volvo.ai.analytic.center.entity.TmCorpusReport;
import com.volvo.ai.analytic.center.service.TmCorpusReportService;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.PostMapping;
import java.util.Date;
@Slf4j
@Component
@@ -13,9 +20,26 @@ import org.springframework.stereotype.Component;
enableMsgTrace = true)
public class CorpushIsLikeConsumer implements RocketMQListener<MessageExt>{
@Autowired
private TmCorpusReportService tmCorpusReportService;
@Override
@PostMapping
public void onMessage(MessageExt messageExt) {
log.info("CorpushIsLikeConsumer message: " + messageExt);
String message = new String(messageExt.getBody());
try {
// 保存报告
JSONObject execDifyFlow = JSONObject.parseObject(message);
Integer like = execDifyFlow.getInteger("isLike");
tmCorpusReportService.saveTmCorpusReport(TmCorpusReport.builder()
.aiAnalysisRequestId(execDifyFlow.getString("aiAnalysisRequestId"))
.isLike(like)
.updateBy("LTO")
.updateTime(new Date())
.build());
} catch (Exception e) {
log.info(" 电话语料处理保存报告异常processItem{} ", e);
}
}
}

View File

@@ -0,0 +1,9 @@
package com.volvo.ai.analytic.center.service;
import com.baomidou.mybatisplus.extension.service.IService;
import com.volvo.ai.analytic.center.entity.TmCorpusReport;
public interface TmCorpusReportService extends IService<TmCorpusReport> {
boolean saveTmCorpusReport(TmCorpusReport tmCorpusReport);
}

View File

@@ -62,13 +62,14 @@ public class DiFyServiceImpl implements DiFyService{
.businessRequest(JSONObject.toJSONString(""))
.difyAgentKey(diFyReq.getFlowId())
.difyRequest(JSON.toJSONString(diFyReq))
.difyResponse(JSON.toJSONString(""))
.aiAnalysisRequestType(businessType)
.build());
log.info("execDifyFlow dify request data:{}",map);
JSONObject difyResult = diFyFeign.runWorkflows("Bearer "+diFyReq.getFlowId(),map);
JSONObject data = difyResult.getJSONObject("data");
log.info("execDifyFlow dify response data:{}",data);
log.info("execDifyFlow dify response aiID:{} , data:{}",aiAnalysisRequestId,data);
data.put("aiAnalysisRequestId",aiAnalysisRequestId);
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder()

View File

@@ -0,0 +1,30 @@
package com.volvo.ai.analytic.center.service.impl;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.volvo.ai.analytic.center.entity.TmCorpusReport;
import com.volvo.ai.analytic.center.mapper.TmCorpusReportMapper;
import com.volvo.ai.analytic.center.service.TmCorpusReportService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
@Slf4j
@Service
public class TmCorpusReportServiceImpl extends ServiceImpl<TmCorpusReportMapper, TmCorpusReport> implements TmCorpusReportService {
@Autowired
private TmCorpusReportMapper tmCorpusReportMapper;
@Override
public boolean saveTmCorpusReport(TmCorpusReport tmCorpusReport) {
LambdaQueryWrapper<TmCorpusReport> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(TmCorpusReport::getAiAnalysisRequestId, tmCorpusReport.getAiAnalysisRequestId());
TmCorpusReport oldTmCorpusReport= tmCorpusReportMapper.selectOne(queryWrapper);
if (oldTmCorpusReport == null) {
return tmCorpusReportMapper.insert(tmCorpusReport) > 0;
} else {
return tmCorpusReportMapper.update(tmCorpusReport, queryWrapper) > 0;
}
}
}

View File

@@ -17,10 +17,7 @@ import com.volvo.ai.analytic.center.feign.RemoteCarModelClient;
import com.volvo.ai.analytic.center.mapper.TmOdsVdqwExternalcontactMapper;
import com.volvo.ai.analytic.center.mapper.TmOdsVdqwMessagearchivingMapper;
import com.volvo.ai.analytic.center.mapper.TmOdsVdqwWorkuserinfoMapper;
import com.volvo.ai.analytic.center.service.DataMaskingRuleService;
import com.volvo.ai.analytic.center.service.DiFyService;
import com.volvo.ai.analytic.center.service.TmOdsVdqwMessagearchivingService;
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
import com.volvo.ai.analytic.center.service.*;
import com.volvo.ai.analytic.center.utils.ConstantStr;
import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil;
import lombok.extern.slf4j.Slf4j;
@@ -37,10 +34,7 @@ import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.time.LocalDate;
import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.*;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@@ -86,6 +80,9 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
@Autowired
private TmTelephoneCorpusService tmTelephoneCorpusService;
@Autowired
private TmCorpusReportService tmCorpusReportService;
@Override
public void runQiWeiCorpusDify(String paramJson) {
@@ -128,7 +125,7 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
try {
processItem(item, finalStatTime, finalEndTime);
} catch (Exception e) {
log.error("处理消息失败: FromUserId={}, AcceptUserId={}, 异常: {}",
log.error("处理企微语料失败: FromUserId={}, AcceptUserId={}, 异常: {}",
item.getFromUserId(), item.getAcceptUserId(), e.getMessage(), e);
}
}, executor))
@@ -148,69 +145,91 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
String unionId = getUnionId(Arrays.asList(item.getFromUserId(), item.getAcceptUserId()));
String userId = getUserId(Arrays.asList(item.getFromUserId(), item.getAcceptUserId()));
log.info("企微查询信息unionId{}, userId{}", unionId, userId);
if (StringUtils.isEmpty(unionId) || StringUtils.isNotBlank(userId)) {
if (StringUtils.isEmpty(unionId) || StringUtils.isEmpty(userId)) {
log.info("企微查询信息为空 ");
return;
}
if (StringUtils.isNotBlank(item.getFromUserId()) && StringUtils.isNotBlank(item.getAcceptUserId())) {
List<OdsVdqwMessageOTD> contetnList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByFromUserIdAndAcceptUserId(statTime, endTime, Arrays.asList(item.getFromUserId(), item.getAcceptUserId()));
List<OdsVdqwMessageOTD> contetnList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByFromUserIdAndAcceptUserId(statTime, endTime, Arrays.asList(item.getFromUserId(), item.getAcceptUserId()) );
OdsVdqwMessageOTD maxMsgTimeItem = contetnList.stream()
.max((o1, o2) -> o1.getMsgTime().compareTo(o2.getMsgTime()))
.orElse(null);
Map<String, Object> inputMap = new HashMap<>();
DiFyReq diFyImageReq = new DiFyReq();
diFyImageReq.setUser(ConstantStr.corpus_user);
diFyImageReq.setFlowId(qiweiToken);
List<DataMaskingRule> maskingRuleItems = dataMaskingRuleService.getDataMaskingRuleListByApplicationChannel(Constant.CHANNEL_DCC);
RunMaskingRuleInput runMaskingRuleInput = new RunMaskingRuleInput();
runMaskingRuleInput.setDataMaskingRules(maskingRuleItems);
StringBuffer chatList = new StringBuffer();
contetnList.forEach(contentItem -> {
Map<String, Object> inputMap = new HashMap<>();
DiFyReq diFyImageReq = new DiFyReq();
diFyImageReq.setUser(ConstantStr.corpus_user);
diFyImageReq.setFlowId(qiweiToken);
JSONObject contentJson = JSONObject.parseObject(contentItem.getContent());
String content = contentJson.getString("content");
List<DataMaskingRule> maskingRuleItems = dataMaskingRuleService.getDataMaskingRuleListByApplicationChannel(Constant.CHANNEL_DCC);
RunMaskingRuleInput runMaskingRuleInput = new RunMaskingRuleInput();
runMaskingRuleInput.setDataMaskingRules(maskingRuleItems);
// 拼接 role 和 text
String chat = contentItem.getFromUserId().concat(":").concat(content);
runMaskingRuleInput.setOldStr(chat);
String corpusChat = dataMaskingRuleService.runMaskingRule(runMaskingRuleInput);
inputMap.put("chat", corpusChat);
inputMap.put("model", tmTelephoneCorpusService.getCarModelList());
diFyImageReq.setInputs(inputMap);
// 获取配置
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode());
log.info("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);
Map<String, String> ltoMap = new HashMap<>();
ltoMap.put("analysisRecordId", execDifyFlow.getString("aiAnalysisRequestId"));
ltoMap.put("analysisScene", "1");
ltoMap.put("unionId", unionId);
ltoMap.put("consultantId", userId);
ltoMap.put("communicateDate", DateUtil.format(maxMsgTimeItem.getMsgTime(), DatePattern.NORM_DATETIME_PATTERN));
ltoMap.put("analysisResult", resultStrOne);
ltoMap.put("analysisDetail", resultStrTwo);
// 发送MQ
log.info("send mq {}", ltoMap);
rocketMqTemplate.asyncSend(topic, MessageBuilder.withPayload(JSON.toJSONString(ltoMap)).build(),
new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
log.info("企微语料发送MQ成功 消息体:{}", JSON.toJSONString(ltoMap));
}
@Override
public void onException(Throwable e) {
log.error("企微语料发送MQ异常 消息体:{}, 异常:", JSON.toJSONString(ltoMap), e);
}
}, 10000);
}
chatList.append(corpusChat).append("\n");
});
inputMap.put("chat", chatList.toString());
// inputMap.put("model", tmTelephoneCorpusService.getCarModelList());
diFyImageReq.setInputs(inputMap);
// 获取配置
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode());
log.info("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);
Map<String, String> ltoMap = new HashMap<>();
ltoMap.put("analysisRecordId", execDifyFlow.getString("aiAnalysisRequestId"));
ltoMap.put("analysisScene", "1");
ltoMap.put("unionId", unionId);
ltoMap.put("consultantId", userId);
ltoMap.put("communicateDate", DateUtil.format(maxMsgTimeItem.getMsgTime(), DatePattern.NORM_DATETIME_PATTERN));
ltoMap.put("analysisResult", resultStrOne.replace("#",""));
ltoMap.put("analysisDetail", resultStrTwo);
// 发送MQ
log.info("send mq {}", ltoMap);
rocketMqTemplate.asyncSend(topic, MessageBuilder.withPayload(JSON.toJSONString(ltoMap)).build(),
new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
log.info("企微语料发送MQ成功 消息体:{}", JSON.toJSONString(ltoMap));
}
@Override
public void onException(Throwable e) {
log.error("企微语料发送MQ异常 消息体:{}, 异常:", JSON.toJSONString(ltoMap), e);
}
}, 10000);
try {
// 保存报告
tmCorpusReportService.saveTmCorpusReport(TmCorpusReport.builder()
.aiAnalysisRequestId(execDifyFlow.getString("aiAnalysisRequestId"))
.corpusType(1l)
.corpusTime(maxMsgTimeItem.getMsgTime())
.userId(userId)
.unionId(unionId)
.reportTitle(resultStrOne)
.reportInfo(resultStrTwo)
.isLike(0)
.isDeleted(0)
.version(0)
.createBy("system")
.updateBy("")
.createSqlby("")
.updateSqlby("")
.createTime(new Date())
.build());
} catch (Exception e) {
log.info(" 企业语料处理保存报告异常processItem{} ", e);
}
}
}
}

View File

@@ -1,5 +1,7 @@
package com.volvo.ai.analytic.center.service.impl;
import cn.hutool.core.date.DatePattern;
import cn.hutool.core.date.DateUtil;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
@@ -12,6 +14,7 @@ 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.ResultDTO;
import com.volvo.ai.analytic.center.entity.DataMaskingRule;
import com.volvo.ai.analytic.center.entity.TmCorpusReport;
import com.volvo.ai.analytic.center.entity.TmTelephoneCorpus;
import com.volvo.ai.analytic.center.enums.BizEnum;
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
@@ -19,6 +22,7 @@ import com.volvo.ai.analytic.center.feign.RemoteCarModelClient;
import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper;
import com.volvo.ai.analytic.center.service.DataMaskingRuleService;
import com.volvo.ai.analytic.center.service.DiFyService;
import com.volvo.ai.analytic.center.service.TmCorpusReportService;
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
import com.volvo.ai.analytic.center.utils.ConstantStr;
import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil;
@@ -65,6 +69,8 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
@Autowired
private DataMaskingRuleService dataMaskingRuleService;
@Autowired
private TmCorpusReportService tmCorpusReportService;
@Override
@Transactional
@@ -91,6 +97,7 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
JSONObject jsonObject = JSONObject.parseObject( aicorpusTelephone.getDisplay());
JSONArray segments = jsonObject.getJSONArray("segments");
// 遍历 segments
StringBuffer chatList = new StringBuffer();
segments.stream()
.map(segment -> (JSONObject) segment)
.forEach(segment -> {
@@ -103,10 +110,11 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
String chat = role + ": " + text;
runMaskingRuleInput.setOldStr(chat);
String corpusChat = dataMaskingRuleService.runMaskingRule(runMaskingRuleInput);
inputMap.put("chat",corpusChat);
chatList.append(corpusChat).append("\n");
});
inputMap.put("chat",chatList.toString());
inputMap.put("model",getCarModelList());
diFyImageReq.setInputs(inputMap);
// 获取配置
@@ -123,7 +131,7 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
// ltoMap.put("unionId", aicorpusTelephone.getSourceId());
ltoMap.put("recordId", aicorpusTelephone.getSourceId());
ltoMap.put("communicateDate", jsonObject.getString("start_time"));
ltoMap.put("analysisResult", resultStrOne);
ltoMap.put("analysisResult", resultStrOne.replace("#",""));
ltoMap.put("analysisDetail", resultStrTwo);
// 发送MQ
@@ -139,7 +147,28 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
log.error("电话语料发送MQ异常 消息体:{}, 异常:", JSON.toJSONString(ltoMap), e);
}
}, 10000);
try {
// 保存报告
Date corpusTime = DateUtil.parseUTC(jsonObject.getString("start_time"));
tmCorpusReportService.saveTmCorpusReport(TmCorpusReport.builder()
.aiAnalysisRequestId(execDifyFlow.getString("aiAnalysisRequestId"))
.corpusType(2l)
.corpusTime(corpusTime)
.userId( aicorpusTelephone.getSourceId())
.reportTitle(resultStrOne)
.reportInfo(resultStrTwo)
.isLike(0)
.isDeleted(0)
.version(0)
.createBy("system")
.updateBy("")
.createSqlby("")
.updateSqlby("")
.createTime(new Date())
.build());
} catch (Exception e) {
log.info(" 电话语料处理保存报告异常processItem{} ", e);
}
}
}
}

View File

@@ -23,7 +23,7 @@
<select id="queryOdsVdqwMessageByData" resultType="com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD" >
SELECT
tovm.from_user_id as fromUserUnId,
tovm.from_user_id as fromUserId,
tovm.accept_user_id as acceptUserId,
tovm.msg_time as msgTime,
tovm.chat_type as chatType
@@ -50,12 +50,12 @@
WHERE
tovm.msg_time between #{statTime} and #{endTime}
AND tovm.from_user_id IN
<foreach collection="fromUserId" item="userIdList" open="(" separator="," close=")">
#{userIdList}
<foreach collection="userIdList" item="userId" open="(" separator="," close=")">
#{userId}
</foreach>
AND tovm.accept_user_id IN
<foreach collection="fromUserId" item="userIdList" open="(" separator="," close=")">
#{userIdList}
<foreach collection="userIdList" item="userId" open="(" separator="," close=")">
#{userId}
</foreach>
AND tovm.chat_type = 0
and tovm.msg_type='text'