Merge branch 'uat' into dev

# Conflicts:
#	ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/DiFyServiceImpl.java
This commit is contained in:
lxu75
2025-03-10 20:27:48 +08:00
15 changed files with 389 additions and 3 deletions

View File

@@ -0,0 +1,41 @@
package com.volvo.ai.analytic.center.dto.req;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class CarModelReqDTO {
//@ApiModelProperty("燃料类型")
private Integer fuelType;
//@ApiModelProperty("是否包含主图 0 不包含 1 包含 默认不包含")
private Integer isContainMainImage;
//@ApiModelProperty("车型类型")
private Integer modelType;
//@ApiModelProperty("车型代码")
private String modelCode;
//@ApiModelProperty("年款")
private String modelYear;
//@ApiModelProperty("是否有效 10041001 有效 10041002 无效")
private Integer isValid;
//@ApiModelProperty("是否是直售 10041001 是 10041002 否")
private Integer isDirectModel;
//@ApiModelProperty("是否在售(10041001:是 10041002: 否)")
private Integer onSale;
//@ApiModelProperty("车系ID")
private Integer seriesId;
//@ApiModelProperty("车型名称模糊查询")
private String modelName;
}

View File

@@ -0,0 +1,65 @@
package com.volvo.ai.analytic.center.dto.resp;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class CarModelRespDTO {
// ("燃料类型")
private Integer fuelType;
// ("公司代码")
private String companyCode;
// @ApiModelProperty("主图")
private String imageUrl;
// @ApiModelProperty("车型代码")
private String modelCode;
// @ApiModelProperty("车型描述(中文)")
private String modelDescriptionChinese;
// ("是否直售车型:10041001:是10041002:否")
private Integer isDirectModel;
// ("车型描述(英文)")
private String modelDescriptionEnglish;
// ("车型名称")
private String modelName;
// ("车型名称英文")
private String modelNameEn;
// ("车型类型")
private Integer modelType;
// ("年款")
private String modelYear;
// ("车系代码")
private String seriesCode;
// ("车型id")
private Integer id;
// ("排序")
private Integer sort;
// ("是否有效 10041001 有效 10041002 无效")
private Integer isValid;
// ("是否在售(10041001:是 10041002: 否)")
private Integer onSale;
//@ApiModelProperty("车型名称 沃世界名称")
private String modelNameC;
//@ApiModelProperty("备注")
private String modelRemark;
}

View File

@@ -0,0 +1,19 @@
package com.volvo.ai.analytic.center.dto.resp;
import lombok.Data;
import java.io.Serializable;
@Data
public class ResultDTO<T> implements Serializable {
private static final long serialVersionUID = -1179271389084311472L;
private String returnCode;
private String returnMessage;
private String errMsg;
private T data;
}

View File

