From 29cfd301c6e0753e8626647bdde1c9e16e4fd8b0 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 22 Mar 2023 18:46:13 -0400 Subject: [PATCH] Better container properties sync * If consumer properties are directly overridden, we need to sync that with container properties. This commit tries to centralize this synching. --- .../binder/PulsarMessageChannelBinder.java | 5 +--- .../PulsarExtendedBindingPropertiesTests.java | 23 +++++++++++++++++++ ...bstractPulsarListenerContainerFactory.java | 2 ++ .../listener/PulsarContainerProperties.java | 15 ++++++++++++ 4 files changed, 41 insertions(+), 4 deletions(-) diff --git a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java index 47cd8f44..2a353f9c 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarMessageChannelBinder.java @@ -20,7 +20,6 @@ import java.util.Optional; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; -import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.common.schema.SchemaType; import org.springframework.cloud.stream.binder.AbstractMessageChannelBinder; @@ -163,9 +162,6 @@ public class PulsarMessageChannelBinder extends } var subscriptionName = PulsarBinderUtils.subscriptionName(properties.getExtension(), destination); containerProperties.setSubscriptionName(subscriptionName); - if (properties.getExtension().getSubscriptionType() != SubscriptionType.Exclusive) { - containerProperties.setSubscriptionType(properties.getExtension().getSubscriptionType()); - } var baseConsumerProps = new ConsumerConfigProperties().buildProperties(); var binderConsumerProps = this.binderConfigProps.getConsumer().buildProperties(); @@ -173,6 +169,7 @@ public class PulsarMessageChannelBinder extends var mergedConsumerProps = PulsarBinderUtils.mergePropertiesWithPrecedence(baseConsumerProps, binderConsumerProps, bindingConsumerProps); containerProperties.getPulsarConsumerProperties().putAll(mergedConsumerProps); + containerProperties.updateContainerProperties(); var container = new DefaultPulsarMessageListenerContainer<>(this.pulsarConsumerFactory, containerProperties); messageDrivenChannelAdapter.setMessageListenerContainer(container); diff --git a/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarExtendedBindingPropertiesTests.java b/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarExtendedBindingPropertiesTests.java index 20e7def6..3a2e56a5 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarExtendedBindingPropertiesTests.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarExtendedBindingPropertiesTests.java @@ -24,6 +24,7 @@ import java.util.Map; import org.apache.pulsar.client.api.ProducerAccessMode; import org.apache.pulsar.client.api.SubscriptionMode; +import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.impl.conf.ConfigurationDataUtils; import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; import org.apache.pulsar.client.impl.conf.ProducerConfigurationData; @@ -34,6 +35,7 @@ import org.springframework.boot.context.properties.bind.Bindable; import org.springframework.boot.context.properties.bind.Binder; import org.springframework.boot.context.properties.source.ConfigurationPropertySource; import org.springframework.boot.context.properties.source.MapConfigurationPropertySource; +import org.springframework.pulsar.listener.PulsarContainerProperties; import org.springframework.pulsar.spring.cloud.stream.binder.properties.PulsarExtendedBindingProperties; /** @@ -111,4 +113,25 @@ public class PulsarExtendedBindingPropertiesTests { // @formatter:on } + @Test + void extendedBindingsArePropagatedToContainerProperties() { + Map props = new HashMap<>(); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.consumer.subscription-name", "my-foo-sbscription"); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.consumer.subscription-type", "Shared"); + + bind(props); + + var bindingConsumerProps = properties.getExtendedConsumerProperties("my-foo").buildProperties(); + PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); + pulsarContainerProperties.getPulsarConsumerProperties().putAll(bindingConsumerProps); + + assertThat(pulsarContainerProperties.getSubscriptionName()).isNull(); + assertThat(pulsarContainerProperties.getSubscriptionType()).isEqualTo(SubscriptionType.Exclusive); + + pulsarContainerProperties.updateContainerProperties(); + + assertThat(pulsarContainerProperties.getSubscriptionName()).isEqualTo("my-foo-sbscription"); + assertThat(pulsarContainerProperties.getSubscriptionType()).isEqualTo(SubscriptionType.Shared); + } + } 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 b1ee5420..846cec3a 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 @@ -178,6 +178,8 @@ public abstract class AbstractPulsarListenerContainerFactory