diff --git a/spring-pulsar-docs/src/main/asciidoc/pulsar.adoc b/spring-pulsar-docs/src/main/asciidoc/pulsar.adoc index b32c82a6..db461c26 100644 --- a/spring-pulsar-docs/src/main/asciidoc/pulsar.adoc +++ b/spring-pulsar-docs/src/main/asciidoc/pulsar.adoc @@ -536,17 +536,17 @@ Here is an example. ==== [source, java] ---- -@PulsarListener(subscriptionName = "hello-pulsar-partitioned-subscription", topics = "hello-pulsar-partitioned", subscriptionType = "failover") +@PulsarListener(subscriptionName = "hello-pulsar-partitioned-subscription", topics = "hello-pulsar-partitioned", subscriptionType = SubscriptionType.Failover) public void listen1(String foo) { System.out.println("Message Received 1: " + foo); } -@PulsarListener(subscriptionName = "hello-pulsar-partitioned-subscription", topics = "hello-pulsar-partitioned", subscriptionType = "failover") +@PulsarListener(subscriptionName = "hello-pulsar-partitioned-subscription", topics = "hello-pulsar-partitioned", subscriptionType = SubscriptionType.Failover) public void listen2(String foo) { System.out.println("Message Received 2: " + foo); } -@PulsarListener(subscriptionName = "hello-pulsar-partitioned-subscription", topics = "hello-pulsar-partitioned", subscriptionType = "failover") +@PulsarListener(subscriptionName = "hello-pulsar-partitioned-subscription", topics = "hello-pulsar-partitioned", subscriptionType = SubscriptionType.Failover) public void listen3(String foo) { System.out.println("Message Received 3: " + foo); } @@ -563,12 +563,12 @@ Here is an example. ==== [source, java] ---- -@PulsarListener(subscriptionName = "hello-pulsar-shared-subscription", topics = "hello-pulsar-partitioned", subscriptionType = "shared") +@PulsarListener(subscriptionName = "hello-pulsar-shared-subscription", topics = "hello-pulsar-partitioned", subscriptionType = SubscriptionType.Shared) public void listen1(String foo) { System.out.println("Message Received 1: " + foo); } -@PulsarListener(subscriptionName = "hello-pulsar-shared-subscription", topics = "hello-pulsar-partitioned", subscriptionType = "shared") +@PulsarListener(subscriptionName = "hello-pulsar-shared-subscription", topics = "hello-pulsar-partitioned", subscriptionType = SubscriptionType.Shared) public void listen2(String foo) { System.out.println("Message Received 2: " + foo); } diff --git a/spring-pulsar-sample-apps/src/main/java/app2/FailoverConsumerApp.java b/spring-pulsar-sample-apps/src/main/java/app2/FailoverConsumerApp.java index e82b7243..1beec714 100644 --- a/spring-pulsar-sample-apps/src/main/java/app2/FailoverConsumerApp.java +++ b/spring-pulsar-sample-apps/src/main/java/app2/FailoverConsumerApp.java @@ -20,6 +20,7 @@ import java.io.Serial; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageRouter; +import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.api.TopicMetadata; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -57,19 +58,19 @@ public class FailoverConsumerApp { } @PulsarListener(subscriptionName = "failover-subscription-demo", topics = "failover-demo-topic", - subscriptionType = "failover") + subscriptionType = SubscriptionType.Failover) void listen1(String foo) { this.logger.info("failover-listen1 : " + foo); } @PulsarListener(subscriptionName = "failover-subscription-demo", topics = "failover-demo-topic", - subscriptionType = "failover") + subscriptionType = SubscriptionType.Failover) void listen2(String foo) { this.logger.info("failover-listen2 : " + foo); } @PulsarListener(subscriptionName = "failover-subscription-demo", topics = "failover-demo-topic", - subscriptionType = "failover") + subscriptionType = SubscriptionType.Failover) void listen(String foo) { this.logger.info("failover-listen3 : " + foo); } 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 ca9b5101..399d2897 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 @@ -22,6 +22,7 @@ import java.lang.annotation.Retention; import java.lang.annotation.RetentionPolicy; import java.lang.annotation.Target; +import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.common.schema.SchemaType; import org.springframework.messaging.handler.annotation.MessageMapping; @@ -76,7 +77,7 @@ public @interface PulsarListener { * Pulsar subscription type for this listener. * @return the {@code subscriptionType} for this listener */ - String subscriptionType() default ""; + SubscriptionType subscriptionType() default SubscriptionType.Exclusive; SchemaType schemaType() default SchemaType.NONE; 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 0830fa7c..96190716 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 @@ -36,12 +36,10 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.BiFunction; -import java.util.stream.Collectors; import org.apache.commons.logging.LogFactory; import org.apache.pulsar.client.api.DeadLetterPolicy; import org.apache.pulsar.client.api.RedeliveryBackoff; -import org.apache.pulsar.client.api.SubscriptionType; import org.springframework.aop.framework.Advised; import org.springframework.aop.support.AopUtils; @@ -347,7 +345,7 @@ public class PulsarListenerAnnotationBeanPostProcessor endpoint.setId(getEndpointId(pulsarListener)); endpoint.setTopics(topics); endpoint.setTopicPattern(topicPattern); - endpoint.setSubscriptionType(getEndpointSubscriptionType(pulsarListener)); + endpoint.setSubscriptionType(pulsarListener.subscriptionType()); endpoint.setSchemaType(pulsarListener.schemaType()); endpoint.setAckMode(pulsarListener.ackMode()); @@ -511,32 +509,14 @@ public class PulsarListenerAnnotationBeanPostProcessor if (StringUtils.hasText(pulsarListener.subscriptionName())) { return resolveExpressionAsString(pulsarListener.subscriptionName(), "subscriptionName"); } - else { - return GENERATED_ID_PREFIX + this.counter.getAndIncrement(); - } - } - - private SubscriptionType getEndpointSubscriptionType(PulsarListener pulsarListener) { - final String subscriptionType = pulsarListener.subscriptionType().toLowerCase(); - if (StringUtils.hasText(subscriptionType)) { - return switch (subscriptionType) { - case "exclusive" -> SubscriptionType.Exclusive; - case "failover" -> SubscriptionType.Failover; - case "shared" -> SubscriptionType.Shared; - case "key_shared" -> SubscriptionType.Key_Shared; - default -> SubscriptionType.Exclusive; - }; - } - return null; + return GENERATED_ID_PREFIX + this.counter.getAndIncrement(); } private String getEndpointId(PulsarListener pulsarListener) { if (StringUtils.hasText(pulsarListener.id())) { return resolveExpressionAsString(pulsarListener.id(), "id"); } - else { - return GENERATED_ID_PREFIX + this.counter.getAndIncrement(); - } + return GENERATED_ID_PREFIX + this.counter.getAndIncrement(); } private String getTopicPattern(PulsarListener pulsarListener) { @@ -642,8 +622,7 @@ public class PulsarListenerAnnotationBeanPostProcessor } PulsarListeners anns = AnnotationUtils.findAnnotation(clazz, PulsarListeners.class); if (anns != null) { - listeners - .addAll(Arrays.stream(anns.value()).map(anno -> enhance(clazz, anno)).collect(Collectors.toList())); + listeners.addAll(Arrays.stream(anns.value()).map(anno -> enhance(clazz, anno)).toList()); } return listeners; } @@ -657,8 +636,7 @@ public class PulsarListenerAnnotationBeanPostProcessor } PulsarListeners anns = AnnotationUtils.findAnnotation(method, PulsarListeners.class); if (anns != null) { - listeners.addAll( - Arrays.stream(anns.value()).map(anno -> enhance(method, anno)).collect(Collectors.toList())); + listeners.addAll(Arrays.stream(anns.value()).map(anno -> enhance(method, anno)).toList()); } return listeners; } @@ -667,11 +645,8 @@ public class PulsarListenerAnnotationBeanPostProcessor if (this.enhancer == null) { return ann; } - else { - return AnnotationUtils.synthesizeAnnotation( - this.enhancer.apply(AnnotationUtils.getAnnotationAttributes(ann), element), PulsarListener.class, - null); - } + return AnnotationUtils.synthesizeAnnotation( + this.enhancer.apply(AnnotationUtils.getAnnotationAttributes(ann), element), PulsarListener.class, null); } private void addFormatters(FormatterRegistry registry) { @@ -817,9 +792,7 @@ public class PulsarListenerAnnotationBeanPostProcessor || target.equals(byte.class) || target.equals(Long.class) || target.equals(Integer.class) || target.equals(Short.class) || target.equals(Byte.class); } - else { - return false; - } + return false; } } 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 e2df4570..8d5ff3c9 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 @@ -34,6 +34,7 @@ import org.apache.pulsar.client.api.Messages; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.RedeliveryBackoff; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.impl.MultiplierRedeliveryBackoff; import org.apache.pulsar.client.impl.schema.AvroSchema; import org.apache.pulsar.client.impl.schema.JSONSchema; @@ -224,7 +225,7 @@ public class PulsarListenerTests extends AbstractContainerBaseTests { } @PulsarListener(id = "bar", topics = "concurrency-on-pl", subscriptionName = "subscription-2", - subscriptionType = "failover", concurrency = "3") + subscriptionType = SubscriptionType.Failover, concurrency = "3") void listen2(String message) { latch1.countDown(); } @@ -264,7 +265,7 @@ public class PulsarListenerTests extends AbstractContainerBaseTests { @PulsarListener(id = "withNegRedeliveryBackoff", subscriptionName = "withNegRedeliveryBackoffSubscription", topics = "withNegRedeliveryBackoff-test-topic", negativeAckRedeliveryBackoff = "redeliveryBackoff", - subscriptionType = "Shared") + subscriptionType = SubscriptionType.Shared) void listen(String msg) { nackRedeliveryBackoffLatch.countDown(); throw new RuntimeException("fail " + msg); @@ -300,8 +301,8 @@ public class PulsarListenerTests extends AbstractContainerBaseTests { @PulsarListener(id = "withAckTimeoutRedeliveryBackoff", subscriptionName = "withAckTimeoutRedeliveryBackoffSubscription", topics = "withAckTimeoutRedeliveryBackoff-test-topic", - ackTimeoutRedeliveryBackoff = "ackTimeoutRedeliveryBackoff", subscriptionType = "Shared", - properties = { "ackTimeoutMillis=1" }) + ackTimeoutRedeliveryBackoff = "ackTimeoutRedeliveryBackoff", + subscriptionType = SubscriptionType.Shared, properties = { "ackTimeoutMillis=1" }) void listen(String msg) { ackTimeoutRedeliveryBackoffLatch.countDown(); throw new RuntimeException(); @@ -344,8 +345,7 @@ public class PulsarListenerTests extends AbstractContainerBaseTests { throw new RuntimeException("fail " + msg); } - @PulsarListener(id = "pceh-dltListener", subscriptionType = "dltListenerSubscription", - topics = "pceht-topic-pceht-subscription-DLT") + @PulsarListener(id = "pceh-dltListener", topics = "pceht-topic-pceht-subscription-DLT") void listenDlq(String msg) { dltLatch.countDown(); } @@ -381,14 +381,14 @@ public class PulsarListenerTests extends AbstractContainerBaseTests { static class DeadLetterPolicyConfig { @PulsarListener(id = "deadLetterPolicyListener", subscriptionName = "deadLetterPolicySubscription", - topics = "dlpt-topic-1", deadLetterPolicy = "deadLetterPolicy", subscriptionType = "Shared", - properties = { "ackTimeoutMillis=1" }) + topics = "dlpt-topic-1", deadLetterPolicy = "deadLetterPolicy", + subscriptionType = SubscriptionType.Shared, properties = { "ackTimeoutMillis=1" }) void listen(String msg) { latch.countDown(); throw new RuntimeException("fail " + msg); } - @PulsarListener(id = "dlqListener", subscriptionType = "dlqListenerSubscription", topics = "dlpt-dlq-topic") + @PulsarListener(id = "dlqListener", topics = "dlpt-dlq-topic") void listenDlq(String msg) { dlqLatch.countDown(); }