修改配置验证数据
This commit is contained in:
@@ -87,7 +87,7 @@ public class KafkaConfig {
|
|||||||
new ConcurrentKafkaListenerContainerFactory<>();
|
new ConcurrentKafkaListenerContainerFactory<>();
|
||||||
factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(props));
|
factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(props));
|
||||||
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
|
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
|
||||||
factory.setBatchListener(true);
|
// factory.setBatchListener(true);
|
||||||
return factory;
|
return factory;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -46,7 +46,8 @@ public class NameplateKafkaConsumer {
|
|||||||
|
|
||||||
@KafkaListener(topics = "${analyticCenterKafka.consumer.topic}",
|
@KafkaListener(topics = "${analyticCenterKafka.consumer.topic}",
|
||||||
groupId = "${analyticCenterKafka.consumer.group}" ,
|
groupId = "${analyticCenterKafka.consumer.group}" ,
|
||||||
containerFactory = "analyticCenterConsumerFactory")
|
containerFactory = "analyticCenterConsumerFactory",
|
||||||
|
concurrency = "3")
|
||||||
public void listen(String recordMessages, Acknowledgment ack) {
|
public void listen(String recordMessages, Acknowledgment ack) {
|
||||||
long startTime = System.currentTimeMillis();
|
long startTime = System.currentTimeMillis();
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user