From 9d48b8df1b3987709fbc2cf230299f0ff4f12ed4 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 16 Sep 2022 18:22:52 -0400 Subject: [PATCH] Batch receive properties not copied properly * When Boot provides properties for batch receive, they are not copied downstream through the ConcurrentPulsarListenerContainerFactory. Fixing this issue. --- ...bstractPulsarListenerContainerFactory.java | 20 ++++++++++++++ ...ntPulsarMessageListenerContainerTests.java | 27 +++++++++++++++++++ 2 files changed, 47 insertions(+) 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 ad086fb7..31a4a966 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 @@ -16,10 +16,14 @@ package org.springframework.pulsar.config; +import java.util.Arrays; + import org.apache.commons.logging.LogFactory; import org.apache.pulsar.client.api.Schema; +import org.springframework.beans.BeanWrapper; import org.springframework.beans.BeansException; +import org.springframework.beans.PropertyAccessorFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; @@ -171,6 +175,9 @@ public abstract class AbstractPulsarListenerContainerFactory requestProperties) { + BeanWrapper wrappedSource = PropertyAccessorFactory.forBeanPropertyAccess(source); + BeanWrapper wrappedTarget = PropertyAccessorFactory.forBeanPropertyAccess(target); + + requestProperties.forEach(p -> wrappedTarget.setPropertyValue(p, wrappedSource.getPropertyValue(p))); + } + } 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 08968bbd..c5915d4c 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 @@ -38,6 +38,8 @@ import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.impl.MultiplierRedeliveryBackoff; import org.junit.jupiter.api.Test; +import org.springframework.pulsar.config.ConcurrentPulsarListenerContainerFactory; +import org.springframework.pulsar.config.PulsarListenerEndpoint; import org.springframework.pulsar.core.PulsarConsumerFactory; import org.springframework.util.backoff.BackOff; @@ -47,6 +49,31 @@ import org.springframework.util.backoff.BackOff; */ public class ConcurrentPulsarMessageListenerContainerTests { + @Test + @SuppressWarnings("unchecked") + void createConcurrentContainerFromFactoryAndVerifyBatchReceivePolicy() { + ConcurrentPulsarListenerContainerFactory factory = new ConcurrentPulsarListenerContainerFactory<>(); + final PulsarConsumerFactory pulsarConsumerFactory = mock(PulsarConsumerFactory.class); + factory.setPulsarConsumerFactory(pulsarConsumerFactory); + + PulsarContainerProperties containerProperties = factory.getContainerProperties(); + containerProperties.setBatchTimeout(60_000); + containerProperties.setMaxNumMessages(120); + containerProperties.setMaxNumBytes(32000); + + factory.setConcurrency(1); + + PulsarListenerEndpoint pulsarListenerEndpoint = mock(PulsarListenerEndpoint.class); + when(pulsarListenerEndpoint.getConcurrency()).thenReturn(1); + + final ConcurrentPulsarMessageListenerContainer concurrentContainer = factory + .createListenerContainer(pulsarListenerEndpoint); + final PulsarContainerProperties pulsarContainerProperties = concurrentContainer.getContainerProperties(); + assertThat(pulsarContainerProperties.getBatchTimeout()).isEqualTo(60_000); + assertThat(pulsarContainerProperties.getMaxNumMessages()).isEqualTo(120); + assertThat(pulsarContainerProperties.getMaxNumBytes()).isEqualTo(32_000); + } + @Test void deadLetterPolicyAppliedOnChildContainer() throws Exception { PulsarListenerMockComponents env = setupListenerMockComponents(SubscriptionType.Shared);