【问题标题】:@KafkaListner - Skip older messages aka Receive only new messages@KafkaListner - 跳过旧消息,也就是只接收新消息
【发布时间】: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


【解决方案1】:

让您的侦听器实现ConsumerSeekAware 并在分配分区时搜索到最后。请参阅this answer,特别是 EDIT2 之后的示例;请改用seekToEnd

或者,您可以使用auto.offset.reset=latest,并在每次应用启动时使用唯一的group.id

还有

factory.setConcurrency(customConfiguration.getConcurrency());

使用synchronized 监听器方法会破坏并发;其他线程将等待同步。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-04-03
    • 2021-06-11
    • 1970-01-01
    • 1970-01-01
    • 2023-03-16
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多