From 1843fedc26098ab9509f409bf0afd417ec7e3198 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 24 Aug 2022 19:13:21 -0400 Subject: [PATCH] Concurrent message listener container improvements --- .../PulsarAnnotationDrivenConfiguration.java | 1 + ...bstractPulsarListenerContainerFactory.java | 13 ++++++ .../AbstractPulsarListenerEndpoint.java | 1 - ...currentPulsarMessageListenerContainer.java | 6 +++ ...ntPulsarMessageListenerContainerTests.java | 41 +++++++++++++++---- 5 files changed, 54 insertions(+), 8 deletions(-) diff --git a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAnnotationDrivenConfiguration.java b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAnnotationDrivenConfiguration.java index 31d38384..8792a504 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAnnotationDrivenConfiguration.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/main/java/org/springframework/pulsar/autoconfigure/PulsarAnnotationDrivenConfiguration.java @@ -54,6 +54,7 @@ public class PulsarAnnotationDrivenConfiguration { factory.setPulsarConsumerFactory(pulsarConsumerFactory1); final PulsarContainerProperties containerProperties = factory.getContainerProperties(); + containerProperties.setSubscriptionType(this.pulsarProperties.getConsumer().getSubscriptionType()); PropertyMapper map = PropertyMapper.get().alwaysApplyingWhenNonNull(); PulsarProperties.Listener properties = this.pulsarProperties.getListener(); diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarListenerContainerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarListenerContainerFactory.java index 4e71bd3f..e9c78be0 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarListenerContainerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/AbstractPulsarListenerContainerFactory.java @@ -142,10 +142,23 @@ public abstract class AbstractPulsarListenerContainerFactory /** * Set the concurrency for this endpoint's container. * @param concurrency the concurrency. - * @since 2.2 */ public void setConcurrency(Integer concurrency) { this.concurrency = concurrency; 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 59d66f9f..16244b17 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 @@ -19,6 +19,8 @@ package org.springframework.pulsar.listener; import java.util.ArrayList; import java.util.List; +import org.apache.pulsar.client.api.SubscriptionType; + import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationEventPublisher; import org.springframework.core.task.AsyncTaskExecutor; @@ -67,6 +69,10 @@ public class ConcurrentPulsarMessageListenerContainer extends AbstractPulsarM if (!isRunning()) { PulsarContainerProperties containerProperties = getContainerProperties(); + if (containerProperties.getSubscriptionType() == SubscriptionType.Exclusive && this.concurrency > 1) { + throw new IllegalStateException("concurrency > 1 is not allowed on Exclusive subscription type"); + } + setRunning(true); for (int i = 0; i < this.concurrency; i++) { 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 e4754471..87753605 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 @@ -16,6 +16,7 @@ package org.springframework.pulsar.listener; +import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.times; @@ -28,6 +29,7 @@ import org.apache.pulsar.client.api.BatchReceivePolicy; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Messages; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.SubscriptionType; import org.junit.jupiter.api.Test; import org.springframework.pulsar.core.AbstractContainerBaseTests; @@ -40,7 +42,36 @@ public class ConcurrentPulsarMessageListenerContainerTests extends AbstractConta @Test @SuppressWarnings("unchecked") - void basicConcurrency() throws Exception { + void basicConcurrencyTesting() throws Exception { + PulsarConsumerFactory pulsarConsumerFactory = mock(PulsarConsumerFactory.class); + Consumer consumer = mock(Consumer.class); + + when(pulsarConsumerFactory.createConsumer(any(Schema.class), any(BatchReceivePolicy.class), any(Map.class))) + .thenReturn(consumer); + + when(consumer.batchReceive()).thenReturn(mock(Messages.class)); + + PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); + pulsarContainerProperties.setSchema(Schema.STRING); + pulsarContainerProperties.setSubscriptionType(SubscriptionType.Failover); + pulsarContainerProperties.setMessageListener((PulsarRecordMessageListener) (cons, msg) -> { + }); + + ConcurrentPulsarMessageListenerContainer container = new ConcurrentPulsarMessageListenerContainer<>( + pulsarConsumerFactory, pulsarContainerProperties); + + container.setConcurrency(3); + + container.start(); + + verify(pulsarConsumerFactory, times(3)).createConsumer(any(Schema.class), any(BatchReceivePolicy.class), + any(Map.class)); + verify(consumer, times(3)).batchReceive(); + } + + @Test + @SuppressWarnings("unchecked") + void exclusiveSubscriptionMustUseSingleThread() throws Exception { PulsarConsumerFactory pulsarConsumerFactory = mock(PulsarConsumerFactory.class); Consumer consumer = mock(Consumer.class); @@ -59,12 +90,8 @@ public class ConcurrentPulsarMessageListenerContainerTests extends AbstractConta container.setConcurrency(3); - container.start(); - - verify(pulsarConsumerFactory, times(3)).createConsumer(any(Schema.class), any(BatchReceivePolicy.class), - any(Map.class)); - - verify(consumer, times(3)).batchReceive(); + assertThatThrownBy(container::start).isInstanceOf(IllegalStateException.class) + .hasMessage("concurrency > 1 is not allowed on Exclusive subscription type"); } }