企微语料ai解析开发

This commit is contained in:
zren25
2025-03-10 18:58:32 +08:00
parent d7271ea431
commit 3e8a7f8e5a
19 changed files with 1203 additions and 19 deletions

View File

@@ -208,6 +208,11 @@
<version>3.22.3.1</version>
</dependency>
<dependency>
<groupId>org.apache.poi</groupId>
<artifactId>poi-ooxml</artifactId>
<version>5.2.3</version>
</dependency>
</dependencies>

View File

@@ -0,0 +1,168 @@
package com.volvo.ai.analytic.center.controller;
import com.alibaba.fastjson.JSONObject;
import com.obs.services.model.ObsObject;
import com.volvo.ai.analytic.center.feign.DiFyFeign;
import com.volvo.ai.analytic.center.utils.ObsUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.multipart.MultipartFile;
import java.io.FileInputStream;
import java.io.FileOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.util.HashMap;
import java.util.Map;
import org.apache.poi.ss.usermodel.*;
import org.apache.poi.xssf.usermodel.XSSFWorkbook;
@Slf4j
@RestController
@RequestMapping("/dify")
public class AiRunDataController {
@Autowired
private DiFyFeign diFyFeign;
@PostMapping("/data")
public void fileUpload(String filePath) throws IOException {
JSONObject jsonObjectResult;
try {
// 文件路径
// String filePath = "D:\\cpq\\本月的外呼时长大于10s的明细- 首次到店时间在1月2月 (1).xlsx";
// String filePath = "D:\\cpq\\test20250308.xls";
// 打开文件
FileInputStream file = new FileInputStream(filePath);
Workbook workbook = new XSSFWorkbook(file);
Sheet sheet = workbook.getSheetAt(0); // 获取第一个工作表
// 读取内容
for (Row row : sheet) {
try {
Cell carModel = row.getCell(4);
Cell corpus = row.getCell(7);
String keyword = "";
switch (carModel.toString()){
case "S60":
keyword = "保养、置换";
break;
case "S90":
keyword = "保养、置换";
break;
case "XC60":
keyword = "保养、置换";
break;
case "XC40":
keyword = "置换";
break;
case "XC90":
keyword = "终身用车无忧、置换";
break;
case "V60":
keyword= "置换";
case "V90":
keyword= "置换”";
case "S60 T8":
keyword= "置换、购置税全免";
case "S90 T8":
keyword= "置换、购置税全免";
case "XC60 T8":
keyword= "置换、购置税全免";
case "XC90 T8":
keyword= "终身用车无忧、购置税全免";
case "S60 RECHARGE":
keyword= "保养、置换";
case "S60L":
keyword= "保养、置换";
case "S90 RECHARGE":
keyword= "置换、购置税全免";
case "XC60 RECHARGE":
keyword= "置换、购置税全免";
case "XC90 RECHARGE":
keyword= "终身用车无忧、购置税全免";
case "V90 Cross Country":
keyword= "置换";
case "EX30":
keyword= "限时权益、置换";
case "XC40 RECHARGE":
keyword= "置换";
}
if(!keyword.equals("") && corpus != null){
JSONObject runResultJson = null;
try {
Map<String ,Object> tpMpMap = new HashMap<>();
tpMpMap.put("chat",corpus.toString());
tpMpMap.put("keyword",keyword);
Map<String ,Object> reqMap = new HashMap<>();
reqMap.put("inputs",tpMpMap);
reqMap.put("response_mode","blocking");
reqMap.put("user","streaming232");
// log.info("tpMpMap:{}",tpMpMap);
runResultJson = diFyFeign.runWorkflows("Bearer app-P09tlGfUbSi9s0Ms3BbskEQF",reqMap);
JSONObject outputs = runResultJson.getJSONObject("data").getJSONObject("outputs");
log.info("outputs:{}",outputs);
String outputsReplase = outputs.toJSONString().replace("\\n","").replace("\n","").replace("```json","").replace(" ```","").replace("```","");
JSONObject text = JSONObject.parseObject(outputsReplase);
log.info("text:{}",text);
JSONObject data = text.getJSONObject("text");
if(data.containsKey("质检评分")){
// 设置新列的值,这里可以根据需求设置不同的值
Cell da1 = row.getCell(14);
da1.setCellValue(data.getString("质检评分"));
}
if(data.containsKey("客户反馈")){
Cell da2 = row.getCell(15);
da2.setCellValue(data.getJSONObject("客户反馈").getString("兴趣程度"));
Cell da3 = row.getCell(16);
da3.setCellValue(data.getJSONObject("客户反馈").getString("原因"));
}
} catch (Exception e) {
log.info("error:{},{}",row.getCell(1).toString(),e);
JSONObject outputs = runResultJson.getJSONObject("data").getJSONObject("outputs");
log.info("outputs:{}",outputs);
String outputsReplase = outputs.toJSONString().replace("\\n","").replace("\n","").replace("```json","").replace(" ```","").replace("```","");
JSONObject text = JSONObject.parseObject(outputsReplase);
log.info("text:{}",text);
Cell da1 = row.getCell(17);
da1.setCellValue(String.valueOf(text));
}
}
} catch (Exception e) {
log.info("error:{}",e);
}
}
// 写入文件
FileOutputStream outFile = new FileOutputStream(filePath);
workbook.write(outFile);
outFile.close();
workbook.close();
file.close();
System.out.println("数据已追加写入Excel文件。");
} catch (Exception e) {
log.error("error:{}",e);
}
}
}

