【问题标题】:Not able to set error handler for Batch mode In kafka无法在kafka中为批处理模式设置错误处理程序
【发布时间】:2021-06-15 01:07:22
【问题描述】:

这是我的配置代码。

    import java.util.function.Consumer;

import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.springframework.boot.autoconfigure.kafka.ConcurrentKafkaListenerContainerFactoryConfigurer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.listener.MessageListenerContainer;
import org.springframework.kafka.listener.SeekToCurrentBatchErrorHandler;
import org.springframework.web.client.RestTemplate;

@Configuration
public class StreamConfiguration {
    @Bean
    public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory(
            ConcurrentKafkaListenerContainerFactoryConfigurer factoryConfigure,
            ConsumerFactory<Object, Object> kafkaConsumerFactory) {
        ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factoryConfigure.configure(factory, kafkaConsumerFactory);
        factory.setBatchListener(true);

        factory.setBatchErrorHandler(new SeekToCurrentBatchErrorHandler() {

            @Override
            public void handle(Exception thrownException, ConsumerRecords<?, ?> data, Consumer<?, ?> consumer,
                    MessageListenerContainer container) {
                Config.this.ehException = thrownException;
                super.handle(thrownException, data, consumer, container);
            }
        });

        return factory;
    }

    

    @Bean
    public RestTemplate restTemplate() {
        return new RestTemplate();
    }

}

这是我的消费者代码

 @KafkaListener(id = "#{'${spring.kafka.listener.id}'}", topics = "#{'${spring.kafka.consumer.topic}'}")
    public void getTopics(@RequestBody List<Request> model) {

        streamProcessor.runParallel(model.parallelStream());


    }

我在处理异常时遇到错误,指出 Consumer 类型的参数数量不正确;它不能用 arguments 参数化。 而且我对导入要导入的配置感到困惑(Config.this.ehException = throwException;),因为有两个选项 apache common 和 apache client admin。

请帮助我无法设置批处理错误处理程序,并且由于反序列化错误而处于无限循环中:(((。 我正在使用 Java8

【问题讨论】:

    标签: apache-kafka spring-kafka


    【解决方案1】:

    你不能在那里使用@RequestBody - 这是一个仅限网络的注释。

    删除它。

    如果您仍然遇到问题,请将完整的堆栈跟踪添加到您的问题中。

    【讨论】:

    • 谢谢您:)。我能够使用您的答案之一移动偏移量的失败记录。我仍然难以通过 application.properties 设置 errorHandlerDeserliazer。您能指点我吗为此实现了示例。这将是非常糟糕的。
    • 我对这个问题的回答stackoverflow.com/questions/63236346/… 有一个例子(接近底部)。
    • 非常感谢 Gary 这给了我一张清晰的照片。
    猜你喜欢
    • 2021-08-31
    • 2023-04-10
    • 1970-01-01
    • 2016-09-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2010-12-17
    • 2021-05-24
    相关资源
    最近更新 更多