@@ -31,6 +31,9 @@ public class AiAnalysisRequestLogs extends BaseEntity {
@TableField("business_response") @TableField("business_response")
private String businessResponse; // JSON 字符串 private String businessResponse; // JSON 字符串
@TableField("dify_response")
private String difyResponse;
@TableField("quest_status") @TableField("quest_status")
private String questStatus; private String questStatus;

View File

@@ -14,7 +14,9 @@ import java.util.Objects;
@Getter @Getter
@ToString @ToString
public enum BizEnum { public enum BizEnum {
SUCCESS(200, "操作成功"),
FAIL(500, "操作失败"),
BAD_REQUEST(4000, "参数不合法"), BAD_REQUEST(4000, "参数不合法"),
METHOD_NOT_ALLOWED(4001, "方法不允许"), METHOD_NOT_ALLOWED(4001, "方法不允许"),
MISS_PARAM(4002, "参数缺失"), MISS_PARAM(4002, "参数缺失"),

View File

@@ -5,7 +5,8 @@ import lombok.Getter;
@Getter @Getter
public enum BusinessTypeEnum { public enum BusinessTypeEnum {
COMMUNITYTARGET("CommunityTarget", "社区舆情分析") COMMUNITYTARGET("CommunityTarget", "社区舆情分析"),
SMART_ASSISTANT("SMART_ASSISTANT", "智能助手")
; ;
private String code; private String code;

View File

@@ -0,0 +1,18 @@
package com.volvo.ai.analytic.center.feign;
import com.volvo.ai.analytic.center.dto.req.CarModelReqDTO;
import com.volvo.ai.analytic.center.dto.resp.CarModelRespDTO;
import com.volvo.ai.analytic.center.dto.resp.ResultDTO;
import org.springframework.cloud.openfeign.FeignClient;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import java.util.List;
@FeignClient(url = "${mse-in.url.domain}", contextId = "basic-data-client", value = "basicDataClient", path = "/2b/2b-basicdata-service")
public interface RemoteCarModelClient {
@PostMapping(name = "查询车型", path = "/model/list")
ResultDTO<List<CarModelRespDTO>> queryCarModelList(@RequestBody CarModelReqDTO carModelReqDTO);
}

View File

@@ -10,6 +10,8 @@ import org.springframework.beans.BeanUtils;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import java.time.LocalDateTime; import java.time.LocalDateTime;
@@ -22,6 +24,7 @@ import java.time.LocalDateTime;
**/ **/
@Slf4j @Slf4j
@Component @Component
@RestController
public class CorpusProcessKafkaConsumer { public class CorpusProcessKafkaConsumer {
@Autowired @Autowired
@@ -29,6 +32,7 @@ public class CorpusProcessKafkaConsumer {
private final ObjectMapper objectMapper = new ObjectMapper(); private final ObjectMapper objectMapper = new ObjectMapper();
@GetMapping("corpusProcessKafkaConsumer")
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") @KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
public void listen(String message) { public void listen(String message) {
try { try {
@@ -44,6 +48,8 @@ public class CorpusProcessKafkaConsumer {
tmTelephoneCorpus.setCreateTime(LocalDateTime.now()); tmTelephoneCorpus.setCreateTime(LocalDateTime.now());
tmTelephoneCorpusService.saveTelephoneCorpus(tmTelephoneCorpus); tmTelephoneCorpusService.saveTelephoneCorpus(tmTelephoneCorpus);
tmTelephoneCorpusService.runDify(aicorpusTelephone);
// 在这里可以添加对解析后的对象的进一步处理逻辑 // 在这里可以添加对解析后的对象的进一步处理逻辑
} catch (Exception e) { } catch (Exception e) {
e.printStackTrace(); e.printStackTrace();

View File

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

View File

@@ -1,9 +1,13 @@
package com.volvo.ai.analytic.center.service; 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.dto.req.DiFyReq;
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
public interface DiFyService { public interface DiFyService {
public Object getDiFyObject(DiFyReq diFyReq); public Object getDiFyObject(DiFyReq diFyReq);
public JSONObject execDifyFlow(DiFyReq diFyReq, String businessType);
} }

View File

@@ -1,6 +1,7 @@
package com.volvo.ai.analytic.center.service; package com.volvo.ai.analytic.center.service;
import com.baomidou.mybatisplus.extension.service.IService; 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 com.volvo.ai.analytic.center.entity.TmTelephoneCorpus;
import java.util.Map; import java.util.Map;
@@ -14,4 +15,7 @@ public interface TmTelephoneCorpusService extends IService<TmTelephoneCorpus> {
void saveTelephoneCorpus(TmTelephoneCorpus tmTelephoneCorpus); void saveTelephoneCorpus(TmTelephoneCorpus tmTelephoneCorpus);
void runDify(AicorpusTelephoneDTO aicorpusTelephone);
} }

View File

@@ -1,13 +1,31 @@
package com.volvo.ai.analytic.center.service.impl; 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.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors; import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
import com.volvo.ai.analytic.center.mapper.AiAnalysisErrorsMapper; import com.volvo.ai.analytic.center.mapper.AiAnalysisErrorsMapper;
import com.volvo.ai.analytic.center.service.AiAnalysisErrorsService; import com.volvo.ai.analytic.center.service.AiAnalysisErrorsService;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
@Slf4j @Slf4j
@Service @Service
public class AiAnalysisErrorsServiceImpl extends ServiceImpl<AiAnalysisErrorsMapper, AiAnalysisErrors> implements AiAnalysisErrorsService { 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;
}
}
} }

View File

@@ -3,12 +3,20 @@ package com.volvo.ai.analytic.center.service.impl;
import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.JSONObject;
import com.volvo.ai.analytic.center.dto.req.DiFyReq; 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.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.service.DiFyService;
import com.volvo.ai.analytic.center.utils.AiAnalysisUtils;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import java.util.Collections;
import java.util.HashMap; import java.util.HashMap;
import java.util.Map; import java.util.Map;
@@ -19,13 +27,19 @@ public class DiFyServiceImpl implements DiFyService{
@Autowired @Autowired
private DiFyFeign diFyFeign; private DiFyFeign diFyFeign;
@Autowired
private AiAnalysisRequestLogsMapper aiAnalysisRequestLogsMapper;
@Autowired
private AiAnalysisErrorsService aiAnalysisErrorsService;
@Override @Override
public Object getDiFyObject(DiFyReq diFyReq) { public Object getDiFyObject(DiFyReq diFyReq) {
Map<String, Object> map = new HashMap<>(); Map<String, Object> map = new HashMap<>();
map.put("inputs",diFyReq.getInputs()); map.put("inputs",diFyReq.getInputs());
map.put("response_mode","blocking"); map.put("response_mode","blocking");
map.put("user",diFyReq.getUser()); map.put("user",diFyReq.getUser());
log.info("请求DiFy入参:{}", JSON.toJSONString(map)); log.info("请求DiFy入参:{}",JSON.toJSONString(map));
JSONObject difyResult = diFyFeign.runWorkflows("Bearer "+diFyReq.getFlowId(),map); JSONObject difyResult = diFyFeign.runWorkflows("Bearer "+diFyReq.getFlowId(),map);
log.info("请求DiFy响应结果:{}",difyResult.toJSONString()); log.info("请求DiFy响应结果:{}",difyResult.toJSONString());
JSONObject data = difyResult.getJSONObject("data"); JSONObject data = difyResult.getJSONObject("data");
@@ -35,4 +49,40 @@ public class DiFyServiceImpl implements DiFyService{
} }
return ""; 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(JSONObject.toJSONString(""))
.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;
}
}
} }

View File

@@ -1,14 +1,37 @@
package com.volvo.ai.analytic.center.service.impl; 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.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.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.mapper.TmTelephoneCorpusMapper;
import com.volvo.ai.analytic.center.service.DiFyService;
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService; import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil;
import lombok.extern.slf4j.Slf4j; 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.stereotype.Service;
import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
import java.util.*;
import java.util.stream.Collectors;
/** /**
* @description 电话语料表-同步表 * @description 电话语料表-同步表
@@ -20,9 +43,112 @@ import org.springframework.transaction.annotation.Transactional;
public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusMapper, TmTelephoneCorpus> implements TmTelephoneCorpusService { 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 @Override
@Transactional @Transactional
public void saveTelephoneCorpus(TmTelephoneCorpus tmTelephoneCorpus) { public void saveTelephoneCorpus(TmTelephoneCorpus tmTelephoneCorpus) {
this.save(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;
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, String> 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.getString("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(){
try {
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;
}
} catch (Exception e) {
log.error("getCarModelList Exception: {}", e);
}
return "C40 RECHARGE、EM90、EX30、S60、S90、V60、V90、XC40、XC40 RECHARGE、XC60、XC90";
}
} }

View File

@@ -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();
}
}