View File

@@ -0,0 +1,21 @@
package com.volvo.ai.analytic.center.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD;
import com.volvo.ai.analytic.center.entity.TmOdsVdqwExternalcontact;
import com.volvo.ai.analytic.center.entity.TmOdsVdqwMessagearchiving;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import java.util.List;
/**
* @description 会话存档消息记录表-湖仓同步表
* @author BEJSON
* @date 2025-03-10
*/
@Mapper
public interface TmOdsVdqwExternalcontactMapper extends BaseMapper<TmOdsVdqwExternalcontact> {
}

View File

@@ -0,0 +1,23 @@
package com.volvo.ai.analytic.center.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD;
import com.volvo.ai.analytic.center.entity.TmOdsVdqwMessagearchiving;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import java.util.List;
/**
* @description 会话存档消息记录表-湖仓同步表
* @author BEJSON
* @date 2025-03-10
*/
@Mapper
public interface TmOdsVdqwMessagearchivingMapper extends BaseMapper<TmOdsVdqwMessagearchiving> {
List<OdsVdqwMessageOTD> queryOdsVdqwMessageByData(@Param("statTime") String statTime, @Param("endTime") String endTime);
List<OdsVdqwMessageOTD> queryOdsVdqwMessageByFromUserIdAndAcceptUserId(@Param("statTime") String statTime, @Param("endTime") String endTime, @Param("userIdList") List userIdList);
}

View File

@@ -0,0 +1,19 @@
package com.volvo.ai.analytic.center.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD;
import com.volvo.ai.analytic.center.entity.TmOdsVdqwMessagearchiving;
import com.volvo.ai.analytic.center.entity.TmOdsVdqwWorkuserinfo;
import org.apache.ibatis.annotations.Mapper;
import org.apache.ibatis.annotations.Param;
import java.util.List;
/**
* @description 会话存档消息记录表-湖仓同步表
* @author BEJSON
* @date 2025-03-10
*/
@Mapper
public interface TmOdsVdqwWorkuserinfoMapper extends BaseMapper<TmOdsVdqwWorkuserinfo> {
}

View File

@@ -48,7 +48,7 @@ public class CorpusProcessKafkaConsumer {
tmTelephoneCorpus.setCreateTime(LocalDateTime.now());
tmTelephoneCorpusService.saveTelephoneCorpus(tmTelephoneCorpus);
tmTelephoneCorpusService.runDify(aicorpusTelephone);
tmTelephoneCorpusService.runTelephoneCorpusDify(aicorpusTelephone);
// 在这里可以添加对解析后的对象的进一步处理逻辑
} catch (Exception e) {

View File

@@ -4,4 +4,6 @@ import com.baomidou.mybatisplus.extension.service.IService;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
public interface AiAnalysisRequestLogsService extends IService<AiAnalysisRequestLogs> {
boolean saveAiAnalysisRequestLogs(AiAnalysisRequestLogs aiAnalysisRequestLogs);
}

View File

@@ -9,5 +9,5 @@ public interface DiFyService {
public Object getDiFyObject(DiFyReq diFyReq);
public JSONObject execDifyFlow(DiFyReq diFyReq, String businessType);
public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType);
}

View File

