请求AI解析代码提交

This commit is contained in:
zren25
2025-04-15 18:06:27 +08:00
parent 766a8729a4
commit 0f9e21e86c
15 changed files with 621 additions and 156 deletions

View File

@@ -0,0 +1,33 @@
package com.volvo.ai.analytic.center.dto.req;
import lombok.Data;
import org.springframework.validation.annotation.Validated;
import javax.validation.constraints.NotNull;
/**
*
* @ClassName: AnalysisRequest
* @author: renzhen
* @Description: Ai解析请求参数
* @date: 2025-04-15 13:39
*/
@Data
public class AnalysisReq {
// 请求的所有参数
@NotNull(message = "请求语料不能为空")
private Object data;
// 回调地址
@NotNull(message = "回调地址不能为空")
private String callbackUrl;
// aiId
private String aiAnalysisRequestId;
// 业务类型
@NotNull(message = "业务类型不能为空")
private String aiAnalysisRequestType;
}

View File

@@ -0,0 +1,24 @@
package com.volvo.ai.analytic.center.dto.resp;
import lombok.Data;
/**
*
* @ClassName: AnalysisRequest
* @author: renzhen
* @Description: 工作流处理结果
* @date: 2025-04-15 13:39
*/
@Data
public class AnalysisDifyResultDTO{
// 响应
private String workflowRunId;
private String workflowAppId;
private String workUserId;
// 分析中心唯一ID
private String aiAnalysisRequestId;
private String difyResponse;
private String aiAnalysisRequestType;
}

View File

@@ -0,0 +1,57 @@
package com.volvo.ai.analytic.center.dto.resp;
import com.volvo.common.core.constant.CommonConstants;
import com.volvo.common.core.util.ResultMsg;
import lombok.Data;
/**
*
* @ClassName: AnalysisRequest
* @author: renzhen
* @Description: Ai解析响应
* @date: 2025-04-15 13:39
*/
@Data
public class AnalysisResp<T> {
// 响应
private T data;
// 分析中心唯一ID
private String aiAnalysisRequestId;
private int code;
private String msg;
public static <T> AnalysisResp<T> success(String message) {
return (AnalysisResp<T>) analysisResp((Object)null, CommonConstants.SUCCESS, message);
}
public static <T> AnalysisResp<T> success(T data, String aiAnalysisRequestId) {
AnalysisResp<T> analysisResp = new AnalysisResp();
analysisResp.setAiAnalysisRequestId(aiAnalysisRequestId);
analysisResp.setData(data);
analysisResp.setMsg("ok");
analysisResp.setCode(CommonConstants.SUCCESS);
return analysisResp;
}
public static <T> AnalysisResp<T> failed(String message) {
return (AnalysisResp<T>) analysisResp((Object)null, CommonConstants.FAIL, message);
}
private static <T> AnalysisResp<T> analysisResp(T data, int code, String msg) {
AnalysisResp<T> analysisResp = new AnalysisResp();
analysisResp.setCode(code);
analysisResp.setData(data);
analysisResp.setMsg(msg);
return analysisResp;
}
private static <T> AnalysisResp<T> analysisResp(T data, String aiAnalysisRequestId) {
AnalysisResp<T> analysisResp = new AnalysisResp();
analysisResp.setAiAnalysisRequestId(aiAnalysisRequestId);
analysisResp.setData(data);
analysisResp.setMsg("ok");
analysisResp.setCode(CommonConstants.SUCCESS);
return analysisResp;
}
}

View File

