去除mock
This commit is contained in:
@@ -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<String, String> 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<Object> 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 消息已发送");
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<String> 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());
|
||||
|
||||
@@ -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<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){
|
||||
log.info(" testKakfa Message received " + record.value() + ", partition " + record.partition() + ", offset=" + record.offset() + ", timestamp=" + record.timestamp());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user