@@ -0,0 +1,16 @@
package com.volvo.ai.analytic.center.service;
import com.baomidou.mybatisplus.extension.service.IService;
import com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD;
import com.volvo.ai.analytic.center.entity.TmOdsVdqwMessagearchiving;
import java.util.*;
/**
* @description 电话语料表-同步表
* @author BEJSON
* @date 2025-03-04
*/
public interface TmOdsVdqwMessagearchivingService extends IService<TmOdsVdqwMessagearchiving> {
void runQiWeiCorpusDify();
}

View File

@@ -17,5 +17,9 @@ public interface TmTelephoneCorpusService extends IService<TmTelephoneCorpus> {
void saveTelephoneCorpus(TmTelephoneCorpus tmTelephoneCorpus);
void runDify(AicorpusTelephoneDTO aicorpusTelephone);
void runTelephoneCorpusDify(AicorpusTelephoneDTO aicorpusTelephone);
String getCarModelList();
}

View File

@@ -1,13 +1,31 @@
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.AiAnalysisErrors;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
import com.volvo.ai.analytic.center.mapper.AiAnalysisRequestLogsMapper;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
@Slf4j
@Service
public class AiAnalysisRequestLogsServiceImpl extends ServiceImpl<AiAnalysisRequestLogsMapper, AiAnalysisRequestLogs> implements AiAnalysisRequestLogsService {
@Autowired
private AiAnalysisRequestLogsMapper aiAnalysisRequestLogsMapper;
@Override
public boolean saveAiAnalysisRequestLogs(AiAnalysisRequestLogs aiAnalysisRequestLogs) {
LambdaQueryWrapper<AiAnalysisRequestLogs> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisRequestLogs.getAiAnalysisRequestId());
AiAnalysisRequestLogs oldAiAnalysisRequestLogs= aiAnalysisRequestLogsMapper.selectOne(queryWrapper);
if (oldAiAnalysisRequestLogs == null) {
return aiAnalysisRequestLogsMapper.insert(aiAnalysisRequestLogs) > 0;
} else {
return aiAnalysisRequestLogsMapper.update(aiAnalysisRequestLogs, queryWrapper) > 0;
}
}
}

View File

