diff --git a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java index 2a353f9c..288c733f 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java @@ -17,6 +17,7 @@ package org.springframework.pulsar.spring.cloud.stream.binder; import java.util.Optional; +import java.util.Set; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; @@ -137,7 +138,7 @@ public class PulsarMessageChannelBinder extends protected MessageProducer createConsumerEndpoint(ConsumerDestination destination, String group, ExtendedConsumerProperties properties) { var containerProperties = new PulsarContainerProperties(); - containerProperties.setTopics(new String[] { destination.getName() }); + containerProperties.setTopics(Set.of(destination.getName())); var inboundHeaderMapper = determineInboundHeaderMapper(properties); 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 a406ecd9..8eac0ffb 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 @@ -18,6 +18,7 @@ package org.springframework.pulsar.config; import java.util.Arrays; import java.util.Collection; +import java.util.HashSet; import org.springframework.lang.Nullable; import org.springframework.pulsar.core.PulsarConsumerFactory; @@ -62,7 +63,7 @@ public class ConcurrentPulsarListenerContainerFactory properties.setTopicResolver(this.getContainerProperties().getTopicResolver()); if (!CollectionUtils.isEmpty(endpoint.getTopics())) { - properties.setTopics(endpoint.getTopics().toArray(new String[0])); + properties.setTopics(new HashSet<>(endpoint.getTopics())); } if (StringUtils.hasText(endpoint.getTopicPattern())) { 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 def261e5..df5536c0 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 @@ -20,6 +20,7 @@ import java.lang.reflect.Method; import java.util.Arrays; import java.util.List; import java.util.Optional; +import java.util.Set; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.DeadLetterPolicy; @@ -159,7 +160,7 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo || StringUtils.hasText(pulsarContainerProperties.getTopicsPattern()); if (!hasTopicInfo) { topicResolver.resolveTopic(null, messageType.getRawClass(), () -> null) - .ifResolved((topic) -> pulsarContainerProperties.setTopics(new String[] { topic })); + .ifResolved((topic) -> pulsarContainerProperties.setTopics(Set.of(topic))); } container.setNegativeAckRedeliveryBackoff(this.negativeAckRedeliveryBackoff); 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 f3ae0029..9a6ad280 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 @@ -17,7 +17,6 @@ package org.springframework.pulsar.listener; import java.util.ArrayList; -import java.util.Arrays; import java.util.HashMap; import java.util.HashSet; import java.util.List; @@ -327,9 +326,8 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess } } if (!currentProperties.containsKey("topicNames")) { - String[] topics = this.containerProperties.getTopics(); - Set listenerDefinedTopics = new HashSet<>(Arrays.stream(topics).toList()); - if (!listenerDefinedTopics.isEmpty()) { + Set listenerDefinedTopics = this.containerProperties.getTopics(); + if (!this.containerProperties.getTopics().isEmpty()) { currentProperties.put("topicNames", listenerDefinedTopics); } } 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 dc092a3a..a28a5a4e 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,6 +18,7 @@ package org.springframework.pulsar.listener; import java.time.Duration; import java.util.Properties; +import java.util.Set; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; @@ -48,7 +49,7 @@ public class PulsarContainerProperties { private Duration consumerStartTimeout = DEFAULT_CONSUMER_START_TIMEOUT; - private String[] topics; + private Set topics; private String topicsPattern; @@ -83,7 +84,7 @@ public class PulsarContainerProperties { private Properties pulsarConsumerProperties = new Properties(); public PulsarContainerProperties(String... topics) { - this.topics = topics.clone(); + this.topics = Set.of(topics); this.topicsPattern = null; this.schemaResolver = new DefaultSchemaResolver(); this.topicResolver = new DefaultTopicResolver(); @@ -186,11 +187,11 @@ public class PulsarContainerProperties { this.consumerStartTimeout = consumerStartTimeout; } - public String[] getTopics() { + public Set getTopics() { return this.topics; } - public void setTopics(String[] topics) { + public void setTopics(Set topics) { this.topics = topics; } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java index 1db59dcc..3816207d 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/DefaultPulsarMessageListenerContainerTests.java @@ -29,6 +29,7 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.Condition; @@ -293,7 +294,7 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> dlqLatch.countDown()); dlqContainerProperties.setSchema(Schema.INT32); dlqContainerProperties.setSubscriptionType(SubscriptionType.Shared); - dlqContainerProperties.setTopics(new String[] { "dpmlct-016-dlq-topic" }); + dlqContainerProperties.setTopics(Set.of("dpmlct-016-dlq-topic")); DefaultPulsarMessageListenerContainer dlqContainer = new DefaultPulsarMessageListenerContainer<>( pulsarConsumerFactory, dlqContainerProperties); dlqContainer.start(); @@ -349,7 +350,7 @@ class DefaultPulsarMessageListenerContainerTests implements PulsarTestContainerS .setMessageListener((PulsarRecordMessageListener) (consumer, msg) -> dlqLatch.countDown()); dlqContainerProperties.setSchema(Schema.INT32); dlqContainerProperties.setSubscriptionType(SubscriptionType.Shared); - dlqContainerProperties.setTopics(new String[] { "dlq-topic" }); + dlqContainerProperties.setTopics(Set.of("dlq-topic")); DefaultPulsarMessageListenerContainer dlqContainer = new DefaultPulsarMessageListenerContainer<>( pulsarConsumerFactory, dlqContainerProperties); dlqContainer.start();