diff --git a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java index c9500a95..eb761d37 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarAutoConfigurationTests.java @@ -162,7 +162,7 @@ class PulsarAutoConfigurationTests { @Test void customPulsarListenerAnnotationBeanPostProcessorIsRespected() { - PulsarListenerAnnotationBeanPostProcessor listenerAnnotationBeanPostProcessor = mock( + PulsarListenerAnnotationBeanPostProcessor listenerAnnotationBeanPostProcessor = mock( PulsarListenerAnnotationBeanPostProcessor.class); this.contextRunner .withBean("org.springframework.pulsar.config.internalPulsarListenerAnnotationProcessor", diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListener.java b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListener.java index 29ac394b..2647326b 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListener.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListener.java @@ -186,4 +186,12 @@ public @interface PulsarListener { */ AckMode ackMode() default AckMode.BATCH; + /** + * The bean name or a 'SpEL' expression that resolves to a + * {@link org.springframework.pulsar.listener.PulsarConsumerErrorHandler} which is + * used as a Spring provided mechanism to handle errors from processing the message. + * @return the bean name for the consumer error handler or an empty string. + */ + String pulsarConsumerErrorHandler() default ""; + } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListenerAnnotationBeanPostProcessor.java b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListenerAnnotationBeanPostProcessor.java index 531af750..fc663201 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListenerAnnotationBeanPostProcessor.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/annotation/PulsarListenerAnnotationBeanPostProcessor.java @@ -85,6 +85,7 @@ import org.springframework.pulsar.config.PulsarListenerContainerFactory; import org.springframework.pulsar.config.PulsarListenerEndpoint; import org.springframework.pulsar.config.PulsarListenerEndpointRegistrar; import org.springframework.pulsar.config.PulsarListenerEndpointRegistry; +import org.springframework.pulsar.listener.PulsarConsumerErrorHandler; import org.springframework.util.Assert; import org.springframework.util.ReflectionUtils; import org.springframework.util.StringUtils; @@ -107,8 +108,7 @@ import org.springframework.validation.Validator; * fine-grained control over endpoints registration. See {@link EnablePulsar} Javadoc for * complete usage details. * - * @param the key type. - * @param the value type. + * @param the payload type. * @author Soby Chacko * @author Chris Bono * @author Alexander Preuß @@ -120,7 +120,7 @@ import org.springframework.validation.Validator; * @see PulsarListenerEndpoint * @see MethodPulsarListenerEndpoint */ -public class PulsarListenerAnnotationBeanPostProcessor +public class PulsarListenerAnnotationBeanPostProcessor implements BeanPostProcessor, Ordered, ApplicationContextAware, InitializingBean, SmartInitializingSingleton { private final LogAccessor logger = new LogAccessor(LogFactory.getLog(getClass())); @@ -366,6 +366,24 @@ public class PulsarListenerAnnotationBeanPostProcessor resolveNegativeAckRedeliveryBackoff(endpoint, pulsarListener); resolveDeadLetterPolicy(endpoint, pulsarListener); + resolvePulsarConsumerErrorHandler(endpoint, pulsarListener); + } + + @SuppressWarnings({ "rawtypes" }) + private void resolvePulsarConsumerErrorHandler(MethodPulsarListenerEndpoint endpoint, + PulsarListener pulsarListener) { + Object pulsarConsumerErrorHandler = resolveExpression(pulsarListener.pulsarConsumerErrorHandler()); + if (pulsarConsumerErrorHandler instanceof PulsarConsumerErrorHandler) { + endpoint.setPulsarConsumerErrorHandler((PulsarConsumerErrorHandler) pulsarConsumerErrorHandler); + } + else { + String pulsarConsumerErrorHandlerBeanName = resolveExpressionAsString( + pulsarListener.pulsarConsumerErrorHandler(), "pulsarConsumerErrorHandler"); + if (StringUtils.hasText(pulsarConsumerErrorHandlerBeanName)) { + endpoint.setPulsarConsumerErrorHandler( + this.beanFactory.getBean(pulsarConsumerErrorHandlerBeanName, PulsarConsumerErrorHandler.class)); + } + } } private void resolveNegativeAckRedeliveryBackoff(MethodPulsarListenerEndpoint endpoint, diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java index 3c05f6fb..a3e9c387 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/MethodPulsarListenerEndpoint.java @@ -46,6 +46,7 @@ import org.springframework.messaging.handler.invocation.InvocableHandlerMethod; import org.springframework.pulsar.core.SchemaUtils; import org.springframework.pulsar.listener.Acknowledgement; import org.springframework.pulsar.listener.ConcurrentPulsarMessageListenerContainer; +import org.springframework.pulsar.listener.PulsarConsumerErrorHandler; import org.springframework.pulsar.listener.PulsarContainerProperties; import org.springframework.pulsar.listener.PulsarMessageListenerContainer; import org.springframework.pulsar.listener.adapter.HandlerAdapter; @@ -83,6 +84,9 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo private DeadLetterPolicy deadLetterPolicy; + @SuppressWarnings("rawtypes") + private PulsarConsumerErrorHandler pulsarConsumerErrorHandler; + public void setBean(Object bean) { this.bean = bean; } @@ -189,6 +193,7 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo container.setNegativeAckRedeliveryBackoff(this.negativeAckRedeliveryBackoff); container.setDeadLetterPolicy(this.deadLetterPolicy); + container.setPulsarConsumerErrorHandler(this.pulsarConsumerErrorHandler); return messageListener; } @@ -271,4 +276,9 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo this.deadLetterPolicy = deadLetterPolicy; } + @SuppressWarnings("rawtypes") + public void setPulsarConsumerErrorHandler(PulsarConsumerErrorHandler pulsarConsumerErrorHandler) { + this.pulsarConsumerErrorHandler = pulsarConsumerErrorHandler; + } + } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/AbstractPulsarMessageListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/AbstractPulsarMessageListenerContainer.java index 7fec5375..e79e3286 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/AbstractPulsarMessageListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/AbstractPulsarMessageListenerContainer.java @@ -65,7 +65,7 @@ public abstract class AbstractPulsarMessageListenerContainer implements Pulsa protected DeadLetterPolicy deadLetterPolicy; - private PulsarConsumerErrorHandler pulsarConsumerErrorHandler; + protected PulsarConsumerErrorHandler pulsarConsumerErrorHandler; @SuppressWarnings("unchecked") protected AbstractPulsarMessageListenerContainer(PulsarConsumerFactory pulsarConsumerFactory, @@ -204,7 +204,8 @@ public abstract class AbstractPulsarMessageListenerContainer implements Pulsa return this.pulsarConsumerErrorHandler; } - public void setPulsarConsumerErrorHandler(PulsarConsumerErrorHandler pulsarConsumerErrorHandler) { + @SuppressWarnings({ "rawtypes", "unchecked" }) + public void setPulsarConsumerErrorHandler(PulsarConsumerErrorHandler pulsarConsumerErrorHandler) { this.pulsarConsumerErrorHandler = pulsarConsumerErrorHandler; } diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/ConcurrentPulsarMessageListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/ConcurrentPulsarMessageListenerContainer.java index e6f74847..5fee1289 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/ConcurrentPulsarMessageListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/ConcurrentPulsarMessageListenerContainer.java @@ -111,6 +111,7 @@ public class ConcurrentPulsarMessageListenerContainer extends AbstractPulsarM } container.setNegativeAckRedeliveryBackoff(this.negativeAckRedeliveryBackoff); container.setDeadLetterPolicy(this.deadLetterPolicy); + container.setPulsarConsumerErrorHandler(this.pulsarConsumerErrorHandler); } @Override diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarMessageListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarMessageListenerContainer.java index 547b1336..a9b28d9f 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarMessageListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarMessageListenerContainer.java @@ -49,4 +49,7 @@ public interface PulsarMessageListenerContainer extends SmartLifecycle, Disposab void setDeadLetterPolicy(DeadLetterPolicy deadLetterPolicy); + @SuppressWarnings("rawtypes") + void setPulsarConsumerErrorHandler(PulsarConsumerErrorHandler pulsarConsumerErrorHandler); + } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/ConcurrentPulsarMessageListenerContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/ConcurrentPulsarMessageListenerContainerTests.java index 8f1dd868..9025641e 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/ConcurrentPulsarMessageListenerContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/ConcurrentPulsarMessageListenerContainerTests.java @@ -39,6 +39,7 @@ import org.apache.pulsar.client.impl.MultiplierRedeliveryBackoff; import org.junit.jupiter.api.Test; import org.springframework.pulsar.core.PulsarConsumerFactory; +import org.springframework.util.backoff.BackOff; /** * @author Soby Chacko @@ -74,6 +75,22 @@ public class ConcurrentPulsarMessageListenerContainerTests { assertThat(childContainer.getNegativeAckRedeliveryBackoff()).isEqualTo(redeliveryBackoff); } + @Test + @SuppressWarnings({ "unchecked", "rawtypes" }) + void pulsarConsumerErrorHandlerAppliedOnChildContainer() throws Exception { + PulsarListenerMockComponents env = setupListenerMockComponents(SubscriptionType.Shared); + ConcurrentPulsarMessageListenerContainer concurrentContainer = env.concurrentContainer(); + + PulsarConsumerErrorHandler pulsarConsumerErrorHandler = new DefaultPulsarConsumerErrorHandler( + mock(PulsarMessageRecovererFactory.class), mock(BackOff.class)); + concurrentContainer.setPulsarConsumerErrorHandler(pulsarConsumerErrorHandler); + + concurrentContainer.start(); + + final DefaultPulsarMessageListenerContainer childContainer = concurrentContainer.getContainers().get(0); + assertThat(childContainer.getPulsarConsumerErrorHandler()).isEqualTo(pulsarConsumerErrorHandler); + } + @Test @SuppressWarnings("unchecked") void basicConcurrencyTesting() throws Exception { diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java index 37a0f4b7..5191f463 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/PulsarListenerTests.java @@ -66,6 +66,7 @@ import org.springframework.pulsar.core.PulsarTopic; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; +import org.springframework.util.backoff.FixedBackOff; /** * @author Soby Chacko @@ -279,6 +280,50 @@ public class PulsarListenerTests extends AbstractContainerBaseTests { } + @Nested + @ContextConfiguration(classes = PulsarConsumerErrorHandlerTest.PulsarConsumerErrorHandlerConfig.class) + class PulsarConsumerErrorHandlerTest { + + private static CountDownLatch pulsarConsumerErrorHandlerLatch = new CountDownLatch(11); + + private static CountDownLatch dltLatch = new CountDownLatch(1); + + @Test + void pulsarListenerWithNackRedeliveryBackoff(@Autowired PulsarListenerEndpointRegistry registry) + throws Exception { + pulsarTemplate.send("pceht-topic", "hello john doe"); + assertThat(pulsarConsumerErrorHandlerLatch.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(dltLatch.await(10, TimeUnit.SECONDS)).isTrue(); + } + + @EnablePulsar + @Configuration + static class PulsarConsumerErrorHandlerConfig { + + @PulsarListener(id = "pceht-id", subscriptionName = "pceht-subscription", topics = "pceht-topic", + pulsarConsumerErrorHandler = "pulsarConsumerErrorHandler") + void listen(String msg) { + pulsarConsumerErrorHandlerLatch.countDown(); + throw new RuntimeException("fail " + msg); + } + + @PulsarListener(id = "pceh-dltListener", subscriptionType = "dltListenerSubscription", + topics = "pceht-topic-pceht-subscription-DLT") + void listenDlq(String msg) { + dltLatch.countDown(); + } + + @Bean + public PulsarConsumerErrorHandler pulsarConsumerErrorHandler( + PulsarTemplate pulsarTemplate) { + return new DefaultPulsarConsumerErrorHandler<>( + new PulsarDeadLetterPublishingRecoverer<>(pulsarTemplate), new FixedBackOff(100, 10)); + } + + } + + } + @Nested @ContextConfiguration(classes = DeadLetterPolicyTest.DeadLetterPolicyConfig.class) class DeadLetterPolicyTest {