From b14187005e22a2e12d19b429ad5ed7e2aa19cb45 Mon Sep 17 00:00:00 2001 From: zren25 Date: Tue, 1 Apr 2025 00:56:31 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E6=94=B9kafka=E5=88=86=E5=8C=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../center/mq/CorpusDccMqConsumer.java | 2 +- .../center/mq/CorpusProcessKafkaProducer.java | 26 +++++++++++++++++-- 2 files changed, 25 insertions(+), 3 deletions(-) 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 6f1f183..c5c4dd2 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 @@ -33,7 +33,7 @@ import org.springframework.web.bind.annotation.RestController; @RocketMQMessageListener(consumerGroup = "${rocketmq.consumer.corpus.dcctopicgroup}", topic = "${rocketmq.consumer.corpus.dcctopic}", instanceName = "CorpusDccMqConsumer1", - consumeThreadNumber = 10, + consumeThreadNumber = 40, enableMsgTrace = true) public class CorpusDccMqConsumer implements RocketMQListener { 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 7839959..5f35be0 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.PostMapping; @@ -58,8 +60,28 @@ public class CorpusProcessKafkaProducer { @Resource private RocketMQTemplate rocketMqTemplate; - @PostMapping("corpusProcessKafkaProducer") - @KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") + @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 {