@@ -5,11 +5,9 @@ import com.alibaba.fastjson.JSONObject;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
import com.volvo.ai.analytic.center.feign.DiFyFeign;
import com.volvo.ai.analytic.center.mapper.AiAnalysisErrorsMapper;
import com.volvo.ai.analytic.center.mapper.AiAnalysisRequestLogsMapper;
import com.volvo.ai.analytic.center.service.AiAnalysisErrorsService;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.DiFyService;
import com.volvo.ai.analytic.center.utils.AiAnalysisUtils;
import lombok.extern.slf4j.Slf4j;
@@ -28,7 +26,7 @@ public class DiFyServiceImpl implements DiFyService{
private DiFyFeign diFyFeign;
@Autowired
private AiAnalysisRequestLogsMapper aiAnalysisRequestLogsMapper;
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@Autowired
private AiAnalysisErrorsService aiAnalysisErrorsService;
@@ -50,24 +48,32 @@ public class DiFyServiceImpl implements DiFyService{
}
@Override
public JSONObject execDifyFlow(DiFyReq diFyReq, String businessType) {
public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType) {
String aiAnalysisRequestId = AiAnalysisUtils.getAiAnalysisRequestId(businessType);
try {
Map<String, Object> map = new HashMap<>();
map.put("inputs",diFyReq.getInputs());
map.put("response_mode","blocking");
map.put("user",diFyReq.getUser());
JSONObject difyResult = diFyFeign.runWorkflows("Bearer "+diFyReq.getFlowId(),map);
JSONObject data = difyResult.getJSONObject("data");
log.info("execDifyFlow dify response data:{}",data);
// 保存请求日志
aiAnalysisRequestLogsMapper.insert(AiAnalysisRequestLogs.builder()
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder()
.aiAnalysisRequestId(aiAnalysisRequestId)
.businessRequest(JSONObject.toJSONString(""))
.difyAgentKey(diFyReq.getFlowId())
.difyRequest(JSON.toJSONString(diFyReq))
.aiAnalysisRequestType(businessType)
.build());
JSONObject difyResult = diFyFeign.runWorkflows("Bearer "+diFyReq.getFlowId(),map);
JSONObject data = difyResult.getJSONObject("data");
log.info("execDifyFlow dify response data:{}",data);
data.put("aiAnalysisRequestId",aiAnalysisRequestId);
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder()
.aiAnalysisRequestId(aiAnalysisRequestId)
.businessRequest(JSONObject.toJSONString(""))
.difyResponse(data.toJSONString())
.build());

View File

@@ -0,0 +1,188 @@
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.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.volvo.ai.analytic.center.constant.Constant;
import com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import com.volvo.ai.analytic.center.dto.req.RunMaskingRuleInput;
import com.volvo.ai.analytic.center.entity.*;
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
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.utils.FlowResultSplitUtil;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/**
* @description 企微语料表-同步表
* @author rz
* @date 2025-03-04
*/
@Slf4j
@Service
public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwMessagearchivingMapper, TmOdsVdqwMessagearchiving> implements TmOdsVdqwMessagearchivingService {
@Autowired
private TmOdsVdqwMessagearchivingMapper tmOdsVdqwMessagearchivingMapper;
@Autowired
private TmOdsVdqwExternalcontactMapper tmOdsVdqwExternalcontactMapper;
@Autowired
private TmOdsVdqwWorkuserinfoMapper tmOdsVdqwWorkuserinfoMapper;
@Autowired
private DiFyService diFyService;
@Resource
private RocketMQTemplate rocketMqTemplate;
@Value("${rocketmq.corpusTelephone.topic}")
private String topic;
@Autowired
private RemoteCarModelClient remoteCarModelClient;
@Autowired
private DataMaskingRuleService dataMaskingRuleService;
@Autowired
private TmTelephoneCorpusService tmTelephoneCorpusService;
@Override
public void runQiWeiCorpusDify() {
String statTime = "2025-03-01";
String endTime = "2025-03-02";
List<OdsVdqwMessageOTD> messageList = tmOdsVdqwMessagearchivingMapper.queryOdsVdqwMessageByData(statTime,endTime);
messageList.stream().forEach(item->{
log.info("企微语料内容FromUserId{}, AcceptUserId{}",item.getFromUserId(), item.getAcceptUserId());
// 1vdqw_workuserinfo 这个表对应是 B端认证中心userId
// 2vdqw_externalcontact 这个表对应是 企微客户unionId
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)){
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()));
OdsVdqwMessageOTD maxMsgTimeItem = contetnList.stream()
.max((o1, o2) -> o1.getMsgTime().compareTo(o2.getMsgTime()))
.orElse(null);
contetnList.stream().forEach(contentItem->{
Map<String, Object> inputMap = new HashMap();
DiFyReq diFyImageReq = new DiFyReq();
diFyImageReq.setUser("11111");
diFyImageReq.setFlowId("app-peJXSjHjVKdkYxjdOUuPnZ5b");
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);
}
});
}
});
}
private String getUnionId(List<String> userIds){
LambdaQueryWrapper<TmOdsVdqwExternalcontact> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.in(TmOdsVdqwExternalcontact::getExternalUserId, userIds);
queryWrapper.eq(TmOdsVdqwExternalcontact::getIsDeleted, 0);
List<TmOdsVdqwExternalcontact> oldAiAnalysisRequestLogs= tmOdsVdqwExternalcontactMapper.selectList(queryWrapper);
if(CollectionUtils.isNotEmpty(oldAiAnalysisRequestLogs)){
return oldAiAnalysisRequestLogs.get(0).getUnionId();
}
return "";
}
private String getUserId(List<String> userIds){
LambdaQueryWrapper<TmOdsVdqwWorkuserinfo> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.in(TmOdsVdqwWorkuserinfo::getUserId, userIds);
queryWrapper.eq(TmOdsVdqwWorkuserinfo::getIsDeleted, 0);
List<TmOdsVdqwWorkuserinfo> tmOdsVdqwWorkuserinfoList = tmOdsVdqwWorkuserinfoMapper.selectList(queryWrapper);
if(CollectionUtils.isNotEmpty(tmOdsVdqwWorkuserinfoList)){
return tmOdsVdqwWorkuserinfoList.get(0).getMiddleUserId().toString();
}
return "";
}
}

View File

