提交验证kafka方法
This commit is contained in:
@@ -2,7 +2,9 @@ package com.volvo.ai.analytic.center.dto.corpus;
|
||||
|
||||
|
||||
import com.fasterxml.jackson.annotation.JsonRawValue;
|
||||
import lombok.Data;
|
||||
|
||||
@Data
|
||||
public class AicorpusTelephoneDTO {
|
||||
private String appid;
|
||||
private String audioUrl;
|
||||
@@ -17,6 +19,8 @@ public class AicorpusTelephoneDTO {
|
||||
private int platformType;
|
||||
private String displayUrl;
|
||||
|
||||
private String aiAnalysisRequestId;
|
||||
|
||||
// Getters and Setters
|
||||
public String getAppid() {
|
||||
return appid;
|
||||
|
||||
@@ -4,6 +4,7 @@ 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.PageDto;
|
||||
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.DiFyReq;
|
||||
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
|
||||
@@ -11,6 +12,7 @@ import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
|
||||
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
|
||||
import com.volvo.ai.analytic.center.enums.CategoryEnum;
|
||||
import com.volvo.ai.analytic.center.mapper.AiAnalysisRequestLogsMapper;
|
||||
import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper;
|
||||
import com.volvo.ai.analytic.center.service.AiAnalysisErrorsService;
|
||||
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
|
||||
import com.volvo.ai.analytic.center.service.DiFyService;
|
||||
@@ -20,14 +22,18 @@ import com.xxl.job.core.handler.annotation.XxlJob;
|
||||
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.Component;
|
||||
import org.springframework.web.bind.annotation.PostMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
|
||||
import javax.annotation.Resource;
|
||||
import java.util.Arrays;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -59,6 +65,13 @@ public class CorpusFailJob {
|
||||
@Autowired
|
||||
private TmTelephoneCorpusService tmTelephoneCorpusService;
|
||||
|
||||
@Value("${rocketmq.producer.corpus.dcctopic}")
|
||||
private String dccMqTipic;
|
||||
|
||||
@Autowired
|
||||
private TmTelephoneCorpusMapper tmTelephoneCorpusMapper;
|
||||
|
||||
|
||||
/**
|
||||
* 企微语料处理
|
||||
*/
|
||||
@@ -66,7 +79,7 @@ public class CorpusFailJob {
|
||||
@PostMapping("corpusFailTask")
|
||||
public void corpusFailTask() {
|
||||
|
||||
log.info("语料解析失败重试处理");
|
||||
log.info(" 解析失败重试处理");
|
||||
Integer total = aiAnalysisErrorsService.queryCountAnalysisErrorList(BusinessTypeEnum.SMART_ASSISTANT.getCode());
|
||||
log.info("语料解析失败重试处理数据量:{}", total);
|
||||
int totalPages = PageDto.getTotalPages(total, pageSize);
|
||||
@@ -88,45 +101,64 @@ public class CorpusFailJob {
|
||||
if (null != oldAiAnalysisRequestLogs) {
|
||||
DiFyReq diFyReq = JSONObject.parseObject(oldAiAnalysisRequestLogs.getDifyRequest(), DiFyReq.class);
|
||||
CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class);
|
||||
Map<String, String> ltoMap = new HashMap<>();
|
||||
ltoMap.put("analysisRecordId", oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
|
||||
ltoMap.put("analysisScene", corpusReportDTO.getAnalysisScene() + "");
|
||||
String tag = "";
|
||||
if (null != corpusReportDTO && corpusReportDTO.getAnalysisScene() == 1) { //企微
|
||||
ltoMap.put("unionId", corpusReportDTO.getUnionId());
|
||||
ltoMap.put("consultantId", corpusReportDTO.getUserId());
|
||||
tag = CategoryEnum.ENTERPRISE_WECHAT.getCode();
|
||||
|
||||
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyReq);
|
||||
log.info(" corpusFailTask 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);
|
||||
if (StringUtils.isBlank(resultStrOne) || StringUtils.isBlank(resultStrTwo)) {
|
||||
log.info(" corpusFailTask 企微语料解析为空,text:{}", text);
|
||||
return;
|
||||
}
|
||||
Map<String, String> ltoMap = new HashMap<>();
|
||||
ltoMap.put("analysisRecordId", oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
|
||||
ltoMap.put("analysisScene", corpusReportDTO.getAnalysisScene() + "");
|
||||
String tag = "";
|
||||
if (null != corpusReportDTO && corpusReportDTO.getAnalysisScene() == 1) { //企微
|
||||
ltoMap.put("unionId", corpusReportDTO.getUnionId());
|
||||
ltoMap.put("consultantId", corpusReportDTO.getUserId());
|
||||
tag = CategoryEnum.ENTERPRISE_WECHAT.getCode();
|
||||
} else {
|
||||
ltoMap.put("recordId", corpusReportDTO.getUserId());
|
||||
tag = CategoryEnum.PHONE_VOICE.getCode();
|
||||
}
|
||||
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyReq);
|
||||
log.info(" corpusFailTask 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);
|
||||
if (StringUtils.isBlank(resultStrOne) || StringUtils.isBlank(resultStrTwo)) {
|
||||
log.info(" corpusFailTask 企微语料解析为空,text:{}", text);
|
||||
return;
|
||||
}
|
||||
|
||||
ltoMap.put("communicateDate", corpusReportDTO.getCorpusTime());
|
||||
ltoMap.put("analysisResult", resultStrOne.replace("#", ""));
|
||||
ltoMap.put("analysisDetail", resultStrTwo);
|
||||
// 发送MQ
|
||||
log.info("corpusFailTask send mq {}", ltoMap);
|
||||
tmTelephoneCorpusService.sendMq(tag, JSONObject.toJSONString(ltoMap));
|
||||
// 保存报告
|
||||
oldAiAnalysisRequestLogs.setBusinessResponse(JSONObject.toJSONString(ltoMap));
|
||||
|
||||
ltoMap.put("communicateDate", corpusReportDTO.getCorpusTime());
|
||||
ltoMap.put("analysisResult", resultStrOne.replace("#", ""));
|
||||
ltoMap.put("analysisDetail", resultStrTwo);
|
||||
// 发送MQ
|
||||
log.info("corpusFailTask send mq {}", ltoMap);
|
||||
tmTelephoneCorpusService.sendMq(tag, JSONObject.toJSONString(ltoMap));
|
||||
// 保存报告
|
||||
oldAiAnalysisRequestLogs.setBusinessResponse(JSONObject.toJSONString(ltoMap));
|
||||
|
||||
oldAiAnalysisRequestLogs.setDifyResponse(JSON.toJSONString(execDifyFlow));
|
||||
updateDiFyRequest(oldAiAnalysisRequestLogs, oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
|
||||
aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1);
|
||||
aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("1");
|
||||
updateAiAnalysisErrors(aiAnalysisErrors, oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
|
||||
|
||||
}
|
||||
oldAiAnalysisRequestLogs.setDifyResponse(JSON.toJSONString(execDifyFlow));
|
||||
updateDiFyRequest(oldAiAnalysisRequestLogs, oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
|
||||
aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1);
|
||||
aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("1");
|
||||
updateAiAnalysisErrors(aiAnalysisErrors, oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
|
||||
} else {
|
||||
List<AicorpusTelephoneDTO> dccDtoList = tmTelephoneCorpusMapper.queryTelephoneCorpusBySourceIds( Arrays.asList(corpusReportDTO.getRecordId()));
|
||||
if(CollectionUtils.isNotEmpty(dccDtoList)){
|
||||
AicorpusTelephoneDTO dccDto = dccDtoList.get(0);
|
||||
dccDto.setAiAnalysisRequestId(oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
|
||||
String message = JSONObject.toJSONString(dccDto);
|
||||
rocketMqTemplate.asyncSend(dccMqTipic, MessageBuilder.withPayload(message).build(),
|
||||
new SendCallback() {
|
||||
@Override
|
||||
public void onSuccess(SendResult sendResult) {
|
||||
log.info("dcc 失败补偿 发送MQ成功 消息体:{}", message);
|
||||
}
|
||||
@Override
|
||||
public void onException(Throwable e) {
|
||||
log.error("dcc 失败补偿 送MQ异常 消息体:{}, 异常:", message, e);
|
||||
}
|
||||
}, 10000);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
} catch (Exception e) {
|
||||
log.info("语料解析失败补偿异常:{}", e);
|
||||
|
||||
@@ -3,7 +3,10 @@ package com.volvo.ai.analytic.center.mapper;
|
||||
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
|
||||
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
|
||||
import org.apache.ibatis.annotations.Mapper;
|
||||
import org.apache.ibatis.annotations.Param;
|
||||
|
||||
@Mapper
|
||||
public interface AiAnalysisRequestLogsMapper extends BaseMapper<AiAnalysisRequestLogs> {
|
||||
|
||||
public AiAnalysisRequestLogs queryAiAnalysisRequestLogsByBusinessReponse(@Param("sourceId") String sourceId);
|
||||
}
|
||||
|
||||
@@ -5,12 +5,16 @@ 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.corpus.DisplayDTO;
|
||||
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
|
||||
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
|
||||
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;
|
||||
|
||||
@@ -24,6 +28,7 @@ import org.springframework.web.bind.annotation.RestController;
|
||||
|
||||
@Slf4j
|
||||
@Component
|
||||
@RefreshScope
|
||||
@RestController
|
||||
@RocketMQMessageListener(consumerGroup = "${rocketmq.consumer.corpus.dcctopicgroup}",
|
||||
topic = "${rocketmq.consumer.corpus.dcctopic}",
|
||||
@@ -35,6 +40,13 @@ public class CorpusDccMqConsumer implements RocketMQListener<MessageExt> {
|
||||
@Autowired
|
||||
private TmTelephoneCorpusService tmTelephoneCorpusService;
|
||||
|
||||
@Autowired
|
||||
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
|
||||
|
||||
@Value("${dify.corpus.checkDccRepeat}")
|
||||
private String checkDccRepeat;
|
||||
|
||||
|
||||
private final ObjectMapper objectMapper = new ObjectMapper();
|
||||
|
||||
@Override
|
||||
@@ -45,10 +57,16 @@ public class CorpusDccMqConsumer implements RocketMQListener<MessageExt> {
|
||||
String message = new String(messageExt.getBody());
|
||||
log.info("dcc_mq message: " + message);
|
||||
AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class);
|
||||
DisplayDTO display = objectMapper.readValue(aicorpusTelephone.getDisplay(), DisplayDTO.class);
|
||||
log.info("corpusDccMqProducer display getSegments: {}", display.getSegments());
|
||||
tmTelephoneCorpusService.runTelephoneCorpusDify(aicorpusTelephone);
|
||||
log.info("corpusDccMqProducer mq 处理完成: {}", aicorpusTelephone.getSourceId());
|
||||
AiAnalysisRequestLogs aiAnalysisRequestLogs = null;
|
||||
if(checkDccRepeat.equals("true")){
|
||||
aiAnalysisRequestLogs = aiAnalysisRequestLogsService.queryAiAnalysisRequestLogsByBusinessReponse(aicorpusTelephone.getSourceId());
|
||||
log.info("corpusDccMqProducer 消费校验是否已经发送: {}", aiAnalysisRequestLogs);
|
||||
}
|
||||
|
||||
if(null == aiAnalysisRequestLogs){
|
||||
tmTelephoneCorpusService.runTelephoneCorpusDify(aicorpusTelephone);
|
||||
log.info("corpusDccMqProducer mq 处理完成: {}", aicorpusTelephone.getSourceId());
|
||||
}
|
||||
log.info("dcc_mq 处理完成,耗时:{}", System.currentTimeMillis() - startTime);
|
||||
} catch (JsonProcessingException e) {
|
||||
log.info(" dcc mq 处理失败:{}", e.getMessage());
|
||||
|
||||
@@ -9,6 +9,7 @@ import com.volvo.ai.analytic.center.entity.TmTelephoneCorpus;
|
||||
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.commons.collections.CollectionUtils;
|
||||
import org.apache.kafka.clients.consumer.OffsetAndTimestamp;
|
||||
import org.apache.rocketmq.client.producer.SendCallback;
|
||||
import org.apache.rocketmq.client.producer.SendResult;
|
||||
import org.apache.rocketmq.spring.core.RocketMQTemplate;
|
||||
@@ -17,16 +18,17 @@ 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.kafka.annotation.KafkaListener;
|
||||
import org.springframework.kafka.annotation.PartitionOffset;
|
||||
import org.springframework.kafka.annotation.TopicPartition;
|
||||
import org.springframework.kafka.listener.ConsumerSeekAware;
|
||||
import org.springframework.messaging.support.MessageBuilder;
|
||||
import org.springframework.stereotype.Component;
|
||||
import org.springframework.web.bind.annotation.PostMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
|
||||
import javax.annotation.Resource;
|
||||
import java.time.LocalDateTime;
|
||||
import java.util.Collections;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
|
||||
@@ -42,7 +44,7 @@ import java.util.concurrent.Executors;
|
||||
@Component
|
||||
@RestController
|
||||
@RefreshScope
|
||||
public class CorpusProcessKafkaProducer {
|
||||
public class CorpusProcessKafkaProducer implements ConsumerSeekAware {
|
||||
|
||||
@Autowired
|
||||
private TmTelephoneCorpusService tmTelephoneCorpusService;
|
||||
@@ -55,28 +57,29 @@ public class CorpusProcessKafkaProducer {
|
||||
@Resource
|
||||
private RocketMQTemplate rocketMqTemplate;
|
||||
|
||||
@PostMapping("corpusProcessKafkaConsumer")
|
||||
@KafkaListener(
|
||||
topicPartitions = @TopicPartition(
|
||||
topic = "${spring.kafka.topic}",
|
||||
partitions = {"0", "1","2", "3","4", "5","6", "7","8", "9", "10", "11"},
|
||||
partitionOffsets = {
|
||||
@PartitionOffset(partition = "0", initialOffset = "2792520"),
|
||||
@PartitionOffset(partition = "1", initialOffset = "2596153"),
|
||||
@PartitionOffset(partition = "2", initialOffset = "2536889"),
|
||||
@PartitionOffset(partition = "3", initialOffset = "2782173"),
|
||||
@PartitionOffset(partition = "4", initialOffset = "2616677"),
|
||||
@PartitionOffset(partition = "5", initialOffset = "2529585"),
|
||||
@PartitionOffset(partition = "6", initialOffset = "2761950"),
|
||||
@PartitionOffset(partition = "7", initialOffset = "2581957"),
|
||||
@PartitionOffset(partition = "8", initialOffset = "2530225"),
|
||||
@PartitionOffset(partition = "9", initialOffset = "2803179"),
|
||||
@PartitionOffset(partition = "10", initialOffset = "2599129"),
|
||||
@PartitionOffset(partition = "11", initialOffset = "2546277")
|
||||
}
|
||||
),
|
||||
groupId = "${spring.kafka.group}"
|
||||
)
|
||||
// @PostMapping("corpusProcessKafkaConsumer")
|
||||
// @KafkaListener(
|
||||
// topicPartitions = @TopicPartition(
|
||||
// topic = "${spring.kafka.topic}",
|
||||
// partitions = {"0", "1","2", "3","4", "5","6", "7","8", "9", "10", "11"},
|
||||
// partitionOffsets = {
|
||||
// @PartitionOffset(partition = "0", initialOffset = "2792520"),
|
||||
// @PartitionOffset(partition = "1", initialOffset = "2596153"),
|
||||
// @PartitionOffset(partition = "2", initialOffset = "2536889"),
|
||||
// @PartitionOffset(partition = "3", initialOffset = "2782173"),
|
||||
// @PartitionOffset(partition = "4", initialOffset = "2616677"),
|
||||
// @PartitionOffset(partition = "5", initialOffset = "2529585"),
|
||||
// @PartitionOffset(partition = "6", initialOffset = "2761950"),
|
||||
// @PartitionOffset(partition = "7", initialOffset = "2581957"),
|
||||
// @PartitionOffset(partition = "8", initialOffset = "2530225"),
|
||||
// @PartitionOffset(partition = "9", initialOffset = "2803179"),
|
||||
// @PartitionOffset(partition = "10", initialOffset = "2599129"),
|
||||
// @PartitionOffset(partition = "11", initialOffset = "2546277")
|
||||
// }
|
||||
// ),
|
||||
// groupId = "${spring.kafka.group}"
|
||||
// )
|
||||
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
|
||||
public void listen(List<String> recordMessages) {
|
||||
long startTime = System.currentTimeMillis();
|
||||
try {
|
||||
@@ -136,5 +139,9 @@ public class CorpusProcessKafkaProducer {
|
||||
log.error("CorpusProcessKafkaProducer 电话语料 解析JSON出错: {}" , e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,69 @@
|
||||
package com.volvo.ai.analytic.center.mq;
|
||||
|
||||
import org.apache.kafka.clients.consumer.*;
|
||||
import org.apache.kafka.common.PartitionInfo;
|
||||
import org.apache.kafka.common.TopicPartition;
|
||||
import org.apache.kafka.common.serialization.StringDeserializer;
|
||||
import org.apache.rocketmq.spring.core.RocketMQTemplate;
|
||||
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.PostMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
|
||||
import javax.annotation.Resource;
|
||||
import java.time.Duration;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Properties;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
@Component
|
||||
@RestController
|
||||
@RefreshScope
|
||||
public class TestKafkaListener {
|
||||
|
||||
@Value("${spring.kafka.topic}")
|
||||
private String topics;
|
||||
@Value("${spring.kafka.group}")
|
||||
private String groupId;
|
||||
@Value("${spring.kafka.bootstrap-servers}")
|
||||
private String bootstrapServer;
|
||||
|
||||
@Resource
|
||||
private RocketMQTemplate rocketMqTemplate;
|
||||
|
||||
@PostMapping("corpusProcessKafkaConsumer")
|
||||
public void testKakfa (){
|
||||
|
||||
Properties props = new Properties();
|
||||
props.put("bootstrap.servers", bootstrapServer);
|
||||
props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
|
||||
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
|
||||
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
|
||||
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
|
||||
|
||||
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
|
||||
|
||||
List<PartitionInfo> partitions = consumer.partitionsFor(topics);
|
||||
List<TopicPartition> topicPartitionList = partitions
|
||||
.stream()
|
||||
.map(info -> new TopicPartition(topics, info.partition()))
|
||||
.collect(Collectors.toList());
|
||||
consumer.assign(topicPartitionList);
|
||||
|
||||
Map<TopicPartition, Long> partitionTimestampMap = topicPartitionList.stream()
|
||||
.collect(Collectors.toMap(tp -> tp, tp -> 1742832000000L));
|
||||
Map<TopicPartition, OffsetAndTimestamp> partitionOffsetMap = consumer.offsetsForTimes(partitionTimestampMap);
|
||||
partitionOffsetMap.forEach((tp, offsetAndTimestamp) -> consumer.seek(tp, offsetAndTimestamp.offset()));
|
||||
|
||||
boolean keepOnReading = true;
|
||||
while(keepOnReading){
|
||||
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
|
||||
for (ConsumerRecord<String, String> record : records){
|
||||
System.out.println(" testKakfa Message received " + record.value() + ", partition " + record.partition() + ", offset=" + record.offset() + ", timestamp=" + record.timestamp());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -6,4 +6,6 @@ import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
|
||||
public interface AiAnalysisRequestLogsService extends IService<AiAnalysisRequestLogs> {
|
||||
|
||||
boolean saveAiAnalysisRequestLogs(AiAnalysisRequestLogs aiAnalysisRequestLogs);
|
||||
|
||||
AiAnalysisRequestLogs queryAiAnalysisRequestLogsByBusinessReponse(String sourceId);
|
||||
}
|
||||
|
||||
@@ -8,7 +8,7 @@ public interface DiFyService {
|
||||
|
||||
public Object getDiFyObject(DiFyReq diFyReq);
|
||||
|
||||
public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType, String businessData);
|
||||
public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType, String businessData, String aiAnalysisRequestId);
|
||||
public JSONObject executeDifyFlow(DiFyReq diFyReq);
|
||||
|
||||
}
|
||||
|
||||
@@ -31,4 +31,9 @@ public class AiAnalysisRequestLogsServiceImpl extends ServiceImpl<AiAnalysisRequ
|
||||
return aiAnalysisRequestLogsMapper.update(aiAnalysisRequestLogs, queryWrapper) > 0;
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public AiAnalysisRequestLogs queryAiAnalysisRequestLogsByBusinessReponse(String sourceId) {
|
||||
return aiAnalysisRequestLogsMapper.queryAiAnalysisRequestLogsByBusinessReponse(sourceId);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -11,6 +11,7 @@ import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
|
||||
import com.volvo.ai.analytic.center.service.DiFyService;
|
||||
import com.volvo.ai.analytic.center.utils.AiAnalysisUtils;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
@@ -48,8 +49,8 @@ public class DiFyServiceImpl implements DiFyService{
|
||||
}
|
||||
|
||||
@Override
|
||||
public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType, String businessData) {
|
||||
String aiAnalysisRequestId = AiAnalysisUtils.getAiAnalysisRequestId(businessType);
|
||||
public JSONObject executeDifyFlow(DiFyReq diFyReq, String businessType, String businessData, String oldAiAnalysisRequestId) {
|
||||
String aiAnalysisRequestId = StringUtils.isEmpty(oldAiAnalysisRequestId)? AiAnalysisUtils.getAiAnalysisRequestId(businessType):oldAiAnalysisRequestId;
|
||||
try {
|
||||
Map<String, Object> map = new HashMap<>();
|
||||
map.put("inputs",diFyReq.getInputs());
|
||||
|
||||
@@ -203,7 +203,7 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl<TmOdsVdqwM
|
||||
corpusReportDTO.setUserId(userId);
|
||||
corpusReportDTO.setAnalysisScene(1l);
|
||||
// 获取配置
|
||||
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode(), JSONObject.toJSONString(corpusReportDTO));
|
||||
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode(), JSONObject.toJSONString(corpusReportDTO),null);
|
||||
log.info("runDify execDifyFlow {}", execDifyFlow);
|
||||
if (null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")) {
|
||||
String text = execDifyFlow.getJSONObject("outputs").getString("text");
|
||||
|
||||
@@ -141,7 +141,7 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
|
||||
corpusReportDTO.setRecordId(aicorpusTelephone.getSourceId());
|
||||
corpusReportDTO.setAnalysisScene(2l);
|
||||
// 获取配置
|
||||
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode(), JSONObject.toJSONString(corpusReportDTO));
|
||||
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode(), JSONObject.toJSONString(corpusReportDTO),null);
|
||||
log.info("runDify execDifyFlow {}",execDifyFlow);
|
||||
if(null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")){
|
||||
|
||||
|
||||
@@ -0,0 +1,19 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
|
||||
<mapper namespace="com.volvo.ai.analytic.center.mapper.AiAnalysisRequestLogsMapper">
|
||||
|
||||
<select id="queryAiAnalysisRequestLogsByBusinessReponse" resultType="com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs">
|
||||
|
||||
SELECT
|
||||
ai_analysis_request_id as aiAnalysisRequestId,
|
||||
ai_analysis_request_type as aiAnalysisRequestType,
|
||||
business_response as businessResponse
|
||||
FROM
|
||||
tt_ai_analysis_request_logs
|
||||
WHERE
|
||||
ai_analysis_request_type = 'SMART_ASSISTANT'
|
||||
AND business_response !='' AND business_response is not null
|
||||
AND JSON_EXTRACT( CAST( business_response AS JSON ), '$.recordId' ) = #{sourceId} limit 1
|
||||
</select>
|
||||
|
||||
</mapper>
|
||||
Reference in New Issue
Block a user