From a21aa1025182775ad431e989734873fac0ffe3dc Mon Sep 17 00:00:00 2001 From: zren25 Date: Mon, 31 Mar 2025 21:09:55 +0800 Subject: [PATCH] =?UTF-8?q?=E6=8F=90=E4=BA=A4=E9=AA=8C=E8=AF=81kafka?= =?UTF-8?q?=E6=96=B9=E6=B3=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../dto/corpus/AicorpusTelephoneDTO.java | 4 + .../ai/analytic/center/job/CorpusFailJob.java | 104 ++++++++++++------ .../mapper/AiAnalysisRequestLogsMapper.java | 3 + .../center/mq/CorpusDccMqConsumer.java | 26 ++++- .../center/mq/CorpusProcessKafkaProducer.java | 57 +++++----- .../analytic/center/mq/TestKafkaListener.java | 69 ++++++++++++ .../service/AiAnalysisRequestLogsService.java | 2 + .../analytic/center/service/DiFyService.java | 2 +- .../AiAnalysisRequestLogsServiceImpl.java | 5 + .../center/service/impl/DiFyServiceImpl.java | 5 +- .../TmOdsVdqwMessagearchivingServiceImpl.java | 2 +- .../impl/TmTelephoneCorpusServiceImpl.java | 2 +- .../mapper/AiAnalysisRequestLogsMapper.xml | 19 ++++ 13 files changed, 230 insertions(+), 70 deletions(-) create mode 100644 ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/TestKafkaListener.java create mode 100644 ai-analytic-center-biz/src/main/resources/mapper/AiAnalysisRequestLogsMapper.xml diff --git a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/corpus/AicorpusTelephoneDTO.java b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/corpus/AicorpusTelephoneDTO.java index 73a5560..caf3a40 100644 --- a/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/corpus/AicorpusTelephoneDTO.java +++ b/ai-analytic-center-api/src/main/java/com/volvo/ai/analytic/center/dto/corpus/AicorpusTelephoneDTO.java @@ -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; diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java index 240c631..5ebd77c 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java @@ -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 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 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 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); diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisRequestLogsMapper.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisRequestLogsMapper.java index 2fc6aeb..9eafa77 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisRequestLogsMapper.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mapper/AiAnalysisRequestLogsMapper.java @@ -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 { + + public AiAnalysisRequestLogs queryAiAnalysisRequestLogsByBusinessReponse(@Param("sourceId") String sourceId); } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusDccMqConsumer.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusDccMqConsumer.java index a3d4373..7cd60cf 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusDccMqConsumer.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusDccMqConsumer.java @@ -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 { @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 { 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()); diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java index 6e7d896..f721549 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/CorpusProcessKafkaProducer.java @@ -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 recordMessages) { long startTime = System.currentTimeMillis(); try { @@ -136,5 +139,9 @@ public class CorpusProcessKafkaProducer { log.error("CorpusProcessKafkaProducer 电话语料 解析JSON出错: {}" , e.getMessage()); } } + + + + } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/TestKafkaListener.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/TestKafkaListener.java new file mode 100644 index 0000000..7e98c60 --- /dev/null +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/TestKafkaListener.java @@ -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 consumer = new KafkaConsumer<>(props); + + List partitions = consumer.partitionsFor(topics); + List topicPartitionList = partitions + .stream() + .map(info -> new TopicPartition(topics, info.partition())) + .collect(Collectors.toList()); + consumer.assign(topicPartitionList); + + Map partitionTimestampMap = topicPartitionList.stream() + .collect(Collectors.toMap(tp -> tp, tp -> 1742832000000L)); + Map partitionOffsetMap = consumer.offsetsForTimes(partitionTimestampMap); + partitionOffsetMap.forEach((tp, offsetAndTimestamp) -> consumer.seek(tp, offsetAndTimestamp.offset())); + + boolean keepOnReading = true; + while(keepOnReading){ + ConsumerRecords records = consumer.poll(Duration.ofMillis(100)); + for (ConsumerRecord record : records){ + System.out.println(" testKakfa Message received " + record.value() + ", partition " + record.partition() + ", offset=" + record.offset() + ", timestamp=" + record.timestamp()); + } + } + } +} + \ No newline at end of file diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiAnalysisRequestLogsService.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiAnalysisRequestLogsService.java index 4d7b7a5..d7bdef1 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiAnalysisRequestLogsService.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/AiAnalysisRequestLogsService.java @@ -6,4 +6,6 @@ import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs; public interface AiAnalysisRequestLogsService extends IService { boolean saveAiAnalysisRequestLogs(AiAnalysisRequestLogs aiAnalysisRequestLogs); + + AiAnalysisRequestLogs queryAiAnalysisRequestLogsByBusinessReponse(String sourceId); } 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 0600d46..7159e60 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 @@ -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); } diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisRequestLogsServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisRequestLogsServiceImpl.java index a4b4776..0c2c050 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisRequestLogsServiceImpl.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/AiAnalysisRequestLogsServiceImpl.java @@ -31,4 +31,9 @@ public class AiAnalysisRequestLogsServiceImpl extends ServiceImpl 0; } } + + @Override + public AiAnalysisRequestLogs queryAiAnalysisRequestLogsByBusinessReponse(String sourceId) { + return aiAnalysisRequestLogsMapper.queryAiAnalysisRequestLogsByBusinessReponse(sourceId); + } } 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 0833aa5..d5b49d2 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 @@ -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 map = new HashMap<>(); map.put("inputs",diFyReq.getInputs()); diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java index 397acf2..c375b6d 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmOdsVdqwMessagearchivingServiceImpl.java @@ -203,7 +203,7 @@ public class TmOdsVdqwMessagearchivingServiceImpl extends ServiceImpl + + + + + + \ No newline at end of file