diff --git a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/req/AnalysisReq.java b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/req/AnalysisReq.java new file mode 100644 index 0000000..0c0b811 --- /dev/null +++ b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/req/AnalysisReq.java @@ -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; + +} diff --git a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/resp/AnalysisDifyResultDTO.java b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/resp/AnalysisDifyResultDTO.java new file mode 100644 index 0000000..b47b7ea --- /dev/null +++ b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/resp/AnalysisDifyResultDTO.java @@ -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; +} diff --git a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/resp/AnalysisResp.java b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/resp/AnalysisResp.java new file mode 100644 index 0000000..21a11ab --- /dev/null +++ b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/resp/AnalysisResp.java @@ -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 { + + // 响应 + private T data; + + // 分析中心唯一ID + private String aiAnalysisRequestId; + + private int code; + private String msg; + + public static AnalysisResp success(String message) { + return (AnalysisResp) analysisResp((Object)null, CommonConstants.SUCCESS, message); + } + public static AnalysisResp success(T data, String aiAnalysisRequestId) { + AnalysisResp analysisResp = new AnalysisResp(); + analysisResp.setAiAnalysisRequestId(aiAnalysisRequestId); + analysisResp.setData(data); + analysisResp.setMsg("ok"); + analysisResp.setCode(CommonConstants.SUCCESS); + return analysisResp; + } + public static AnalysisResp failed(String message) { + return (AnalysisResp) analysisResp((Object)null, CommonConstants.FAIL, message); + } + + private static AnalysisResp analysisResp(T data, int code, String msg) { + AnalysisResp analysisResp = new AnalysisResp(); + analysisResp.setCode(code); + analysisResp.setData(data); + analysisResp.setMsg(msg); + return analysisResp; + } + + private static AnalysisResp analysisResp(T data, String aiAnalysisRequestId) { + AnalysisResp analysisResp = new AnalysisResp(); + analysisResp.setAiAnalysisRequestId(aiAnalysisRequestId); + analysisResp.setData(data); + analysisResp.setMsg("ok"); + analysisResp.setCode(CommonConstants.SUCCESS); + return analysisResp; + } +} diff --git a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/entity/AiAnalysisRequestLogs.java b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/entity/AiAnalysisRequestLogs.java index 9ef50ac..69402b7 100644 --- a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/entity/AiAnalysisRequestLogs.java +++ b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/entity/AiAnalysisRequestLogs.java @@ -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; diff --git a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/entity/TcBusinessType.java b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/entity/TcBusinessType.java new file mode 100644 index 0000000..fd36142 --- /dev/null +++ b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/entity/TcBusinessType.java @@ -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; +} diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/AiDifyResultController.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/AiAnalysisDifyController.java similarity index 62% rename from ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/AiDifyResultController.java rename to ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/AiAnalysisDifyController.java index 393ee94..d2b908e 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/AiDifyResultController.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/AiAnalysisDifyController.java @@ -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 aiAnalyze(@RequestBody AnalysisReq analysisReq) { + log.info("aiAnalyze data: {}",analysisReq); + return AnalysisResp.success("ok"); + } + } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TcBusinessTypeMapper.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TcBusinessTypeMapper.java new file mode 100644 index 0000000..d381fdd --- /dev/null +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/TcBusinessTypeMapper.java @@ -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 { + + +} diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/AnalysisDifyCallbackMqConsumer.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/AnalysisDifyCallbackMqConsumer.java new file mode 100644 index 0000000..35dbf3e --- /dev/null +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/AnalysisDifyCallbackMqConsumer.java @@ -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 { + + + @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 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()); + } + } + +} + diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/AnalysisDifyMqConsumer.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/AnalysisDifyMqConsumer.java new file mode 100644 index 0000000..74be96f --- /dev/null +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/AnalysisDifyMqConsumer.java @@ -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 { + + + @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 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()); + } + } + +} + diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiAnalysisDifyService.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiAnalysisDifyService.java new file mode 100644 index 0000000..2a12f06 --- /dev/null +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiAnalysisDifyService.java @@ -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); + +} diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiDifyResultService.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiDifyResultService.java deleted file mode 100644 index 677011e..0000000 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiDifyResultService.java +++ /dev/null @@ -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); - -} diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/DiFyService.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/DiFyService.java index 7159e60..3bfba4a 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/DiFyService.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/DiFyService.java @@ -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 asyncExecuteDifyFlow(DiFyReq diFyReq); + } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisDifyServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisDifyServiceImpl.java new file mode 100644 index 0000000..6a53fa2 --- /dev/null +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisDifyServiceImpl.java @@ -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 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 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 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 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 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 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 queryTcBusinessType() { + LambdaQueryWrapper queryWrapper = new LambdaQueryWrapper<>(); + queryWrapper.eq(TcBusinessType::getIsDeleted, "0"); + List tcBusinessTypeList =tcBusinessTypeMapper.selectList(queryWrapper); + return tcBusinessTypeList.stream() + .collect(Collectors.toMap( + TcBusinessType::getBusinessRequestType, + tcBusinessType -> tcBusinessType + )); + } + + public TcBusinessType queryTcBusinessType(String businessRequestType) { + LambdaQueryWrapper 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); + } + + +} diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiDifyResultServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiDifyResultServiceImpl.java deleted file mode 100644 index acca000..0000000 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiDifyResultServiceImpl.java +++ /dev/null @@ -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 ltoMap = null; - // 特殊处理 - List 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 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 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 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 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 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; - } - - -} diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/DiFyServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/DiFyServiceImpl.java index d5b49d2..661cfc5 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/DiFyServiceImpl.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/DiFyServiceImpl.java @@ -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 asyncExecuteDifyFlow(DiFyReq diFyReq) { + Map map = new HashMap<>(); + map.put("inputs",diFyReq.getInputs()); + map.put("user",diFyReq.getUser()); + return CompletableFuture.supplyAsync(() -> diFyFeign.runWorkflows("Bearer "+diFyReq.getFlowId(),map)); + + } + }