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 83599cbf..31d38384 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 @@ -23,7 +23,7 @@ import org.springframework.boot.context.properties.PropertyMapper; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.pulsar.annotation.EnablePulsar; -import org.springframework.pulsar.config.DefaultPulsarListenerContainerFactory; +import org.springframework.pulsar.config.ConcurrentPulsarListenerContainerFactory; import org.springframework.pulsar.config.PulsarListenerBeanNames; import org.springframework.pulsar.core.PulsarConsumerFactory; import org.springframework.pulsar.listener.PulsarContainerProperties; @@ -46,9 +46,9 @@ public class PulsarAnnotationDrivenConfiguration { @Bean @ConditionalOnMissingBean(name = "pulsarListenerContainerFactory") - DefaultPulsarListenerContainerFactory pulsarListenerContainerFactory( + ConcurrentPulsarListenerContainerFactory pulsarListenerContainerFactory( ObjectProvider> pulsarConsumerFactory) { - DefaultPulsarListenerContainerFactory factory = new DefaultPulsarListenerContainerFactory<>(); + ConcurrentPulsarListenerContainerFactory factory = new ConcurrentPulsarListenerContainerFactory<>(); final PulsarConsumerFactory pulsarConsumerFactory1 = pulsarConsumerFactory.getIfAvailable(); factory.setPulsarConsumerFactory(pulsarConsumerFactory1); 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 343ebd7a..2cc7098e 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 @@ -36,7 +36,7 @@ import org.springframework.core.annotation.Order; import org.springframework.pulsar.annotation.EnablePulsar; import org.springframework.pulsar.annotation.PulsarBootstrapConfiguration; import org.springframework.pulsar.annotation.PulsarListenerAnnotationBeanPostProcessor; -import org.springframework.pulsar.config.DefaultPulsarListenerContainerFactory; +import org.springframework.pulsar.config.ConcurrentPulsarListenerContainerFactory; import org.springframework.pulsar.config.PulsarClientConfiguration; import org.springframework.pulsar.config.PulsarClientFactoryBean; import org.springframework.pulsar.config.PulsarListenerContainerFactory; @@ -84,12 +84,13 @@ class PulsarAutoConfigurationTests { @Test void defaultBeansAreAutoConfigured() { - this.contextRunner.run((context) -> assertThat(context).hasNotFailed() - .hasSingleBean(PulsarClientConfiguration.class).hasSingleBean(PulsarClientFactoryBean.class) - .hasSingleBean(PulsarProducerFactory.class).hasSingleBean(PulsarTemplate.class) - .hasSingleBean(PulsarConsumerFactory.class).hasSingleBean(DefaultPulsarListenerContainerFactory.class) - .hasSingleBean(PulsarListenerAnnotationBeanPostProcessor.class) - .hasSingleBean(PulsarListenerEndpointRegistry.class)); + this.contextRunner + .run((context) -> assertThat(context).hasNotFailed().hasSingleBean(PulsarClientConfiguration.class) + .hasSingleBean(PulsarClientFactoryBean.class).hasSingleBean(PulsarProducerFactory.class) + .hasSingleBean(PulsarTemplate.class).hasSingleBean(PulsarConsumerFactory.class) + .hasSingleBean(ConcurrentPulsarListenerContainerFactory.class) + .hasSingleBean(PulsarListenerAnnotationBeanPostProcessor.class) + .hasSingleBean(PulsarListenerEndpointRegistry.class)); } @Test 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 deb40c28..a7f89bc9 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 @@ -77,6 +77,8 @@ public abstract class AbstractPulsarListenerEndpoint private Boolean batchListener; + private Integer concurrency; + @Override public void setBeanFactory(BeanFactory beanFactory) throws BeansException { this.beanFactory = beanFactory; @@ -214,4 +216,19 @@ public abstract class AbstractPulsarListenerEndpoint this.schemaType = schemaType; } + @Override + @Nullable + public Integer getConcurrency() { + return this.concurrency; + } + + /** + * 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/config/DefaultPulsarListenerContainerFactory.java b/spring-pulsar/src/main/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactory.java similarity index 66% rename from spring-pulsar/src/main/java/org/springframework/pulsar/config/DefaultPulsarListenerContainerFactory.java rename to spring-pulsar/src/main/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactory.java index fcef2bbf..8deb09be 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/config/DefaultPulsarListenerContainerFactory.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactory.java @@ -21,23 +21,32 @@ import java.util.Collection; import org.apache.pulsar.client.api.SubscriptionType; -import org.springframework.pulsar.listener.DefaultPulsarMessageListenerContainer; +import org.springframework.pulsar.listener.ConcurrentPulsarMessageListenerContainer; import org.springframework.pulsar.listener.PulsarContainerProperties; import org.springframework.util.StringUtils; /** * Concrete implementation for {@link PulsarListenerContainerFactory}. * - * @param container implementation type. * @param message type in the listener. * @author Soby Chacko * @author Chris Bono */ -public class DefaultPulsarListenerContainerFactory - extends AbstractPulsarListenerContainerFactory, T> { +public class ConcurrentPulsarListenerContainerFactory + extends AbstractPulsarListenerContainerFactory, T> { + + private Integer concurrency; + + /** + * Specify the container concurrency. + * @param concurrency the number of consumers to create. + */ + public void setConcurrency(Integer concurrency) { + this.concurrency = concurrency; + } @Override - protected DefaultPulsarMessageListenerContainer createContainerInstance(PulsarListenerEndpoint endpoint) { + protected ConcurrentPulsarMessageListenerContainer createContainerInstance(PulsarListenerEndpoint endpoint) { PulsarContainerProperties properties = new PulsarContainerProperties(); Collection topics = endpoint.getTopics(); @@ -62,18 +71,25 @@ public class DefaultPulsarListenerContainerFactory properties.setSchemaType(endpoint.getSchemaType()); - return new DefaultPulsarMessageListenerContainer(getPulsarConsumerFactory(), properties); + return new ConcurrentPulsarMessageListenerContainer(getPulsarConsumerFactory(), properties); } @Override - protected void initializeContainer(DefaultPulsarMessageListenerContainer instance, + protected void initializeContainer(ConcurrentPulsarMessageListenerContainer instance, PulsarListenerEndpoint endpoint) { super.initializeContainer(instance, endpoint); + Integer conc = endpoint.getConcurrency(); + if (conc != null) { + instance.setConcurrency(conc); + } + else if (this.concurrency != null) { + instance.setConcurrency(this.concurrency); + } } @Override - public DefaultPulsarMessageListenerContainer createContainer(String... topics) { + public ConcurrentPulsarMessageListenerContainer createContainer(String... topics) { PulsarListenerEndpoint endpoint = new PulsarListenerEndpointAdapter() { @Override @@ -82,7 +98,7 @@ public class DefaultPulsarListenerContainerFactory } }; - DefaultPulsarMessageListenerContainer container = createContainerInstance(endpoint); + ConcurrentPulsarMessageListenerContainer container = createContainerInstance(endpoint); initializeContainer(container, endpoint); // customizeContainer(container); return container; 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 55e6249f..b18726d3 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 @@ -38,7 +38,7 @@ import org.springframework.messaging.converter.SmartMessageConverter; import org.springframework.messaging.handler.annotation.support.MessageHandlerMethodFactory; import org.springframework.messaging.handler.invocation.InvocableHandlerMethod; import org.springframework.pulsar.listener.Acknowledgement; -import org.springframework.pulsar.listener.DefaultPulsarMessageListenerContainer; +import org.springframework.pulsar.listener.ConcurrentPulsarMessageListenerContainer; import org.springframework.pulsar.listener.PulsarContainerProperties; import org.springframework.pulsar.listener.PulsarMessageListenerContainer; import org.springframework.pulsar.listener.adapter.HandlerAdapter; @@ -122,7 +122,7 @@ public class MethodPulsarListenerEndpoint extends AbstractPulsarListenerEndpo methodParameter = parameter.get(); } - final DefaultPulsarMessageListenerContainer containerInstance = (DefaultPulsarMessageListenerContainer) container; + final ConcurrentPulsarMessageListenerContainer containerInstance = (ConcurrentPulsarMessageListenerContainer) container; final PulsarContainerProperties pulsarContainerProperties = containerInstance.getPulsarContainerProperties(); final SchemaType schemaType = pulsarContainerProperties.getSchemaType(); if (schemaType != SchemaType.NONE) { 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 0bd4e346..3c951591 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 @@ -58,4 +58,7 @@ public interface PulsarListenerEndpoint { Properties getConsumerProperties(); + @Nullable + Integer getConcurrency(); + } 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 e3b59e4d..7f5ec3d9 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 @@ -23,6 +23,7 @@ import java.util.Properties; import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.common.schema.SchemaType; +import org.springframework.lang.Nullable; import org.springframework.pulsar.listener.PulsarMessageListenerContainer; import org.springframework.pulsar.support.MessageConverter; @@ -79,4 +80,10 @@ public class PulsarListenerEndpointAdapter implements PulsarListenerEndpoint { return null; } + @Nullable + @Override + public Integer getConcurrency() { + return null; + } + } 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 a5722973..59d66f9f 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 @@ -42,7 +42,7 @@ public class ConcurrentPulsarMessageListenerContainer extends AbstractPulsarM private final List executors = new ArrayList<>(); - protected ConcurrentPulsarMessageListenerContainer(PulsarConsumerFactory pulsarConsumerFactory, + public ConcurrentPulsarMessageListenerContainer(PulsarConsumerFactory pulsarConsumerFactory, PulsarContainerProperties pulsarContainerProperties) { super(pulsarConsumerFactory, pulsarContainerProperties); } @@ -98,13 +98,8 @@ public class ConcurrentPulsarMessageListenerContainer extends AbstractPulsarM AsyncTaskExecutor exec = container.getContainerProperties().getConsumerTaskExecutor(); if (exec == null) { - if ((this.executors.size() > index)) { - exec = this.executors.get(index); - } - else { - exec = new SimpleAsyncTaskExecutor(beanName + "-C-"); - this.executors.add(exec); - } + exec = new SimpleAsyncTaskExecutor(beanName + "-C-"); + this.executors.add(exec); container.getContainerProperties().setConsumerTaskExecutor(exec); } } 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 7f58355d..f05dce4f 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 @@ -186,6 +186,7 @@ public class DefaultPulsarMessageListenerContainer extends AbstractPulsarMess .timeout(pulsarContainerProperties.getBatchTimeout(), TimeUnit.MILLISECONDS).build(); this.consumer = getPulsarConsumerFactory().createConsumer( (Schema) pulsarContainerProperties.getSchema(), batchReceivePolicy, propertiesToOverride); + Assert.state(this.consumer != null, "Unable to create a consumer"); } catch (PulsarClientException e) { DefaultPulsarMessageListenerContainer.this.logger.error(e, () -> "Pulsar client exceptions."); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarMessageListenerContainerTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java similarity index 99% rename from spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarMessageListenerContainerTests.java rename to spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java index d227e20f..62a0e858 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarMessageListenerContainerTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/ConsumerAcknowledgmentTests.java @@ -54,7 +54,7 @@ import org.springframework.util.Assert; /** * @author Soby Chacko */ -class PulsarMessageListenerContainerTests extends AbstractContainerBaseTests { +class ConsumerAcknowledgmentTests extends AbstractContainerBaseTests { @Test void testRecordAck() throws Exception { 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 new file mode 100644 index 00000000..e4754471 --- /dev/null +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/listener/ConcurrentPulsarMessageListenerContainerTests.java @@ -0,0 +1,70 @@ +/* + * Copyright 2022 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.pulsar.listener; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import java.util.Map; + +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.junit.jupiter.api.Test; + +import org.springframework.pulsar.core.AbstractContainerBaseTests; +import org.springframework.pulsar.core.PulsarConsumerFactory; + +/** + * @author Soby Chacko + */ +public class ConcurrentPulsarMessageListenerContainerTests extends AbstractContainerBaseTests { + + @Test + @SuppressWarnings("unchecked") + void basicConcurrency() 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.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(); + } + +} 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 28d54f87..e75be260 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 @@ -32,7 +32,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.pulsar.annotation.EnablePulsar; import org.springframework.pulsar.annotation.PulsarListener; -import org.springframework.pulsar.config.DefaultPulsarListenerContainerFactory; +import org.springframework.pulsar.config.ConcurrentPulsarListenerContainerFactory; import org.springframework.pulsar.config.PulsarClientConfiguration; import org.springframework.pulsar.config.PulsarClientFactoryBean; import org.springframework.pulsar.config.PulsarListenerContainerFactory; @@ -113,7 +113,7 @@ public class PulsarListenerTests extends AbstractContainerBaseTests { @Bean PulsarListenerContainerFactory pulsarListenerContainerFactory( PulsarConsumerFactory pulsarConsumerFactory) { - final DefaultPulsarListenerContainerFactory pulsarListenerContainerFactory = new DefaultPulsarListenerContainerFactory<>(); + final ConcurrentPulsarListenerContainerFactory pulsarListenerContainerFactory = new ConcurrentPulsarListenerContainerFactory<>(); pulsarListenerContainerFactory.setPulsarConsumerFactory(pulsarConsumerFactory); return pulsarListenerContainerFactory; }