From e364dee85a397c5a4f34bcfea729d87343ae45aa Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Alexander=20Preu=C3=9F?= Date: Wed, 31 Aug 2022 21:01:32 +0200 Subject: [PATCH] Implement topic pattern on PulsarListener (#5) (#77) * Implement topicPattern parameter in PulsarListener * Fix Typos --- ...arListenerAnnotationBeanPostProcessor.java | 15 +++++++++---- .../AbstractPulsarListenerEndpoint.java | 18 +++++++++++++-- ...currentPulsarListenerContainerFactory.java | 6 +++++ .../pulsar/config/PulsarListenerEndpoint.java | 3 +++ .../config/PulsarListenerEndpointAdapter.java | 6 +++++ ...DefaultPulsarMessageListenerContainer.java | 15 ++++++++++--- .../listener/PulsarContainerProperties.java | 12 +++++----- ...tPulsarMessageListenerContainerTests.java} | 2 +- .../pulsar/listener/PulsarListenerTests.java | 22 +++++++++++++++++++ 9 files changed, 83 insertions(+), 16 deletions(-) rename spring-pulsar/src/test/java/org/springframework/pulsar/listener/{DefaultPulsarMessageListenerContainierTests.java => DefaultPulsarMessageListenerContainerTests.java} (98%) 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 56f980f8..c10279e6 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 @@ -109,6 +109,7 @@ import org.springframework.validation.Validator; * @param the value type. * @author Soby Chacko * @author Chris Bono + * @author Alexander Preuß * @see PulsarListener * @see EnablePulsar * @see PulsarListenerConfigurer @@ -277,14 +278,15 @@ public class PulsarListenerAnnotationBeanPostProcessor String beanRef = pulsarListener.beanRef(); this.listenerScope.addListener(beanRef, bean); String[] topics = resolveTopics(pulsarListener); - processListener(endpoint, pulsarListener, bean, beanName, topics); + String topicPattern = getTopicPattern(pulsarListener); + processListener(endpoint, pulsarListener, bean, beanName, topics, topicPattern); this.listenerScope.removeListener(beanRef); } protected void processListener(MethodPulsarListenerEndpoint endpoint, PulsarListener PulsarListener, Object bean, - String beanName, String[] topics) { + String beanName, String[] topics, String topicPattern) { - processPulsarListenerAnnotation(endpoint, PulsarListener, bean, topics); + processPulsarListenerAnnotation(endpoint, PulsarListener, bean, topics, topicPattern); String containerFactory = resolve(PulsarListener.containerFactory()); PulsarListenerContainerFactory listenerContainerFactory = resolveContainerFactory(PulsarListener, @@ -335,13 +337,14 @@ public class PulsarListenerAnnotationBeanPostProcessor } private void processPulsarListenerAnnotation(MethodPulsarListenerEndpoint endpoint, - PulsarListener pulsarListener, Object bean, String[] topics) { + PulsarListener pulsarListener, Object bean, String[] topics, String topicPattern) { endpoint.setBean(bean); endpoint.setMessageHandlerMethodFactory(this.messageHandlerMethodFactory); endpoint.setSubscriptionName(getEndpointSubscriptionName(pulsarListener)); endpoint.setId(getEndpointId(pulsarListener)); endpoint.setTopics(topics); + endpoint.setTopicPattern(topicPattern); endpoint.setSubscriptionType(getEndpointSubscriptionType(pulsarListener)); endpoint.setSchemaType(pulsarListener.schemaType()); @@ -464,6 +467,10 @@ public class PulsarListenerAnnotationBeanPostProcessor } } + private String getTopicPattern(PulsarListener pulsarListener) { + return resolveExpressionAsString(pulsarListener.topicPattern(), "topicPattern"); + } + private String resolveExpressionAsString(String value, String attribute) { Object resolved = resolveExpression(value); if (resolved instanceof String) { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarListenerEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarListenerEndpoint.java index 9edc3dcd..631e2763 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarListenerEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarListenerEndpoint.java @@ -41,12 +41,14 @@ import org.springframework.pulsar.listener.PulsarMessageListenerContainer; import org.springframework.pulsar.listener.adapter.PulsarMessagingMessageListenerAdapter; import org.springframework.pulsar.support.MessageConverter; import org.springframework.util.Assert; +import org.springframework.util.StringUtils; /** * Base implementation for {@link PulsarListenerEndpoint}. * * @param Message payload type. * @author Soby Chacko + * @author Alexander Preuß */ public abstract class AbstractPulsarListenerEndpoint implements PulsarListenerEndpoint, BeanFactoryAware, InitializingBean { @@ -63,6 +65,8 @@ public abstract class AbstractPulsarListenerEndpoint private final Collection topics = new ArrayList<>(); + private String topicPattern; + private BeanFactory beanFactory; private BeanExpressionResolver resolver; @@ -97,8 +101,8 @@ public abstract class AbstractPulsarListenerEndpoint @Override public void afterPropertiesSet() { boolean topicsEmpty = getTopics().isEmpty(); - if (!topicsEmpty) { - throw new IllegalStateException("Topics or topicPartitions must be provided but not both for " + this); + if (!topicsEmpty && !StringUtils.hasText(getTopicPattern())) { + throw new IllegalStateException("Topics or topicPattern must be provided but not both for " + this); } } @@ -148,6 +152,16 @@ public abstract class AbstractPulsarListenerEndpoint return Collections.unmodifiableCollection(this.topics); } + public void setTopicPattern(String topicPattern) { + Assert.notNull(topicPattern, "'topicPattern' must not be null"); + this.topicPattern = topicPattern; + } + + @Override + public String getTopicPattern() { + return this.topicPattern; + } + @Override @Nullable public Boolean getAutoStartup() { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactory.java index 8deb09be..e74dfd6e 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactory.java @@ -31,6 +31,7 @@ import org.springframework.util.StringUtils; * @param message type in the listener. * @author Soby Chacko * @author Chris Bono + * @author Alexander Preuß */ public class ConcurrentPulsarListenerContainerFactory extends AbstractPulsarListenerContainerFactory, T> { @@ -50,12 +51,17 @@ public class ConcurrentPulsarListenerContainerFactory PulsarContainerProperties properties = new PulsarContainerProperties(); Collection topics = endpoint.getTopics(); + String topicPattern = endpoint.getTopicPattern(); if (!topics.isEmpty()) { final String[] topics1 = topics.toArray(new String[0]); properties.setTopics(topics1); } + if (StringUtils.hasText(topicPattern)) { + properties.setTopicsPattern(topicPattern); + } + final String subscriptionName = endpoint.getSubscriptionName(); if (StringUtils.hasText(subscriptionName)) { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpoint.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpoint.java index 3c951591..4b1d470c 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpoint.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpoint.java @@ -32,6 +32,7 @@ import org.springframework.pulsar.support.MessageConverter; * endpoints programmatically. * * @author Soby Chacko + * @author Alexander Preuß */ public interface PulsarListenerEndpoint { @@ -46,6 +47,8 @@ public interface PulsarListenerEndpoint { Collection getTopics(); + String getTopicPattern(); + @Nullable Boolean getAutoStartup(); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpointAdapter.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpointAdapter.java index 7f5ec3d9..9413c0e1 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpointAdapter.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/PulsarListenerEndpointAdapter.java @@ -31,6 +31,7 @@ import org.springframework.pulsar.support.MessageConverter; * Adapter to avoid having to implement all methods. * * @author Soby Chacko + * @author Alexander Preuß */ public class PulsarListenerEndpointAdapter implements PulsarListenerEndpoint { @@ -54,6 +55,11 @@ public class PulsarListenerEndpointAdapter implements PulsarListenerEndpoint { return Collections.emptyList(); } + @Override + public String getTopicPattern() { + return null; + } + @Override public Boolean getAutoStartup() { return null; diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java index f05dce4f..91653790 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainer.java @@ -56,6 +56,7 @@ import org.springframework.util.StringUtils; * * @param message type. * @author Soby Chacko + * @author Alexander Preuß */ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMessageListenerContainer { @@ -223,6 +224,14 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess propertiesToOverride.put("topicNames", listenerDefinedTopics); } } + + if (!propertiesToOverride.containsKey("topicsPattern")) { + final String topicsPattern = pulsarContainerProperties.getTopicsPattern(); + if (topicsPattern != null) { + propertiesToOverride.put("topicsPattern", topicsPattern); + } + } + if (!propertiesToOverride.containsKey("subscriptionName")) { if (StringUtils.hasText(pulsarContainerProperties.getSubscriptionName())) { propertiesToOverride.put("subscriptionName", pulsarContainerProperties.getSubscriptionName()); @@ -267,7 +276,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } if (this.containerProperties.getAckMode() == PulsarContainerProperties.AckMode.BATCH) { try { - if (isSharedSubsriptionType()) { + if (isSharedSubscriptionType()) { this.consumer.acknowledge(messages); } else { @@ -322,7 +331,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } } - private boolean isSharedSubsriptionType() { + private boolean isSharedSubscriptionType() { return this.containerProperties.getSubscriptionType() == SubscriptionType.Shared || this.containerProperties.getSubscriptionType() == SubscriptionType.Key_Shared; } @@ -331,7 +340,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess if (this.nackableMessages.isEmpty()) { try { if (messages.size() > 0) { - if (isSharedSubsriptionType()) { + if (isSharedSubscriptionType()) { this.consumer.acknowledge(messages); } else { diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarContainerProperties.java b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarContainerProperties.java index 97706264..214fb690 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarContainerProperties.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/listener/PulsarContainerProperties.java @@ -18,7 +18,6 @@ package org.springframework.pulsar.listener; import java.time.Duration; import java.util.Properties; -import java.util.regex.Pattern; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; @@ -31,6 +30,7 @@ import org.springframework.util.Assert; * Contains runtime properties for a listener container. * * @author Soby Chacko + * @author Alexander Preuß */ public class PulsarContainerProperties { @@ -60,7 +60,7 @@ public class PulsarContainerProperties { private String[] topics; - private Pattern topicsPattern; + private String topicsPattern; private String subscriptionName; @@ -91,7 +91,7 @@ public class PulsarContainerProperties { this.topicsPattern = null; } - public PulsarContainerProperties(Pattern topicPattern) { + public PulsarContainerProperties(String topicPattern) { this.topicsPattern = topicPattern; this.topics = null; } @@ -170,7 +170,7 @@ public class PulsarContainerProperties { * @param consumerStartTimeout the consumer start timeout. */ public void setConsumerStartTimeout(Duration consumerStartTimeout) { - Assert.notNull(consumerStartTimeout, "'consumerStartTimout' cannot be null"); + Assert.notNull(consumerStartTimeout, "'consumerStartTimeout' cannot be null"); this.consumerStartTimeout = consumerStartTimeout; } @@ -190,11 +190,11 @@ public class PulsarContainerProperties { this.topics = topics; } - public Pattern getTopicsPattern() { + public String getTopicsPattern() { return this.topicsPattern; } - public void setTopicsPattern(Pattern topicsPattern) { + public void setTopicsPattern(String topicsPattern) { this.topicsPattern = topicsPattern; } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainierTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java similarity index 98% rename from spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainierTests.java rename to spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java index 70593350..2a00b07c 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainierTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java @@ -39,7 +39,7 @@ import org.springframework.pulsar.core.PulsarTemplate; /** * @author Soby Chacko */ -class DefaultPulsarMessageListenerContainierTests extends AbstractContainerBaseTests { +class DefaultPulsarMessageListenerContainerTests extends AbstractContainerBaseTests { @Test void basicDefaultConsumer() 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 0550dfb8..704a75be 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 @@ -55,6 +55,7 @@ import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; /** * @author Soby Chacko + * @author Alexander Preuß */ @SpringJUnitConfig @DirtiesContext @@ -62,6 +63,7 @@ public class PulsarListenerTests extends AbstractContainerBaseTests { static CountDownLatch latch = new CountDownLatch(1); static CountDownLatch latch1 = new CountDownLatch(3); + static CountDownLatch latch2 = new CountDownLatch(3); @Autowired PulsarTemplate pulsarTemplate; @@ -125,6 +127,19 @@ public class PulsarListenerTests extends AbstractContainerBaseTests { @ContextConfiguration(classes = PulsarListenerBasicTestCases.TestPulsarListenersForBasicScenario.class) class PulsarListenerBasicTestCases { + @Test + void testPulsarListenerWithTopicsPattern(@Autowired PulsarListenerEndpointRegistry registry) throws Exception { + PulsarMessageListenerContainer baz = registry.getListenerContainer("baz"); + PulsarContainerProperties containerProperties = baz.getContainerProperties(); + assertThat(containerProperties.getTopicsPattern()).isEqualTo("persistent://public/default/pattern.*"); + + pulsarTemplate.send("persistent://public/default/pattern-1", "hello baz"); + pulsarTemplate.send("persistent://public/default/pattern-2", "hello baz"); + pulsarTemplate.send("persistent://public/default/pattern-3", "hello baz"); + + assertThat(latch2.await(10, TimeUnit.SECONDS)).isTrue(); + } + @Test void testPulsarListenerProvidedConsumerProperties(@Autowired PulsarListenerEndpointRegistry registry) throws Exception { @@ -192,6 +207,13 @@ public class PulsarListenerTests extends AbstractContainerBaseTests { latch1.countDown(); } + @PulsarListener(id = "baz", topicPattern = "persistent://public/default/pattern.*", + subscriptionName = "subscription-3", + properties = { "patternAutoDiscoveryPeriod=5", "subscriptionInitialPosition=Earliest" }) + void listen3(String message) { + latch2.countDown(); + } + } }