From b279945084f9d78a3ce0de4d677deb02b7ff821b Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 14 Dec 2018 13:11:30 -0500 Subject: [PATCH] Process errorHandler in class level KafkaListener The `errorHandler` attributed has been missed from the `KafkaListenerAnnotationBeanPostProcessor.processMultiMethodListeners()` logic. * Move `errorHandler` and `BeanFactory` population into the `processListener()` method in the `KafkaListenerAnnotationBeanPostProcessor` **Cherry-pick to 2.1.x, 2.0.x & 1.3.x** --- ...kaListenerAnnotationBeanPostProcessor.java | 14 ++++----- .../EnableKafkaIntegrationTests.java | 29 +++++++++++++++---- 2 files changed, 31 insertions(+), 12 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java index 1914c4f1..019c26c1 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/annotation/KafkaListenerAnnotationBeanPostProcessor.java @@ -352,7 +352,8 @@ public class KafkaListenerAnnotationBeanPostProcessor Method defaultMethod = null; for (Method method : multiMethods) { Method checked = checkProxy(method, bean); - if (AnnotationUtils.findAnnotation(method, KafkaHandler.class).isDefault()) { // NOSONAR never null + KafkaHandler annotation = AnnotationUtils.findAnnotation(method, KafkaHandler.class); + if (annotation != null && annotation.isDefault()) { final Method toAssert = defaultMethod; Assert.state(toAssert == null, () -> "Only one @KafkaHandler can be marked 'isDefault', found: " + toAssert.toString() + " and " + method.toString()); @@ -363,7 +364,6 @@ public class KafkaListenerAnnotationBeanPostProcessor for (KafkaListener classLevelListener : classLevelListeners) { MultiMethodKafkaListenerEndpoint endpoint = new MultiMethodKafkaListenerEndpoint<>(checkedMethods, defaultMethod, bean); - endpoint.setBeanFactory(this.beanFactory); processListener(endpoint, classLevelListener, bean, bean.getClass(), beanName); } } @@ -372,11 +372,6 @@ public class KafkaListenerAnnotationBeanPostProcessor Method methodToUse = checkProxy(method, bean); MethodKafkaListenerEndpoint endpoint = new MethodKafkaListenerEndpoint<>(); endpoint.setMethod(methodToUse); - endpoint.setBeanFactory(this.beanFactory); - String errorHandlerBeanName = resolveExpressionAsString(kafkaListener.errorHandler(), "errorHandler"); - if (StringUtils.hasText(errorHandlerBeanName)) { - endpoint.setErrorHandler(this.beanFactory.getBean(errorHandlerBeanName, KafkaListenerErrorHandler.class)); - } processListener(endpoint, kafkaListener, bean, methodToUse, beanName); } @@ -459,6 +454,11 @@ public class KafkaListenerAnnotationBeanPostProcessor } } + endpoint.setBeanFactory(this.beanFactory); + String errorHandlerBeanName = resolveExpressionAsString(kafkaListener.errorHandler(), "errorHandler"); + if (StringUtils.hasText(errorHandlerBeanName)) { + endpoint.setErrorHandler(this.beanFactory.getBean(errorHandlerBeanName, KafkaListenerErrorHandler.class)); + } this.registrar.registerEndpoint(endpoint, factory); if (StringUtils.hasText(beanRef)) { this.listenerScope.removeListener(beanRef); diff --git a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java index aec1a316..6f3b2048 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java @@ -341,14 +341,17 @@ public class EnableKafkaIntegrationTests { ConsumerFactory cf = new DefaultKafkaConsumerFactory<>(consumerProps); Consumer consumer = cf.createConsumer(); embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "annotated8reply"); - template.send("annotated8", 0, 1, "foo"); - template.send("annotated8", 0, 1, null); - template.flush(); + this.template.send("annotated8", 0, 1, "foo"); + this.template.send("annotated8", 0, 1, null); + this.template.flush(); assertThat(this.multiListener.latch1.await(60, TimeUnit.SECONDS)).isTrue(); assertThat(this.multiListener.latch2.await(60, TimeUnit.SECONDS)).isTrue(); ConsumerRecord reply = KafkaTestUtils.getSingleRecord(consumer, "annotated8reply"); assertThat(reply.value()).isEqualTo("OK"); consumer.close(); + + template.send("annotated8", 0, 1, "junk"); + assertThat(this.multiListener.errorLatch.await(60, TimeUnit.SECONDS)).isTrue(); } @Test @@ -1259,6 +1262,15 @@ public class EnableKafkaIntegrationTests { }); } + + @Bean + public KafkaListenerErrorHandler consumeMultiMethodException(MultiListenerBean listener) { + return (m, e) -> { + listener.errorLatch.countDown(); + return null; + }; + } + } @Component @@ -1667,16 +1679,23 @@ public class EnableKafkaIntegrationTests { } - @KafkaListener(id = "multi", topics = "annotated8") + @KafkaListener(id = "multi", topics = "annotated8", errorHandler = "consumeMultiMethodException") static class MultiListenerBean { private final CountDownLatch latch1 = new CountDownLatch(1); private final CountDownLatch latch2 = new CountDownLatch(1); + private final CountDownLatch errorLatch = new CountDownLatch(1); + @KafkaHandler public void bar(@NonNull String bar) { - this.latch1.countDown(); + if ("junk".equals(bar)) { + throw new RuntimeException("intentional"); + } + else { + this.latch1.countDown(); + } } @KafkaHandler