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.
This commit is contained in:
@@ -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);
|
||||
|
||||
@@ -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<String, String> 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);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -178,6 +178,8 @@ public abstract class AbstractPulsarListenerContainerFactory<C extends AbstractP
|
||||
.acceptIfNotNull(this.applicationEventPublisher, instance::setApplicationEventPublisher)
|
||||
.acceptIfNotNull(endpoint.getConsumerProperties(),
|
||||
instance.getContainerProperties()::setPulsarConsumerProperties);
|
||||
// Update container properties if there are relevant direct consumer properties
|
||||
instanceProperties.updateContainerProperties();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -42,6 +42,10 @@ public class PulsarContainerProperties {
|
||||
|
||||
private static final Duration DEFAULT_CONSUMER_START_TIMEOUT = Duration.ofSeconds(30);
|
||||
|
||||
private static final String SUBSCRIPTION_NAME = "subscriptionName";
|
||||
|
||||
private static final String SUBSCRIPTION_TYPE = "subscriptionType";
|
||||
|
||||
private Duration consumerStartTimeout = DEFAULT_CONSUMER_START_TIMEOUT;
|
||||
|
||||
private String[] topics;
|
||||
@@ -246,4 +250,15 @@ public class PulsarContainerProperties {
|
||||
this.pulsarConsumerProperties = pulsarConsumerProperties;
|
||||
}
|
||||
|
||||
public void updateContainerProperties() {
|
||||
if (!this.pulsarConsumerProperties.isEmpty()) {
|
||||
if (this.pulsarConsumerProperties.containsKey(SUBSCRIPTION_NAME)) {
|
||||
this.subscriptionName = (String) this.pulsarConsumerProperties.get(SUBSCRIPTION_NAME);
|
||||
}
|
||||
if (this.pulsarConsumerProperties.containsKey(SUBSCRIPTION_TYPE)) {
|
||||
this.subscriptionType = (SubscriptionType) this.pulsarConsumerProperties.get(SUBSCRIPTION_TYPE);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user