【问题标题】:Getting java.lang.ClassCastException when deploying my batch consumer using spring-kafka使用 spring-kafka 部署我的批处理消费者时获取 java.lang.ClassCastException
【发布时间】:2020-11-12 05:23:31
【问题描述】:

我正在使用 spring-kafka 2.2.8 创建批量消费者,并在部署消费者时遇到异常。

java.lang.ClassCastException: org.springframework.kafka.listener.adapter.BatchMessagingMessageListenerAdapter cannot be cast to org.springframework.kafka.listener.MessageListener
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.kafka.config.AbstractKafkaListenerEndpoint.setupMessageListener(AbstractKafkaListenerEndpoint.java:455) ~[spring-kafka-2.2.7.RELEASE.jar!/:2.2.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.kafka.config.AbstractKafkaListenerEndpoint.setupListenerContainer(AbstractKafkaListenerEndpoint.java:433) ~[spring-kafka-2.2.7.RELEASE.jar!/:2.2.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.kafka.config.AbstractKafkaListenerContainerFactory.createListenerContainer(AbstractKafkaListenerContainerFactory.java:310) ~[spring-kafka-2.2.7.RELEASE.jar!/:2.2.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.kafka.config.AbstractKafkaListenerContainerFactory.createListenerContainer(AbstractKafkaListenerContainerFactory.java:62) ~[spring-kafka-2.2.7.RELEASE.jar!/:2.2.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.kafka.config.KafkaListenerEndpointRegistry.createListenerContainer(KafkaListenerEndpointRegistry.java:200) ~[spring-kafka-2.2.7.RELEASE.jar!/:2.2.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.kafka.config.KafkaListenerEndpointRegistry.registerListenerContainer(KafkaListenerEndpointRegistry.java:172) ~[spring-kafka-2.2.7.RELEASE.jar!/:2.2.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.kafka.config.KafkaListenerEndpointRegistry.registerListenerContainer(KafkaListenerEndpointRegistry.java:146) ~[spring-kafka-2.2.7.RELEASE.jar!/:2.2.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.kafka.config.KafkaListenerEndpointRegistrar.registerAllEndpoints(KafkaListenerEndpointRegistrar.java:164) ~[spring-kafka-2.2.7.RELEASE.jar!/:2.2.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.kafka.config.KafkaListenerEndpointRegistrar.afterPropertiesSet(KafkaListenerEndpointRegistrar.java:158) ~[spring-kafka-2.2.7.RELEASE.jar!/:2.2.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.kafka.annotation.KafkaListenerAnnotationBeanPostProcessor.afterSingletonsInstantiated(KafkaListenerAnnotationBeanPostProcessor.java:263) ~[spring-kafka-2.2.7.RELEASE.jar!/:2.2.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.beans.factory.support.DefaultListableBeanFactory.preInstantiateSingletons(DefaultListableBeanFactory.java:862) ~[spring-beans-5.1.9.RELEASE.jar!/:5.1.9.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.context.support.AbstractApplicationContext.finishBeanFactoryInitialization(AbstractApplicationContext.java:877) ~[spring-context-5.1.9.RELEASE.jar!/:5.1.9.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.context.support.AbstractApplicationContext.refresh(AbstractApplicationContext.java:549) ~[spring-context-5.1.9.RELEASE.jar!/:5.1.9.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.boot.web.servlet.context.ServletWebServerApplicationContext.refresh(ServletWebServerApplicationContext.java:141) ~[spring-boot-2.1.7.RELEASE.jar!/:2.1.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.boot.SpringApplication.refresh(SpringApplication.java:743) ~[spring-boot-2.1.7.RELEASE.jar!/:2.1.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.boot.SpringApplication.refreshContext(SpringApplication.java:390) ~[spring-boot-2.1.7.RELEASE.jar!/:2.1.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.boot.SpringApplication.run(SpringApplication.java:312) ~[spring-boot-2.1.7.RELEASE.jar!/:2.1.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.boot.SpringApplication.run(SpringApplication.java:1214) ~[spring-boot-2.1.7.RELEASE.jar!/:2.1.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.boot.SpringApplication.run(SpringApplication.java:1203) ~[spring-boot-2.1.7.RELEASE.jar!/:2.1.7.RELEASE]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at com.abc.xyz.yyy.kafka.integrationtestapi.Application.main(Application.java:10) ~[classes/:?]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) ~[?:1.8.0_202]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) ~[?:1.8.0_202]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) ~[?:1.8.0_202]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at java.lang.reflect.Method.invoke(Method.java:498) ~[?:1.8.0_202]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.boot.loader.MainMethodRunner.run(MainMethodRunner.java:48) ~[app/:?]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.boot.loader.Launcher.launch(Launcher.java:87) ~[app/:?]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.boot.loader.Launcher.launch(Launcher.java:51) ~[app/:?]
   2020-11-11T13:38:31.97-0500 [APP/PROC/WEB/0] OUT     at org.springframework.boot.loader.JarLauncher.main(JarLauncher.java:52) ~[app/:?]

这是我的消费者配置

@Bean
public ConsumerFactory consumerFactory(){
    return new DefaultKafkaConsumerFactory(consumerConfigs(),stringKeyDeserializer(), avroValueDeserializer());
}
@Bean
public RetryPolicy getRetryPolicy(){
    SimpleRetryPolicy simpleRetryPolicy = new SimpleRetryPolicy();
    simpleRetryPolicy.setMaxAttempts(3);
    return simpleRetryPolicy;
}

@Bean
public FixedBackOffPolicy getBackOffPolicy() {
    FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy();
    backOffPolicy.setBackOffPeriod(100);
    return backOffPolicy;
}

@Bean
public RetryTemplate getRetryTemplate(){
    RetryTemplate retryTemplate = new RetryTemplate();
    retryTemplate.setRetryPolicy(getRetryPolicy());
    retryTemplate.setBackOffPolicy(getBackOffPolicy());
    return retryTemplate;
}

@Bean
public ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory(){
    ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory();
    factory.setConsumerFactory(consumerFactory());
    factory.setBatchListener(true);
    factory.setStatefulRetry(true);
    factory.setRetryTemplate(getRetryTemplate());
    return factory;
}

现在我的问题是,

  1. 设置retryTemplate时是否必须设置ErrorHandler?还是我在这里遗漏了导致此异常的其他内容?

【问题讨论】:

    标签: spring-kafka


    【解决方案1】:

    批处理侦听器不支持RetryTemplate(我们不知道批处理中的哪条记录失败)。

    对于这种情况,最好在侦听器中直接使用RetryTemplate

    像这样……

    @KafkaListener(...)
    void listen(List<Foo> in) {
        in.forEach(foo -> {
            this.retryTemplate.execute(context -> {
                // process the foo
            }, context -> {
                // retries exhausted - save somewhere and move on to the next
            });
        });
    }
    

    从2.3.7版本开始,可以配置RetryingBatchErrorHandler,整个批次都会重试。

    从 2.5 版开始,您可以配置 RecoveringBatchErrorHandler,您可以在其中向批处理中的记录失败的处理程序抛出一个特殊异常,以便只重新传递未处理的记录。

    【讨论】:

    • 我很想详细了解这句话——“对于这种情况,最好直接在侦听器中使用 RetryTemplate。”因为我在 org.springframework.kafka.annotation.KafkaListener 接口中没有看到与 RetryTemplate 相关的任何内容
    • 我在答案中添加了一个示例。顺便说一句,在更新的版本中,您会得到一个带有此错误配置的 IllegalStateException,而不是 ClassCastExceptiongithub.com/spring-projects/spring-kafka/issues/1219
    猜你喜欢
    • 2021-05-06
    • 1970-01-01
    • 2019-11-25
    • 1970-01-01
    • 1970-01-01
    • 2017-08-13
    • 1970-01-01
    • 2019-12-31
    • 2017-04-20
    相关资源
    最近更新 更多