From 4eb63000ae858c3c4b7cc5ffd033520ce768579c Mon Sep 17 00:00:00 2001 From: zren25 Date: Tue, 1 Apr 2025 01:16:23 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E6=94=B9=E5=88=86=E5=8C=BA=E8=AE=BE?= =?UTF-8?q?=E7=BD=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../center/mq/CorpusProcessKafkaProducer.java | 45 ++++++++++--------- 1 file changed, 23 insertions(+), 22 deletions(-) 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 4afa417..9da8485 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 @@ -18,6 +18,8 @@ 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.messaging.support.MessageBuilder; import org.springframework.stereotype.Component; import org.springframework.web.bind.annotation.RestController; @@ -57,28 +59,27 @@ public class CorpusProcessKafkaProducer { @Resource private RocketMQTemplate rocketMqTemplate; -// @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( + topicPartitions = @TopicPartition( + topic = "${spring.kafka.topic}", + 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 {