@@ -4,16 +4,20 @@ import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.volvo.ai.analytic.center.constant.Constant;
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
import com.volvo.ai.analytic.center.dto.req.CarModelReqDTO;
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.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.TmTelephoneCorpus;
import com.volvo.ai.analytic.center.enums.BizEnum;
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
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.TmTelephoneCorpusService;
import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil;
@@ -55,6 +59,9 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
@Autowired
private RemoteCarModelClient remoteCarModelClient;
@Autowired
private DataMaskingRuleService dataMaskingRuleService;
@Override
@Transactional
public void saveTelephoneCorpus(TmTelephoneCorpus tmTelephoneCorpus) {
@@ -62,10 +69,15 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
}
@Override
public void runDify(AicorpusTelephoneDTO aicorpusTelephone) {
public void runTelephoneCorpusDify(AicorpusTelephoneDTO aicorpusTelephone) {
if(null != aicorpusTelephone){
List<DataMaskingRule> maskingRuleItems = dataMaskingRuleService.getDataMaskingRuleListByApplicationChannel(Constant.CHANNEL_DCC);
RunMaskingRuleInput runMaskingRuleInput = new RunMaskingRuleInput();
runMaskingRuleInput.setDataMaskingRules(maskingRuleItems);
Map<String, Object> inputMap = new HashMap();
DiFyReq diFyImageReq = new DiFyReq();
diFyImageReq.setUser("11111");
@@ -85,14 +97,16 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
// 拼接 role 和 text
String chat = role + ": " + text;
inputMap.put("chat",chat);
runMaskingRuleInput.setOldStr(chat);
String corpusChat = dataMaskingRuleService.runMaskingRule(runMaskingRuleInput);
inputMap.put("chat",corpusChat);
});
inputMap.put("model",getCarModelList());
diFyImageReq.setInputs(inputMap);
// 获取配置
JSONObject execDifyFlow = diFyService.execDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode());
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode());
log.info("runDify execDifyFlow {}",execDifyFlow);
if(null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")){
@@ -100,7 +114,7 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
String resultStrOne = FlowResultSplitUtil.flowOutputTextSplit(text, "任务1", "任务2");
String resultStrTwo =FlowResultSplitUtil.flowOutputTextSplit(text, "任务2", null);
Map<String, String> ltoMap = new HashMap();
ltoMap.put("analysisRecordId", aicorpusTelephone.getSourceId());
ltoMap.put("analysisRecordId", execDifyFlow.getString("aiAnalysisRequestId"));
ltoMap.put("analysisScene", "2");
// ltoMap.put("unionId", aicorpusTelephone.getSourceId());
ltoMap.put("recordId", aicorpusTelephone.getSourceId());
@@ -114,11 +128,11 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
log.info("发送导入标签值成功 消息体:{}", JSON.toJSONString(ltoMap));
log.info("电话语料发送MQ成功 消息体:{}", JSON.toJSONString(ltoMap));
}
@Override
public void onException(Throwable e) {
log.error("发送导入标签值异常 消息体:{}, 异常:", JSON.toJSONString(ltoMap), e);
log.error("电话语料发送MQ异常 消息体:{}, 异常:", JSON.toJSONString(ltoMap), e);
}
}, 10000);
@@ -126,7 +140,7 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
}
}
private String getCarModelList(){
public String getCarModelList(){
try {
long startTime = System.currentTimeMillis();
CarModelReqDTO carModelReqDTO = new CarModelReqDTO();
@@ -151,4 +165,5 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
}
}

View File

@@ -0,0 +1,49 @@
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper">
<select id="queryOdsVdqwMessageByData" resultType="com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD" >
SELECT
tovm.from_user_id as fromUserUnId,
tovm.accept_user_id as acceptUserId,
tovm.msg_time as msgTime,
tovm.chat_type as chatType
FROM
`tm_ods_vdqw_messagearchiving` tovm
WHERE
tovm.msg_time between #{statTime} and #{endTime}
AND tovm.chat_type = 0
GROUP BY
tovm.from_user_id,
tovm.accept_user_id
</select>
<select id="queryOdsVdqwMessageByFromUserIdAndAcceptUserId" resultType="com.volvo.ai.analytic.center.dto.corpus.OdsVdqwMessageOTD" >
SELECT
tovm.msg_time,
tovm.content,
tovm.from_user_id,
tovm.accept_user_id
FROM
`tm_ods_vdqw_messagearchiving` tovm
WHERE
tovm.msg_time between #{statTime} and #{endTime}
AND tovm.from_user_id IN
<foreach collection="fromUserId" item="userIdList" open="(" separator="," close=")">
#{userIdList}
</foreach>
AND tovm.accept_user_id IN
<foreach collection="fromUserId" item="userIdList" open="(" separator="," close=")">
#{userIdList}
</foreach>
AND tovm.chat_type = 0
and tovm.msg_type='text'
and tovm.action_type='send'
ORDER BY
tovm.msg_time
</select>
</mapper>