【发布时间】:2020-01-29 05:36:33
【问题描述】:
我有工厂 bean:
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(
CustomConsumerRebalanceListener consumerRebalanceListener, ConsumerFactory consumerFactory, CustomConfiguration customConfiguration) {
ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory);
factory.getContainerProperties().setAckOnError(false);
factory.setBatchListener(true);
factory.setStatefulRetry(true);
factory.setBatchErrorHandler(errorHandler());
factory.setConcurrency(customConfiguration.getConcurrency());
factory.getContainerProperties().setConsumerRebalanceListener(consumerRebalanceListener);
return factory;
}
这是我的听众:
@KafkaListener(topics = "${spring.kafka.topic}")
private synchronized void consumeKafkaQueue(@Payload String message, Acknowledgment acknowledgment) {
...
}
如何将侦听器配置为仅接收新消息而不在启动时获取旧消息?
【问题讨论】:
-
我从来没有使用过这个库,但我猜
acknowledgment.acknowledge(); -
Acknowledgment实例仅适用于当前消息;看我的回答。
标签: java spring apache-kafka spring-kafka