【问题标题】:How shutdown KafkaListener when error occurs发生错误时如何关闭KafkaListener
【发布时间】:2020-10-14 14:16:44
【问题描述】:

我就是这样写了一个Listener

@Autowired
private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry;

@KafkaListener(containerFactory = "cdcKafkaListenerContainerFactory", errorHandler = "errorHandler")
public void consume(@Payload String message) throws Exception {
    ...
}

@Bean
public KafkaListenerErrorHandler errorHandler() {
    return ((message, e) -> {
        kafkaListenerEndpointRegistry.stop();
        return null;
    });
}

@KafkaListener 注释中,我指定了我的错误处理程序,它只是停止消费者。 它似乎有效,但我有一些问题要问。

这个范围有内置的errorHandler吗?我读过ContainerStoppingErrorHandler 可以使用,但我无法设置它,因为@KafkaListener 的errorHandler 接受KafkaListenerErrorHandler 类型的bean。

我看到kafkaListenerEndpointRegistry.stop();优雅地停下来。因此,在停止之前已提交已消费消息的分区偏移量。 我知道的是,当kafkaListenerEndpointRegistry.stop(); 被调用并且在监听器被确定关闭之前,另一条消息到达主题时会发生什么? 这条消息被消费了吗?

我想象这个场景

time0: kafkaListenerEndpointRegistry.stop() is called
time1: a message is pushed into the listened topic
time2: kafkaListenerEndpointRegistry.stop() complete graceful stop

我担心消息可能会在 time1 到达。在这种情况下会发生什么?

【问题讨论】:

    标签: apache-kafka spring-kafka


    【解决方案1】:

    不要在监听器中停止容器。

    ContainerStoppingErrorHandler 设置在容器工厂上,而不是注解上。

    如果您使用的是 Spring Boot,只需将错误处理程序声明为 bean,boot 就会将其连接进来。

    否则将错误处理程序添加到连接工厂 bean。

    使用此错误处理程序,抛出异常将立即停止容器。

    【讨论】:

    • 谢谢加里。我将factory.setErrorHandler(new ContainerStoppingErrorHandler()); 添加到 kafka 侦听器容器工厂中。这样我就应该安全了?
    猜你喜欢
    • 2021-11-22
    • 1970-01-01
    • 1970-01-01
    • 2020-04-08
    • 1970-01-01
    • 2010-12-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多