From f1fe84a38297bc9f0e2f5c38a6201c55ef58b206 Mon Sep 17 00:00:00 2001 From: zren25 Date: Mon, 31 Mar 2025 21:56:52 +0800 Subject: [PATCH] =?UTF-8?q?=E5=8E=BB=E9=99=A4mock?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../center/controller/TestController.java | 31 -------- .../center/mq/CorpusProcessKafkaProducer.java | 41 ++++------- .../analytic/center/mq/TestKafkaListener.java | 71 ------------------- 3 files changed, 14 insertions(+), 129 deletions(-) delete mode 100644 ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/TestKafkaListener.java diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/TestController.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/TestController.java index 48071ea..5f316f4 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/TestController.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/TestController.java @@ -1,19 +1,13 @@ package com.volvo.ai.analytic.center.controller; -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.common.core.util.ResultMsg; import io.swagger.annotations.Api; import io.swagger.annotations.ApiOperation; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.spring.core.RocketMQTemplate; 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.core.KafkaTemplate; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; @@ -30,15 +24,6 @@ public class TestController { @Autowired private RocketMQTemplate rocketMQTemplate; - @Autowired - private KafkaTemplate kafkaTemplate; // 注入 KafkaTemplate - - private final ObjectMapper objectMapper = new ObjectMapper(); - - @Value("spring.kafka.producer.topic") - private String kafkaTopic; // 注入 KafkaTemplate - @Value("") - private String kafkaTopicGroup; @PostMapping("/mockMq") @ApiOperation(value = "补偿处理消息") @@ -47,21 +32,5 @@ public class TestController { return ResultMsg.ok("ok"); } - @PostMapping("/mockKafka") - @ApiOperation(value = "生成 Kafka 数据") - public ResultMsg mockKafka(@RequestBody String message) { - try { - AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class); - log.info("mockKafka:{}",message); - for (int i = 0; i < 20; i++){ - aicorpusTelephone.setSourceId(aicorpusTelephone.getSourceId().concat("_"+i)); - kafkaTemplate.send("topic_voc_covert_text_log", JSONObject.toJSONString(aicorpusTelephone)); // 发送 Kafka 消息 - } - log.info("Kafka 消息已发送: {}", message); - } catch (JsonProcessingException e) { - throw new RuntimeException(e); - } - return ResultMsg.ok("Kafka 消息已发送"); - } } 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 f721549..35519cc 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,7 +9,6 @@ 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; @@ -18,7 +17,6 @@ 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.TopicPartition; import org.springframework.kafka.listener.ConsumerSeekAware; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Component; @@ -26,9 +24,8 @@ import org.springframework.web.bind.annotation.RestController; import javax.annotation.Resource; import java.time.LocalDateTime; -import java.util.Collections; +import java.time.format.DateTimeFormatter; import java.util.List; -import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -54,32 +51,13 @@ public class CorpusProcessKafkaProducer implements ConsumerSeekAware { @Value("${rocketmq.producer.corpus.dcctopic}") private String dccMqTipic; + @Value("${dify.corpus.kafkaTimeLimit}") + private String kafkaTimeLimit; + @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}" -// ) -@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") + @KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") public void listen(List recordMessages) { long startTime = System.currentTimeMillis(); try { @@ -95,6 +73,15 @@ public class CorpusProcessKafkaProducer implements ConsumerSeekAware { executor.submit(() -> { try { AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class); + + String transcribeTimeStr = aicorpusTelephone.getTranscribeTime(); + if (transcribeTimeStr != null) { + LocalDateTime transcribeTime = LocalDateTime.parse(transcribeTimeStr, DateTimeFormatter.ISO_LOCAL_DATE_TIME); + LocalDateTime kafkaStartTime = LocalDateTime.parse(kafkaTimeLimit, DateTimeFormatter.ISO_LOCAL_DATE_TIME); + if (transcribeTime.isBefore(kafkaStartTime)) { + return; + } + } log.info("aicorpusTelephone categoryCode:{}, display: {}", aicorpusTelephone.getCategoryCode(), aicorpusTelephone.getDisplay()); DisplayDTO display = objectMapper.readValue(aicorpusTelephone.getDisplay(), DisplayDTO.class); log.info("aicorpusTelephone display getSegments: {}", display.getSegments()); 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 deleted file mode 100644 index 6d5918a..0000000 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/TestKafkaListener.java +++ /dev/null @@ -1,71 +0,0 @@ -package com.volvo.ai.analytic.center.mq; - -import lombok.extern.slf4j.Slf4j; -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; - -@Slf4j -@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){ - log.info(" testKakfa Message received " + record.value() + ", partition " + record.partition() + ", offset=" + record.offset() + ", timestamp=" + record.timestamp()); - } - } - } -} - \ No newline at end of file