Concurrent Message Listener Container
Factory to create a concurrent message listener container. This is the default factory with concurrency set to 1. Users can configure it to increase concurrency.
This commit is contained in:
@@ -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<Object>> pulsarConsumerFactory) {
|
||||
DefaultPulsarListenerContainerFactory<Object, Object> factory = new DefaultPulsarListenerContainerFactory<>();
|
||||
ConcurrentPulsarListenerContainerFactory<Object> factory = new ConcurrentPulsarListenerContainerFactory<>();
|
||||
|
||||
final PulsarConsumerFactory<Object> pulsarConsumerFactory1 = pulsarConsumerFactory.getIfAvailable();
|
||||
factory.setPulsarConsumerFactory(pulsarConsumerFactory1);
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -77,6 +77,8 @@ public abstract class AbstractPulsarListenerEndpoint<K>
|
||||
|
||||
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<K>
|
||||
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;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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 <C> container implementation type.
|
||||
* @param <T> message type in the listener.
|
||||
* @author Soby Chacko
|
||||
* @author Chris Bono
|
||||
*/
|
||||
public class DefaultPulsarListenerContainerFactory<C, T>
|
||||
extends AbstractPulsarListenerContainerFactory<DefaultPulsarMessageListenerContainer<T>, T> {
|
||||
public class ConcurrentPulsarListenerContainerFactory<T>
|
||||
extends AbstractPulsarListenerContainerFactory<ConcurrentPulsarMessageListenerContainer<T>, 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<T> createContainerInstance(PulsarListenerEndpoint endpoint) {
|
||||
protected ConcurrentPulsarMessageListenerContainer<T> createContainerInstance(PulsarListenerEndpoint endpoint) {
|
||||
|
||||
PulsarContainerProperties properties = new PulsarContainerProperties();
|
||||
Collection<String> topics = endpoint.getTopics();
|
||||
@@ -62,18 +71,25 @@ public class DefaultPulsarListenerContainerFactory<C, T>
|
||||
|
||||
properties.setSchemaType(endpoint.getSchemaType());
|
||||
|
||||
return new DefaultPulsarMessageListenerContainer<T>(getPulsarConsumerFactory(), properties);
|
||||
return new ConcurrentPulsarMessageListenerContainer<T>(getPulsarConsumerFactory(), properties);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void initializeContainer(DefaultPulsarMessageListenerContainer<T> instance,
|
||||
protected void initializeContainer(ConcurrentPulsarMessageListenerContainer<T> 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<T> createContainer(String... topics) {
|
||||
public ConcurrentPulsarMessageListenerContainer<T> createContainer(String... topics) {
|
||||
PulsarListenerEndpoint endpoint = new PulsarListenerEndpointAdapter() {
|
||||
|
||||
@Override
|
||||
@@ -82,7 +98,7 @@ public class DefaultPulsarListenerContainerFactory<C, T>
|
||||
}
|
||||
|
||||
};
|
||||
DefaultPulsarMessageListenerContainer<T> container = createContainerInstance(endpoint);
|
||||
ConcurrentPulsarMessageListenerContainer<T> container = createContainerInstance(endpoint);
|
||||
initializeContainer(container, endpoint);
|
||||
// customizeContainer(container);
|
||||
return container;
|
||||
@@ -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<V> 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) {
|
||||
|
||||
@@ -58,4 +58,7 @@ public interface PulsarListenerEndpoint {
|
||||
|
||||
Properties getConsumerProperties();
|
||||
|
||||
@Nullable
|
||||
Integer getConcurrency();
|
||||
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -42,7 +42,7 @@ public class ConcurrentPulsarMessageListenerContainer<T> extends AbstractPulsarM
|
||||
|
||||
private final List<AsyncTaskExecutor> executors = new ArrayList<>();
|
||||
|
||||
protected ConcurrentPulsarMessageListenerContainer(PulsarConsumerFactory<? super T> pulsarConsumerFactory,
|
||||
public ConcurrentPulsarMessageListenerContainer(PulsarConsumerFactory<? super T> pulsarConsumerFactory,
|
||||
PulsarContainerProperties pulsarContainerProperties) {
|
||||
super(pulsarConsumerFactory, pulsarContainerProperties);
|
||||
}
|
||||
@@ -98,13 +98,8 @@ public class ConcurrentPulsarMessageListenerContainer<T> 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);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -186,6 +186,7 @@ public class DefaultPulsarMessageListenerContainer<T> 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.");
|
||||
|
||||
@@ -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 {
|
||||
@@ -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<String> pulsarConsumerFactory = mock(PulsarConsumerFactory.class);
|
||||
Consumer<String> 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<String> 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();
|
||||
}
|
||||
|
||||
}
|
||||
@@ -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<Object> pulsarConsumerFactory) {
|
||||
final DefaultPulsarListenerContainerFactory<?, ?> pulsarListenerContainerFactory = new DefaultPulsarListenerContainerFactory<>();
|
||||
final ConcurrentPulsarListenerContainerFactory<?> pulsarListenerContainerFactory = new ConcurrentPulsarListenerContainerFactory<>();
|
||||
pulsarListenerContainerFactory.setPulsarConsumerFactory(pulsarConsumerFactory);
|
||||
return pulsarListenerContainerFactory;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user