@@ -43,6 +43,18 @@ public class AiAnalysisRequestLogs {
@TableField("dify_agent_key")
private String difyAgentKey;
@TableField("workflowRunId")
private String workflowRunId;
@TableField("workflowAppId")
private String workflowAppId;
@TableField("workUserId")
private String workUserId;
@TableField("callback_url")
private String callbackUrl;
@TableField("is_deleted")
@TableLogic
private Integer isDeleted;

View File

@@ -0,0 +1,73 @@
package com.volvo.ai.analytic.center.entity;
import com.baomidou.mybatisplus.annotation.*;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.Date;
@Data
@Builder
@AllArgsConstructor
@NoArgsConstructor
@TableName("tc_business_type")
public class TcBusinessType {
@TableId(value = "id", type = IdType.AUTO)
private Long id;
@TableField("business_request_type")
private String businessRequestType;
@TableField("business_request_desc")
private String businessRequestDesc;
@TableField("business_type_topic")
private String businessTypeTopic; // JSON 字符串
@TableField("business_type_topic_tag")
private String businessTypeTopicTag;
@TableField("workflow_api_key")
private String workflowApiKey; // JSON 字符串
@TableField("workflow_user")
private String workflowUser;
@TableField("max_retry_count")
private Integer maxRetryCount;
@TableField("is_deleted")
@TableLogic
private Integer isDeleted;
@TableField("versions")
@Version
private Integer versions;
/**
* 创建者
*/
@TableField("create_by")
private String createBy;
/**
* 创建时间
*/
@TableField("create_time")
private Date createTime;
/**
* 更新者
*/
@TableField("update_by")
private String updateBy;
/**
* 更新时间
*/
@TableField("update_time")
private Date updateTime;
}

View File

@@ -1,7 +1,9 @@
package com.volvo.ai.analytic.center.controller;
import com.volvo.ai.analytic.center.service.AiDifyResultService;
import com.volvo.ai.analytic.center.dto.req.AnalysisReq;
import com.volvo.ai.analytic.center.dto.resp.AnalysisResp;
import com.volvo.ai.analytic.center.service.AiAnalysisDifyService;
import com.volvo.common.core.util.ResultMsg;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
@@ -11,21 +13,20 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
@RestController
@Api(tags = "测试类API")
@RestController("analysis")
@Api(tags = "分析中心接口")
@Slf4j
@RefreshScope
public class AiDifyResultController {
public class AiAnalysisDifyController {
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Autowired
private AiDifyResultService AiDifyResultService;
private AiAnalysisDifyService AiDifyResultService;
@@ -37,6 +38,13 @@ public class AiDifyResultController {
return ResultMsg.ok("ok");
}
@PostMapping("/aiAnalyze")
@ApiOperation(value = "Ai解析接口")
public AnalysisResp<Object> aiAnalyze(@RequestBody AnalysisReq analysisReq) {
log.info("aiAnalyze data: {}",analysisReq);
return AnalysisResp.success("ok");
}
}

View File

@@ -0,0 +1,11 @@
package com.volvo.ai.analytic.center.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.volvo.ai.analytic.center.entity.TcBusinessType;
import org.apache.ibatis.annotations.Mapper;
@Mapper
public interface TcBusinessTypeMapper extends BaseMapper<TcBusinessType> {
}

View File

@@ -0,0 +1,82 @@
package com.volvo.ai.analytic.center.mq;
import com.alibaba.fastjson.JSONObject;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import com.volvo.ai.analytic.center.dto.resp.AnalysisDifyResultDTO;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.DiFyService;
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.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.http.ResponseEntity;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.client.RestTemplate;
import java.util.concurrent.CompletableFuture;
/**
* @ClassName AnalysisDifyMqConsumer
* @Description AI解析MQ-Callback处理
* @Author renzhen
* @Date 2025-03-04 10:18
* @Version 1.0
**/
@Slf4j
@Component
@RefreshScope
@RestController
@RocketMQMessageListener(topic = "${rocketmq.consumer.analysisDify.callbackTopic}",consumerGroup = "${rocketmq.consumer.analysisDify.callbackGroup}",
instanceName = "analysisDifyCallbackMqConsumer",
consumeThreadNumber = 40,
enableMsgTrace = true)
public class AnalysisDifyCallbackMqConsumer implements RocketMQListener<MessageExt> {
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@Value("${dify.corpus.checkDccRepeat}")
private String checkDccRepeat;
@Autowired
private DiFyService diFyService;
@Autowired
private RestTemplate restTemplate;
private final ObjectMapper objectMapper = new ObjectMapper();
@Override
public void onMessage(MessageExt messageExt) {
long startTime = System.currentTimeMillis();
try {
log.info("analysisDifyCallbackMqConsumer 当前线程: {}, 线程ID: {}", Thread.currentThread().getName(), Thread.currentThread().getId());
String message = new String(messageExt.getBody());
log.info("analysisDifyCallbackMqConsumer message: " + message);
AnalysisDifyResultDTO analysisResp = JSONObject.parseObject(message, AnalysisDifyResultDTO.class);
AiAnalysisRequestLogs aiAnalysisRequestLogs = aiAnalysisRequestLogsService.queryByAiAnalysisRequestId(analysisResp.getAiAnalysisRequestId());
ResponseEntity<String> response = restTemplate.getForEntity(aiAnalysisRequestLogs.getCallbackUrl(), String.class);
if (response.getStatusCode().is2xxSuccessful()) {
log.info("analysisDifyCallbackMqConsumer aiAnalysisRequestId{},回调请求成功url:{}: " ,analysisResp.getAiAnalysisRequestId(), aiAnalysisRequestLogs.getCallbackUrl());
} else {
log.info("analysisDifyCallbackMqConsumer aiAnalysisRequestId{},回调请求失败url:{}: " ,analysisResp.getAiAnalysisRequestId(), aiAnalysisRequestLogs.getCallbackUrl());
}
log.info(" analysisDifyCallbackMqConsumer耗时{}", System.currentTimeMillis() - startTime);
} catch (Exception e) {
log.info(" analysisDifyCallbackMqConsumer mq 处理失败:{}", e.getMessage());
}
}
}

View File

@@ -0,0 +1,76 @@
package com.volvo.ai.analytic.center.mq;
import com.alibaba.fastjson.JSONObject;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.DiFyService;
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
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.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.RestController;
import java.util.concurrent.CompletableFuture;
/**
* @ClassName AnalysisDifyMqConsumer
* @Description AI解析 MQ处理
* @Author renzhen
* @Date 2025-03-04 10:18
* @Version 1.0
**/
@Slf4j
@Component
@RefreshScope
@RestController
@RocketMQMessageListener(topic = "${rocketmq.consumer.analysisDify.topic}",consumerGroup = "${rocketmq.consumer.analysisDify.group}",
instanceName = "analysisDifyMqConsumer",
consumeThreadNumber = 40,
enableMsgTrace = true)
public class AnalysisDifyMqConsumer implements RocketMQListener<MessageExt> {
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@Value("${dify.corpus.checkDccRepeat}")
private String checkDccRepeat;
@Autowired
private DiFyService diFyService;
private final ObjectMapper objectMapper = new ObjectMapper();
@Override
public void onMessage(MessageExt messageExt) {
long startTime = System.currentTimeMillis();
try {
log.info("analysisDifyMqConsumer 当前线程: {}, 线程ID: {}", Thread.currentThread().getName(), Thread.currentThread().getId());
String message = new String(messageExt.getBody());
log.info("analysisDifyMqConsumer message: " + message);
DiFyReq difyReq = JSONObject.parseObject(message, DiFyReq.class);
CompletableFuture<JSONObject> future = diFyService.asyncExecuteDifyFlow(difyReq);
future.thenAccept(result -> {
JSONObject difyRequest = JSONObject.parseObject(JSONObject.toJSONString(difyReq.getInputs()), JSONObject.class);
String aiAnalysisRequestId = difyRequest.getString("aiAnalysisRequestId");
// 处理异步结果
log.info("异步处理asyncExecuteDifyFlow完成aiAnalysisRequestId: {},处理结果:{}", aiAnalysisRequestId, result);
});
log.info("analysisDifyMqConsumer处理完成耗时{}", System.currentTimeMillis() - startTime);
} catch (Exception e) {
log.info(" dcc mq 处理失败:{}", e.getMessage());
}
}
}

View File

@@ -0,0 +1,13 @@
package com.volvo.ai.analytic.center.service;
import com.volvo.ai.analytic.center.dto.req.AnalysisReq;
import com.volvo.ai.analytic.center.dto.resp.AnalysisResp;
public interface AiAnalysisDifyService {
public boolean updateAiDifyResult(String message);
AnalysisResp aiAnalyze(AnalysisReq analysisReq);
}

View File

@@ -1,9 +0,0 @@
package com.volvo.ai.analytic.center.service;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
public interface AiDifyResultService {
public boolean updateAiDifyResult(String message);
}

View File

@@ -3,6 +3,8 @@ package com.volvo.ai.analytic.center.service;
import com.alibaba.fastjson.JSONObject;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import java.util.concurrent.CompletableFuture;
public interface DiFyService {
@@ -11,4 +13,6 @@ public interface DiFyService {
public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType, String businessData, String aiAnalysisRequestId);
public JSONObject executeDifyFlow(DiFyReq diFyReq);
public CompletableFuture<JSONObject> asyncExecuteDifyFlow(DiFyReq diFyReq);
}

View File

@@ -0,0 +1,212 @@
package com.volvo.ai.analytic.center.service.impl;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
import com.volvo.ai.analytic.center.dto.corpus.CorpusReportDTO;
import com.volvo.ai.analytic.center.dto.req.AnalysisReq;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import com.volvo.ai.analytic.center.dto.resp.AnalysisDifyResultDTO;
import com.volvo.ai.analytic.center.dto.resp.AnalysisResp;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
import com.volvo.ai.analytic.center.entity.TcBusinessType;
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
import com.volvo.ai.analytic.center.enums.CategoryEnum;
import com.volvo.ai.analytic.center.mapper.TcBusinessTypeMapper;
import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper;
import com.volvo.ai.analytic.center.service.AiAnalysisDifyService;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
import com.volvo.ai.analytic.center.utils.AiAnalysisUtils;
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.time.ZonedDateTime;
import java.time.format.DateTimeFormatter;
import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
@Slf4j
@Service
public class AiAnalysisDifyServiceImpl implements AiAnalysisDifyService {
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@Autowired
private TmTelephoneCorpusMapper tmTelephoneCorpusMapper;
@Autowired
private TmTelephoneCorpusService tmTelephoneCorpusService;
@Autowired
private TcBusinessTypeMapper tcBusinessTypeMapper;
@Value("${rocketmq.producer.analysisDify.topic}")
private String analysisDifyTopic;
@Value("${rocketmq.producer.analysisDify.callbackTopic}")
private String callbackTopic;
@Resource
private RocketMQTemplate rocketMqTemplate;
@Override
public boolean updateAiDifyResult(String message) {
if(StringUtils.isNotEmpty(message)){
AnalysisDifyResultDTO analysisResp = JSONObject.parseObject(message, AnalysisDifyResultDTO.class);
AiAnalysisRequestLogs aiAnalysisRequestLogs = new AiAnalysisRequestLogs();
aiAnalysisRequestLogs.setAiAnalysisRequestId(analysisResp.getAiAnalysisRequestId());
aiAnalysisRequestLogs.setDifyResponse(analysisResp.getDifyResponse());
aiAnalysisRequestLogs.setWorkflowRunId(analysisResp.getWorkflowRunId());
aiAnalysisRequestLogs.setWorkflowAppId(analysisResp.getWorkflowRunId());
aiAnalysisRequestLogs.setWorkUserId(analysisResp.getWorkUserId());
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(aiAnalysisRequestLogs);
// 发送 mq
sendMq(callbackTopic, analysisResp);
return true;
}
return false;
}
@Override
public AnalysisResp aiAnalyze(AnalysisReq analysisReq) {
if(null == analysisReq){
log.info("请求Ai解析对象为空");
return AnalysisResp.failed("请求Ai解析对象为空");
}
String aiAnalysisRequestId = StringUtils.isEmpty(analysisReq.getAiAnalysisRequestId())? AiAnalysisUtils.getAiAnalysisRequestId(analysisReq.getAiAnalysisRequestType()):analysisReq.getAiAnalysisRequestId();
Map<String, TcBusinessType> queryTcBusinessType = queryTcBusinessType();
TcBusinessType tcBusinessType = queryTcBusinessType.get(analysisReq.getAiAnalysisRequestType());
if(null == tcBusinessType || StringUtils.isEmpty(tcBusinessType.getWorkflowApiKey())){
log.info("接入业务类型未配置!");
return AnalysisResp.failed("接入业务类型未配置!");
}
DiFyReq diFyReq = new DiFyReq();
diFyReq.setUser(StringUtils.isEmpty(tcBusinessType.getWorkflowUser())?analysisReq.getAiAnalysisRequestType().concat("_USER"):tcBusinessType.getWorkflowUser());
diFyReq.setFlowId(tcBusinessType.getWorkflowApiKey());
// 发送mq消息
JSONObject difyRequest = JSONObject.parseObject(JSONObject.toJSONString(analysisReq.getData()), JSONObject.class);
difyRequest.put("aiAnalysisRequestId",aiAnalysisRequestId);
diFyReq.setInputs(difyRequest);
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder()
.aiAnalysisRequestId(aiAnalysisRequestId)
.businessRequest(JSONObject.toJSONString(analysisReq.getData()))
.difyAgentKey(diFyReq.getFlowId())
.difyRequest(JSON.toJSONString(diFyReq))
.aiAnalysisRequestType(analysisReq.getAiAnalysisRequestType())
.callbackUrl(analysisReq.getCallbackUrl())
.build());
sendMq(analysisDifyTopic, diFyReq);
return AnalysisResp.success(analysisReq.getData(),aiAnalysisRequestId);
}
Map<String, String> sendDccCorpus(AiAnalysisRequestLogs oldAiAnalysisRequestLogs,String difyResponse ){
CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class);
String text = JSONObject.parseObject(difyResponse).getJSONObject("outputs").getString("text");
String resultStrOne = FlowResultSplitUtil.flowOutputTextSplit(text, "任务1", "任务2");
String resultStrTwo =FlowResultSplitUtil.flowOutputTextSplit(text, "任务2", null);
if (StringUtils.isBlank(resultStrOne) || StringUtils.isBlank(resultStrTwo)){
log.info("电话语料解析为空text:{}", text);
return null;
}
List<AicorpusTelephoneDTO> dccDtoList = tmTelephoneCorpusMapper.queryTelephoneCorpusBySourceIds( Arrays.asList(corpusReportDTO.getRecordId()));
if(CollectionUtils.isNotEmpty(dccDtoList)){
AicorpusTelephoneDTO dccDto = dccDtoList.get(0);
JSONObject jsonObject = JSONObject.parseObject( dccDto.getDisplay());
ZonedDateTime zonedDateTime = ZonedDateTime.parse(jsonObject.getString("start_time"));
DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
String formattedDateStartTime = zonedDateTime.format(formatter);
Map<String, String> ltoMap = new HashMap();
ltoMap.put("analysisRecordId", oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
ltoMap.put("analysisScene", "2");
ltoMap.put("recordId", corpusReportDTO.getRecordId());
ltoMap.put("communicateDate", formattedDateStartTime);
ltoMap.put("analysisResult", resultStrOne);
ltoMap.put("analysisDetail", resultStrTwo);
// 发送MQ
log.info("send mq {}",ltoMap);
return ltoMap;
}
return null;
}
public Map<String, String> sendQiweiCorpus(AiAnalysisRequestLogs aiAnalysisRequestLogs,String difyResponse) {
CorpusReportDTO corpusReportDTO = JSONObject.parseObject(aiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class);
JSONObject execDifyFlow = JSONObject.parseObject(difyResponse);
String text = execDifyFlow.getJSONObject("outputs").getString("text");
String resultStrOne = FlowResultSplitUtil.flowOutputTextSplit(text, "任务1", "任务2");
String resultStrTwo = FlowResultSplitUtil.flowOutputTextSplit(text, "任务2", null);
if (StringUtils.isBlank(resultStrOne) || StringUtils.isBlank(resultStrTwo)){
log.info("企微语料解析为空text:{}", text);
return null;
}
Map<String, String> ltoMap = new HashMap<>();
ltoMap.put("analysisRecordId", execDifyFlow.getString("aiAnalysisRequestId"));
ltoMap.put("analysisScene", "1");
ltoMap.put("unionId", corpusReportDTO.getUnionId());
ltoMap.put("consultantId", corpusReportDTO.getUserId());
ltoMap.put("communicateDate", corpusReportDTO.getCorpusTime());
ltoMap.put("analysisResult", resultStrOne);
ltoMap.put("analysisDetail", resultStrTwo);
// 发送MQ
log.info("send mq {}", ltoMap);
return ltoMap;
}
public Map<String, TcBusinessType> queryTcBusinessType() {
LambdaQueryWrapper<TcBusinessType> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(TcBusinessType::getIsDeleted, "0");
List<TcBusinessType> tcBusinessTypeList =tcBusinessTypeMapper.selectList(queryWrapper);
return tcBusinessTypeList.stream()
.collect(Collectors.toMap(
TcBusinessType::getBusinessRequestType,
tcBusinessType -> tcBusinessType
));
}
public TcBusinessType queryTcBusinessType(String businessRequestType) {
LambdaQueryWrapper<TcBusinessType> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(TcBusinessType::getBusinessRequestType, businessRequestType);
queryWrapper.eq(TcBusinessType::getIsDeleted, "0");
TcBusinessType tcBusinessType =tcBusinessTypeMapper.selectOne(queryWrapper);
return tcBusinessType;
}
private void sendMq(String topic, Object message){
rocketMqTemplate.asyncSend(callbackTopic, MessageBuilder.withPayload(message).build(),
new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
log.info("请求AI解析发送MQ成功消息体:{}", message);
}
@Override
public void onException(Throwable e) {
log.error("请求AI解析发送MQ异常消息体:{}, 异常:", message, e);
}
}, 10000);
}
}

