Respect executor set on container props

Makes sure that the ConcurrentPulsarListenerContainerFactory copies the
task executor from the factory container properties to the container
instance properties.

Signed-off-by: Daniel Szabo <cheery.gate6737@milkmail.eu>
This commit is contained in:
Daniel Szabo
2025-05-01 21:29:59 +02:00
committed by Chris Bono
parent bd9aea5858
commit 6dcc813576
3 changed files with 25 additions and 1 deletions

View File

@@ -115,7 +115,6 @@ public abstract class AbstractPulsarListenerContainerFactory<C extends AbstractP
endpoint.setupListenerContainer(instance, this.messageConverter);
initializeContainer(instance, endpoint);
// customizeContainer(instance);
return instance;
}

View File

@@ -37,6 +37,7 @@ import org.springframework.util.StringUtils;
* @author Chris Bono
* @author Alexander Preuß
* @author Vedran Pavic
* @author Daniel Szabo
*/
public class ConcurrentPulsarListenerContainerFactory<T>
extends AbstractPulsarListenerContainerFactory<ConcurrentPulsarMessageListenerContainer<T>, T> {
@@ -80,6 +81,7 @@ public class ConcurrentPulsarListenerContainerFactory<T>
var containerProps = new PulsarContainerProperties();
// Map factory props (defaults) to the container props
containerProps.setConsumerTaskExecutor(factoryProps.getConsumerTaskExecutor());
containerProps.setSchemaResolver(factoryProps.getSchemaResolver());
containerProps.setTopicResolver(factoryProps.getTopicResolver());
containerProps.setSubscriptionType(factoryProps.getSubscriptionType());

View File

@@ -26,6 +26,7 @@ import org.apache.pulsar.client.api.SubscriptionType;
import org.junit.jupiter.api.Nested;
import org.junit.jupiter.api.Test;
import org.springframework.core.task.SimpleAsyncTaskExecutor;
import org.springframework.pulsar.core.PulsarConsumerFactory;
import org.springframework.pulsar.listener.PulsarContainerProperties;
@@ -142,4 +143,26 @@ class ConcurrentPulsarListenerContainerFactoryTests {
}
@Nested
class ConsumerTaskExecutor {
@Test
@SuppressWarnings("unchecked")
void factoryValueCopiedWhenCreatingContainer() {
final var factoryProps = new PulsarContainerProperties();
factoryProps.setConsumerTaskExecutor(new SimpleAsyncTaskExecutor());
final var containerFactory = new ConcurrentPulsarListenerContainerFactory<String>(
mock(PulsarConsumerFactory.class), factoryProps);
final var endpoint = mock(PulsarListenerEndpoint.class);
// Mockito by default returns 0 for Integer
when(endpoint.getConcurrency()).thenReturn(null);
final var createdContainer = containerFactory.createRegisteredContainer(endpoint);
final var containerProperties = createdContainer.getContainerProperties();
assertThat(containerProperties.getConsumerTaskExecutor()).isEqualTo(factoryProps.getConsumerTaskExecutor());
}
}
}