修改kafka分区
This commit is contained in:
@@ -33,7 +33,7 @@ import org.springframework.web.bind.annotation.RestController;
|
|||||||
@RocketMQMessageListener(consumerGroup = "${rocketmq.consumer.corpus.dcctopicgroup}",
|
@RocketMQMessageListener(consumerGroup = "${rocketmq.consumer.corpus.dcctopicgroup}",
|
||||||
topic = "${rocketmq.consumer.corpus.dcctopic}",
|
topic = "${rocketmq.consumer.corpus.dcctopic}",
|
||||||
instanceName = "CorpusDccMqConsumer1",
|
instanceName = "CorpusDccMqConsumer1",
|
||||||
consumeThreadNumber = 10,
|
consumeThreadNumber = 40,
|
||||||
enableMsgTrace = true)
|
enableMsgTrace = true)
|
||||||
public class CorpusDccMqConsumer implements RocketMQListener<MessageExt> {
|
public class CorpusDccMqConsumer implements RocketMQListener<MessageExt> {
|
||||||
|
|
||||||
|
|||||||
@@ -18,6 +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.PostMapping;
|
||||||
@@ -58,8 +60,28 @@ public class CorpusProcessKafkaProducer {
|
|||||||
@Resource
|
@Resource
|
||||||
private RocketMQTemplate rocketMqTemplate;
|
private RocketMQTemplate rocketMqTemplate;
|
||||||
|
|
||||||
@PostMapping("corpusProcessKafkaProducer")
|
@KafkaListener(
|
||||||
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
|
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<ConsumerRecord<String, Object>> recordMessages) {
|
public void listen(List<ConsumerRecord<String, Object>> recordMessages) {
|
||||||
long startTime = System.currentTimeMillis();
|
long startTime = System.currentTimeMillis();
|
||||||
try {
|
try {
|
||||||
|
|||||||
Reference in New Issue
Block a user