View File

@@ -1,141 +0,0 @@
package com.volvo.ai.analytic.center.service.impl;
import cn.hutool.core.date.DatePattern;
import cn.hutool.core.date.DateUtil;
import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
import com.volvo.ai.analytic.center.dto.corpus.CorpusReportDTO;
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
import com.volvo.ai.analytic.center.entity.TtVdqwRecord;
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
import com.volvo.ai.analytic.center.enums.CategoryEnum;
import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.AiDifyResultService;
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.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;
import java.time.ZonedDateTime;
import java.time.format.DateTimeFormatter;
import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@Slf4j
@Service
public class AiDifyResultServiceImpl implements AiDifyResultService {
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@Autowired
private TmTelephoneCorpusMapper tmTelephoneCorpusMapper;
@Autowired
private TmTelephoneCorpusService tmTelephoneCorpusService;
@Override
public boolean updateAiDifyResult(String message) {
JSONObject messageJson = JSONObject.parseObject(message);
String aiAnalysisRequestId = messageJson.getString("aiAnalysisRequestId");
String difyResponse = messageJson.getString("difyResponse");
AiAnalysisRequestLogs oldAiAnalysisRequestLogs = aiAnalysisRequestLogsService.queryByAiAnalysisRequestId(aiAnalysisRequestId);
if (null == oldAiAnalysisRequestLogs) {
log.info("根据aiAnalysisRequestId查询的log为空");
return false;
}
Map<String, String> ltoMap = null;
// 特殊处理
List<String> analysisRequestTypeList = Arrays.asList(BusinessTypeEnum.SMART_ASSISTANT.getCode(), BusinessTypeEnum.SMART_ASSISTANT_QIWEI.getCode());
if(analysisRequestTypeList.contains(oldAiAnalysisRequestLogs.getAiAnalysisRequestType())){
// 结果特殊封装
if(oldAiAnalysisRequestLogs.getAiAnalysisRequestType().equals(BusinessTypeEnum.SMART_ASSISTANT.getCode())){
ltoMap = sendDccCorpus(oldAiAnalysisRequestLogs,difyResponse);
}
if(oldAiAnalysisRequestLogs.getAiAnalysisRequestType().equals(BusinessTypeEnum.SMART_ASSISTANT_QIWEI.getCode())){
ltoMap = sendQiweiCorpus(oldAiAnalysisRequestLogs,difyResponse);
}
}
AiAnalysisRequestLogs aiAnalysisRequestLogs = new AiAnalysisRequestLogs();
aiAnalysisRequestLogs.setAiAnalysisRequestId(aiAnalysisRequestId);
aiAnalysisRequestLogs.setDifyResponse(difyResponse);
aiAnalysisRequestLogs.setBusinessResponse(JSONObject.toJSONString(ltoMap));
aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(aiAnalysisRequestLogs);
// aiAnalysisRequestLogsService.saveAiAnalysisRequestLogs(AiAnalysisRequestLogs.builder().aiAnalysisRequestId(aiAnalysisRequestId).businessResponse(JSONObject.toJSONString(ltoMap)).build());
tmTelephoneCorpusService.sendMq( CategoryEnum.ENTERPRISE_WECHAT.getCode(), JSONObject.toJSONString(ltoMap));
return false;
}
Map<String, String> sendDccCorpus(AiAnalysisRequestLogs oldAiAnalysisRequestLogs,String difyResponse ){
CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class);
String text = JSONObject.parseObject(difyResponse).getJSONObject("outputs").getString("text");
String resultStrOne = FlowResultSplitUtil.flowOutputTextSplit(text, "任务1", "任务2");
String resultStrTwo =FlowResultSplitUtil.flowOutputTextSplit(text, "任务2", null);
if (StringUtils.isBlank(resultStrOne) || StringUtils.isBlank(resultStrTwo)){
log.info("电话语料解析为空text:{}", text);
return null;
}
List<AicorpusTelephoneDTO> dccDtoList = tmTelephoneCorpusMapper.queryTelephoneCorpusBySourceIds( Arrays.asList(corpusReportDTO.getRecordId()));
if(CollectionUtils.isNotEmpty(dccDtoList)){
AicorpusTelephoneDTO dccDto = dccDtoList.get(0);
JSONObject jsonObject = JSONObject.parseObject( dccDto.getDisplay());
ZonedDateTime zonedDateTime = ZonedDateTime.parse(jsonObject.getString("start_time"));
DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
String formattedDateStartTime = zonedDateTime.format(formatter);
Map<String, String> ltoMap = new HashMap();
ltoMap.put("analysisRecordId", oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
ltoMap.put("analysisScene", "2");
ltoMap.put("recordId", corpusReportDTO.getRecordId());
ltoMap.put("communicateDate", formattedDateStartTime);
ltoMap.put("analysisResult", resultStrOne);
ltoMap.put("analysisDetail", resultStrTwo);
// 发送MQ
log.info("send mq {}",ltoMap);
return ltoMap;
}
return null;
}
public Map<String, String> sendQiweiCorpus(AiAnalysisRequestLogs aiAnalysisRequestLogs,String difyResponse) {
CorpusReportDTO corpusReportDTO = JSONObject.parseObject(aiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class);
JSONObject execDifyFlow = JSONObject.parseObject(difyResponse);
String text = execDifyFlow.getJSONObject("outputs").getString("text");
String resultStrOne = FlowResultSplitUtil.flowOutputTextSplit(text, "任务1", "任务2");
String resultStrTwo = FlowResultSplitUtil.flowOutputTextSplit(text, "任务2", null);
if (StringUtils.isBlank(resultStrOne) || StringUtils.isBlank(resultStrTwo)){
log.info("企微语料解析为空text:{}", text);
return null;
}
Map<String, String> ltoMap = new HashMap<>();
ltoMap.put("analysisRecordId", execDifyFlow.getString("aiAnalysisRequestId"));
ltoMap.put("analysisScene", "1");
ltoMap.put("unionId", corpusReportDTO.getUnionId());
ltoMap.put("consultantId", corpusReportDTO.getUserId());
ltoMap.put("communicateDate", corpusReportDTO.getCorpusTime());
ltoMap.put("analysisResult", resultStrOne);
ltoMap.put("analysisDetail", resultStrTwo);
// 发送MQ
log.info("send mq {}", ltoMap);
return ltoMap;
}
}

View File

@@ -17,6 +17,7 @@ import org.springframework.stereotype.Service;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
@Slf4j
@Service
@@ -105,4 +106,13 @@ public class DiFyServiceImpl implements DiFyService{
return data;
}
@Override
public CompletableFuture<JSONObject> asyncExecuteDifyFlow(DiFyReq diFyReq) {
Map<String, Object> map = new HashMap<>();
map.put("inputs",diFyReq.getInputs());
map.put("user",diFyReq.getUser());
return CompletableFuture.supplyAsync(() -> diFyFeign.runWorkflows("Bearer "+diFyReq.getFlowId(),map));
}
}