去掉分区设置

This commit is contained in:
zren25
2025-04-01 01:04:42 +08:00
parent b14187005e
commit 710539eb8a

View File

@@ -18,11 +18,8 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope; import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.kafka.annotation.KafkaListener; 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.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RestController; import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource; import javax.annotation.Resource;
@@ -60,28 +57,28 @@ public class CorpusProcessKafkaProducer {
@Resource @Resource
private RocketMQTemplate rocketMqTemplate; private RocketMQTemplate rocketMqTemplate;
@KafkaListener( // @KafkaListener(
topicPartitions = @TopicPartition( // topicPartitions = @TopicPartition(
topic = "${spring.kafka.topic}", // topic = "${spring.kafka.topic}",
partitions = {"0", "1","2", "3","4", "5","6", "7","8", "9", "10", "11"}, // partitions = {"0", "1","2", "3","4", "5","6", "7","8", "9", "10", "11"},
partitionOffsets = { // partitionOffsets = {
@PartitionOffset(partition = "0", initialOffset = "2792520"), // @PartitionOffset(partition = "0", initialOffset = "2792520"),
@PartitionOffset(partition = "1", initialOffset = "2596153"), // @PartitionOffset(partition = "1", initialOffset = "2596153"),
@PartitionOffset(partition = "2", initialOffset = "2536889"), // @PartitionOffset(partition = "2", initialOffset = "2536889"),
@PartitionOffset(partition = "3", initialOffset = "2782173"), // @PartitionOffset(partition = "3", initialOffset = "2782173"),
@PartitionOffset(partition = "4", initialOffset = "2616677"), // @PartitionOffset(partition = "4", initialOffset = "2616677"),
@PartitionOffset(partition = "5", initialOffset = "2529585"), // @PartitionOffset(partition = "5", initialOffset = "2529585"),
@PartitionOffset(partition = "6", initialOffset = "2761950"), // @PartitionOffset(partition = "6", initialOffset = "2761950"),
@PartitionOffset(partition = "7", initialOffset = "2581957"), // @PartitionOffset(partition = "7", initialOffset = "2581957"),
@PartitionOffset(partition = "8", initialOffset = "2530225"), // @PartitionOffset(partition = "8", initialOffset = "2530225"),
@PartitionOffset(partition = "9", initialOffset = "2803179"), // @PartitionOffset(partition = "9", initialOffset = "2803179"),
@PartitionOffset(partition = "10", initialOffset = "2599129"), // @PartitionOffset(partition = "10", initialOffset = "2599129"),
@PartitionOffset(partition = "11", initialOffset = "2546277") // @PartitionOffset(partition = "11", initialOffset = "2546277")
} // }
), // ),
groupId = "${spring.kafka.group}" // 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<ConsumerRecord<String, Object>> recordMessages) { public void listen(List<ConsumerRecord<String, Object>> recordMessages) {
long startTime = System.currentTimeMillis(); long startTime = System.currentTimeMillis();
try { try {