Merge remote-tracking branch 'origin/dev_20250306_rz' into uat
This commit is contained in:
@@ -10,6 +10,8 @@ import org.springframework.beans.BeanUtils;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.kafka.annotation.KafkaListener;
|
||||
import org.springframework.stereotype.Component;
|
||||
import org.springframework.web.bind.annotation.GetMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
|
||||
import java.time.LocalDateTime;
|
||||
|
||||
@@ -22,6 +24,7 @@ import java.time.LocalDateTime;
|
||||
**/
|
||||
@Slf4j
|
||||
@Component
|
||||
@RestController
|
||||
public class CorpusProcessKafkaConsumer {
|
||||
|
||||
@Autowired
|
||||
@@ -29,6 +32,7 @@ public class CorpusProcessKafkaConsumer {
|
||||
|
||||
private final ObjectMapper objectMapper = new ObjectMapper();
|
||||
|
||||
@GetMapping("corpusProcessKafkaConsumer")
|
||||
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
|
||||
public void listen(String message) {
|
||||
try {
|
||||
@@ -44,6 +48,8 @@ public class CorpusProcessKafkaConsumer {
|
||||
tmTelephoneCorpus.setCreateTime(LocalDateTime.now());
|
||||
tmTelephoneCorpusService.saveTelephoneCorpus(tmTelephoneCorpus);
|
||||
|
||||
tmTelephoneCorpusService.runDify(aicorpusTelephone);
|
||||
|
||||
// 在这里可以添加对解析后的对象的进一步处理逻辑
|
||||
} catch (Exception e) {
|
||||
e.printStackTrace();
|
||||
|
||||
@@ -4,4 +4,6 @@ import com.baomidou.mybatisplus.extension.service.IService;
|
||||
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
|
||||
|
||||
public interface AiAnalysisErrorsService extends IService<AiAnalysisErrors> {
|
||||
|
||||
boolean saveAiAnalysisErrors(AiAnalysisErrors entity);
|
||||
}
|
||||
|
||||
@@ -1,9 +1,13 @@
|
||||
package com.volvo.ai.analytic.center.service;
|
||||
|
||||
import com.alibaba.fastjson.JSONObject;
|
||||
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
|
||||
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
|
||||
|
||||
public interface DiFyService {
|
||||
|
||||
|
||||
public Object getDiFyObject(DiFyReq diFyReq);
|
||||
|
||||
public JSONObject execDifyFlow(DiFyReq diFyReq, String businessType);
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package com.volvo.ai.analytic.center.service;
|
||||
|
||||
import com.baomidou.mybatisplus.extension.service.IService;
|
||||
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
|
||||
import com.volvo.ai.analytic.center.entity.TmTelephoneCorpus;
|
||||
|
||||
import java.util.Map;
|
||||
@@ -14,4 +15,7 @@ public interface TmTelephoneCorpusService extends IService<TmTelephoneCorpus> {
|
||||
|
||||
|
||||
void saveTelephoneCorpus(TmTelephoneCorpus tmTelephoneCorpus);
|
||||
|
||||
|
||||
void runDify(AicorpusTelephoneDTO aicorpusTelephone);
|
||||
}
|
||||
@@ -1,13 +1,31 @@
|
||||
package com.volvo.ai.analytic.center.service.impl;
|
||||
|
||||
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
|
||||
import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
|
||||
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
|
||||
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
|
||||
import com.volvo.ai.analytic.center.mapper.AiAnalysisErrorsMapper;
|
||||
import com.volvo.ai.analytic.center.service.AiAnalysisErrorsService;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
@Slf4j
|
||||
@Service
|
||||
public class AiAnalysisErrorsServiceImpl extends ServiceImpl<AiAnalysisErrorsMapper, AiAnalysisErrors> implements AiAnalysisErrorsService {
|
||||
@Autowired
|
||||
private AiAnalysisErrorsMapper aiAnalysisErrorsMapper;
|
||||
|
||||
@Override
|
||||
public boolean saveAiAnalysisErrors(AiAnalysisErrors entity) {
|
||||
|
||||
LambdaQueryWrapper<AiAnalysisErrors> queryWrapper = new LambdaQueryWrapper<>();
|
||||
queryWrapper.eq(AiAnalysisErrors::getAiAnalysisRequestId, entity.getAiAnalysisRequestId());
|
||||
AiAnalysisErrors oldAiAnalysisErrors = aiAnalysisErrorsMapper.selectOne(queryWrapper);
|
||||
if (oldAiAnalysisErrors == null) {
|
||||
return aiAnalysisErrorsMapper.insert(entity) > 0;
|
||||
} else {
|
||||
return aiAnalysisErrorsMapper.update(entity, queryWrapper) > 0;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,13 +1,22 @@
|
||||
package com.volvo.ai.analytic.center.service.impl;
|
||||
|
||||
import com.alibaba.fastjson.JSON;
|
||||
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.DiFyService;
|
||||
import com.volvo.ai.analytic.center.utils.AiAnalysisUtils;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
@@ -18,6 +27,12 @@ public class DiFyServiceImpl implements DiFyService{
|
||||
@Autowired
|
||||
private DiFyFeign diFyFeign;
|
||||
|
||||
@Autowired
|
||||
private AiAnalysisRequestLogsMapper aiAnalysisRequestLogsMapper;
|
||||
|
||||
@Autowired
|
||||
private AiAnalysisErrorsService aiAnalysisErrorsService;
|
||||
|
||||
@Override
|
||||
public Object getDiFyObject(DiFyReq diFyReq) {
|
||||
Map<String, Object> map = new HashMap<>();
|
||||
@@ -33,4 +48,40 @@ public class DiFyServiceImpl implements DiFyService{
|
||||
}
|
||||
return "";
|
||||
}
|
||||
|
||||
@Override
|
||||
public JSONObject execDifyFlow(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()
|
||||
.aiAnalysisRequestId(aiAnalysisRequestId)
|
||||
.businessRequest("")
|
||||
.difyAgentKey(diFyReq.getFlowId())
|
||||
.difyRequest(JSON.toJSONString(diFyReq))
|
||||
.aiAnalysisRequestType(businessType)
|
||||
.difyResponse(data.toJSONString())
|
||||
.build());
|
||||
|
||||
return data;
|
||||
|
||||
} catch (Exception e) {
|
||||
log.error("dify请求失败",e);
|
||||
aiAnalysisErrorsService.saveAiAnalysisErrors(AiAnalysisErrors.builder()
|
||||
.aiAnalysisRequestId(aiAnalysisRequestId)
|
||||
.aiAnalysisErrorHandlingStatus("0")
|
||||
.aiAnalysisErrorMessage(e.getMessage())
|
||||
.build());
|
||||
return null;
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,14 +1,37 @@
|
||||
package com.volvo.ai.analytic.center.service.impl;
|
||||
|
||||
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.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.resp.CarModelRespDTO;
|
||||
import com.volvo.ai.analytic.center.dto.resp.ResultDTO;
|
||||
import com.volvo.ai.analytic.center.entity.TmTelephoneCorpus;
|
||||
import com.volvo.ai.analytic.center.exception.BizException;
|
||||
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.DiFyService;
|
||||
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.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 org.springframework.transaction.annotation.Transactional;
|
||||
|
||||
import javax.annotation.Resource;
|
||||
import java.util.*;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
|
||||
/**
|
||||
* @description 电话语料表-同步表
|
||||
@@ -20,9 +43,109 @@ import org.springframework.transaction.annotation.Transactional;
|
||||
public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusMapper, TmTelephoneCorpus> implements TmTelephoneCorpusService {
|
||||
|
||||
|
||||
@Autowired
|
||||
private DiFyService diFyService;
|
||||
|
||||
@Resource
|
||||
private RocketMQTemplate rocketMqTemplate;
|
||||
|
||||
@Value("${rocketmq.corpusTelephone.topic}")
|
||||
private String topic;
|
||||
|
||||
@Autowired
|
||||
private RemoteCarModelClient remoteCarModelClient;
|
||||
|
||||
@Override
|
||||
@Transactional
|
||||
public void saveTelephoneCorpus(TmTelephoneCorpus tmTelephoneCorpus) {
|
||||
this.save(tmTelephoneCorpus);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void runDify(AicorpusTelephoneDTO aicorpusTelephone) {
|
||||
|
||||
if(null != aicorpusTelephone){
|
||||
|
||||
Map<String, Object> inputMap = new HashMap();
|
||||
DiFyReq diFyImageReq = new DiFyReq();
|
||||
diFyImageReq.setUser("11111");
|
||||
diFyImageReq.setFlowId("app-peJXSjHjVKdkYxjdOUuPnZ5b");
|
||||
|
||||
|
||||
JSONObject jsonObject = JSONObject.parseObject( aicorpusTelephone.getDisplay());
|
||||
JSONArray segments = jsonObject.getJSONArray("segments");
|
||||
// 遍历 segments
|
||||
segments.stream()
|
||||
.map(segment -> (JSONObject) segment)
|
||||
.forEach(segment -> {
|
||||
JSONObject result = segment.getJSONObject("result");
|
||||
String text = result.getString("text");
|
||||
JSONObject analysisInfo = result.getJSONObject("analysis_info");
|
||||
String role = analysisInfo.getString("role");
|
||||
|
||||
// 拼接 role 和 text
|
||||
String chat = role + ": " + text;
|
||||
System.out.println(chat);
|
||||
inputMap.put("chat",chat);
|
||||
});
|
||||
|
||||
|
||||
inputMap.put("model",getCarModelList());
|
||||
diFyImageReq.setInputs(inputMap);
|
||||
// 获取配置
|
||||
JSONObject execDifyFlow = diFyService.execDifyFlow(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, Object> ltoMap = new HashMap();
|
||||
ltoMap.put("analysisRecordId", aicorpusTelephone.getSourceId());
|
||||
ltoMap.put("analysisScene", "2");
|
||||
// ltoMap.put("unionId", aicorpusTelephone.getSourceId());
|
||||
ltoMap.put("recordId", aicorpusTelephone.getSourceId());
|
||||
ltoMap.put("communicateDate", jsonObject.get("start_time"));
|
||||
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("发送导入标签值成功 消息体:{}", JSON.toJSONString(ltoMap));
|
||||
}
|
||||
@Override
|
||||
public void onException(Throwable e) {
|
||||
log.error("发送导入标签值异常 消息体:{}, 异常:", JSON.toJSONString(ltoMap), e);
|
||||
}
|
||||
}, 10000);
|
||||
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private String getCarModelList(){
|
||||
long startTime = System.currentTimeMillis();
|
||||
CarModelReqDTO carModelReqDTO = new CarModelReqDTO();
|
||||
carModelReqDTO.setOnSale(10041001);
|
||||
carModelReqDTO.setIsValid(10041001);
|
||||
ResultDTO<List<CarModelRespDTO>> cardModelResult= remoteCarModelClient.queryCarModelList(carModelReqDTO);
|
||||
log.info("queryCarModelList 导出查询耗时开始时间:{}",System.currentTimeMillis()-startTime);
|
||||
if (BizEnum.SUCCESS.getCode().toString().equals(cardModelResult.getReturnCode()) && CollectionUtils.isNotEmpty(cardModelResult.getData())) {
|
||||
List<CarModelRespDTO> carModelRespDTOList = cardModelResult.getData();
|
||||
// 拼接 modelName
|
||||
String modelNames = carModelRespDTOList.stream()
|
||||
.map(CarModelRespDTO::getModelName) // 提取 modelName
|
||||
.collect(Collectors.joining(", ")); // 用逗号和空格拼接
|
||||
|
||||
log.info("拼接后的车型名称: {}", modelNames);
|
||||
return modelNames;
|
||||
}
|
||||
return "C40 RECHARGE、EM90、EX30、S60、S90、V60、V90、XC40、XC40 RECHARGE、XC60、XC90";
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
package com.volvo.ai.analytic.center.utils;
|
||||
|
||||
public class FlowResultSplitUtil {
|
||||
|
||||
|
||||
/**
|
||||
* 提取指定任务的内容
|
||||
*
|
||||
* @param text 原始文本
|
||||
* @param startMark 任务起始标记(如 "### 任务1")
|
||||
* @param endMark 任务结束标记(如 "### 任务2"),如果为 null,则提取到文本末尾
|
||||
* @return 任务内容
|
||||
*/
|
||||
public static String flowOutputTextSplit(String text, String startMark, String endMark) {
|
||||
int startIndex = text.indexOf(startMark);
|
||||
if (startIndex == -1) {
|
||||
return "未找到任务起始标记:" + startMark;
|
||||
}
|
||||
|
||||
int endIndex = (endMark != null) ? text.indexOf(endMark) : text.length();
|
||||
if (endIndex == -1) {
|
||||
return "未找到任务结束标记:" + endMark;
|
||||
}
|
||||
|
||||
return text.substring(startIndex, endIndex).trim();
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user