From 6dcc813576a49aec3196e4ad5f8dca23381250b4 Mon Sep 17 00:00:00 2001 From: Daniel Szabo Date: Thu, 1 May 2025 21:29:59 +0200 Subject: [PATCH] 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 --- ...bstractPulsarListenerContainerFactory.java | 1 - ...currentPulsarListenerContainerFactory.java | 2 ++ ...ntPulsarListenerContainerFactoryTests.java | 23 +++++++++++++++++++ 3 files changed, 25 insertions(+), 1 deletion(-) 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 95a23dd8..134da31c 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 @@ -115,7 +115,6 @@ public abstract class AbstractPulsarListenerContainerFactory extends AbstractPulsarListenerContainerFactory, T> { @@ -80,6 +81,7 @@ public class ConcurrentPulsarListenerContainerFactory 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()); diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactoryTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactoryTests.java index b5dfcd45..b0956279 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactoryTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/config/ConcurrentPulsarListenerContainerFactoryTests.java @@ -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( + 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()); + } + + } + }