【问题标题】:Spring Kafka - How to Retry with @KafkaListenerSpring Kafka - 如何使用 @KafkaListener 重试
【发布时间】:2019-01-20 16:44:41
【问题描述】:

来自推特的问题:

只是想找出一个简单的例子,用 spring-kafka 2.1.7 与 KafkaListener 和 AckMode.MANUAL_IMMEDIATE 一起工作,重试最后一条失败的消息。

https://twitter.com/tolbier/status/1028936942447149056

【问题讨论】:

    标签: spring spring-boot spring-kafka


    【解决方案1】:

    通常最好在 Stack Overflow 上提出此类问题(标记为

    有两种方式:

    • RetryTemplate 添加到侦听器容器工厂 - 重试将在内存中执行,您可以设置退避属性。
    • 添加SeekToCurrentErrorHandler,它将重新查找未处理的记录。

    这是一个例子:

    @SpringBootApplication
    public class Twitter1Application {
    
        public static void main(String[] args) {
            SpringApplication.run(Twitter1Application.class, args);
        }
    
        boolean fail = true;
    
        @KafkaListener(id = "foo", topics = "twitter1")
        public void listen(String in, Acknowledgment ack) {
            System.out.println(in);
            if (fail) {
                fail = false;
                throw new RuntimeException("failed");
            }
            ack.acknowledge();
        }
    
        @Bean
        public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory(
                ConcurrentKafkaListenerContainerFactoryConfigurer configurer,
                ConsumerFactory<Object, Object> kafkaConsumerFactory) {
            ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
            configurer.configure(factory, kafkaConsumerFactory);
            factory.getContainerProperties().setErrorHandler(new SeekToCurrentErrorHandler());
            // or factory.setRetryTemplate(aRetryTemplate);
            // and factory.setRecoveryCallback(aRecoveryCallback);
            return factory;
        }
    
        @Bean
        public ApplicationRunner runner(KafkaTemplate<String, String> template) {
            return args -> {
                Thread.sleep(2000);
                template.send("twitter1", "foo");
                template.send("twitter1", "bar");
            };
        }
    
        @Bean
        public NewTopic topic() {
            return new NewTopic("twitter1", 1, (short) 1);
        }
    
    }
    

    spring.kafka.consumer.auto-offset-reset=earliest
    spring.kafka.consumer.enable-auto-commit=false
    
    spring.kafka.listener.ack-mode=manual-immediate
    
    logging.level.org.springframework.kafka=debug
    

    <?xml version="1.0" encoding="UTF-8"?>
    <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
        xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
        <modelVersion>4.0.0</modelVersion>
    
        <groupId>com.example</groupId>
        <artifactId>twitter1</artifactId>
        <version>0.0.1-SNAPSHOT</version>
        <packaging>jar</packaging>
    
        <name>twitter1</name>
        <description>Demo project for Spring Boot</description>
    
        <parent>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-parent</artifactId>
            <version>2.0.4.RELEASE</version>
            <relativePath/> <!-- lookup parent from repository -->
        </parent>
    
        <properties>
            <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
            <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
            <java.version>1.8</java.version>
        </properties>
    
        <dependencies>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter</artifactId>
            </dependency>
            <dependency>
                <groupId>org.springframework.kafka</groupId>
                <artifactId>spring-kafka</artifactId>
            </dependency>
    
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-starter-test</artifactId>
                <scope>test</scope>
            </dependency>
        </dependencies>
    
        <build>
            <plugins>
                <plugin>
                    <groupId>org.springframework.boot</groupId>
                    <artifactId>spring-boot-maven-plugin</artifactId>
                </plugin>
            </plugins>
        </build>
    
    
    </project>
    

    (启动 2.0.4 引入 2.1.8,即当前版本)。

    foo
    2018-08-13 17:36:14.901 ERROR 3945 --- [      foo-0-C-1] essageListenerContainer$ListenerConsumer : Error handler threw an exception
    
    org.springframework.kafka.KafkaException: Seek to current after exception; nested exception is  ...    
    
    2018-08-13 17:36:15.396 DEBUG 3945 --- [      foo-0-C-1] essageListenerContainer$ListenerConsumer : Received: 2 records
    foo
    2018-08-13 17:36:15.398 DEBUG 3945 --- [      foo-0-C-1] essageListenerContainer$ListenerConsumer : Committing: {twitter1-0=OffsetAndMetadata{offset=5, metadata=''}}
    bar
    2018-08-13 17:36:15.403 DEBUG 3945 --- [      foo-0-C-1] essageListenerContainer$ListenerConsumer : Committing: {twitter1-0=OffsetAndMetadata{offset=6, metadata=''}}
    

    在即将发布的 2.2 版本中,可以为错误处理程序配置恢复器,并提供标准恢复器以将失败的记录发布到死信主题。

    Commit hereDocs Here.

    【讨论】:

    • 嗨 - 很抱歉恢复旧线程,但我无法弄清楚如何让 SeekToCurrentErrorHandler 使用该恢复器(Kafka 2.3)。由于我们使用 auto commit=false,recoverer 如何提交失败的记录,使其不再被读取?它是一个简单的 BiConsumer,无法访问任何上下文或确认对象。
    • SeekToCurrentErrorHandler eh = new SeekToCurrentErrorHandler(recoverer, backoff); eh.setCommitRecovered(true);不起作用。即使在调用恢复器之后,偏移量也不会提交。 (对不起 cmets 中的代码)
    • 您应该提出一个新问题并显示更多上下文;不鼓励用新问题评论旧答案。使用 MANUAL_IMMEDIATE ack 模式时使用 commitRecovered。从 SK 版本 2.3.2 开始,STCEH 为 isAckAfterHandle() 返回 true。
    • 我完全理解,但这个问题与我想问的问题完全一样,再次问同样的事情似乎是多余的(并将其标记为重复)。我真的很感谢你花时间回复。我在 containerfactory 上使用 errorHandler.setCommitRecovered(true) 和手动立即确认模式。但是无论如何都不会提交恢复的记录。调用了恢复器 lambda,但它无法确认恢复,错误处理程序必须假定这是所需的操作。但它似乎并没有这样做。
    • 哎哟。我想我遇到了这个:github.com/spring-projects/spring-kafka/issues/1309
    猜你喜欢
    • 2020-02-08
    • 2019-09-16
    • 2018-06-27
    • 2020-06-23
    • 2019-08-24
    • 2020-10-03
    • 1970-01-01
    • 2020-05-27
    • 1970-01-01
    相关资源
    最近更新 更多