Pulsar binder partition count defaults (#339)
* Pulsar binder partition count defaults In Pulsar, partitioned topics are optional. If we set the default partition count to be 1, then it unnecessarily creates a partitoned topic. Therefore, we need to set the default partition count as 0 in the binder and extended binding properties. * PR review
This commit is contained in:
@@ -18,6 +18,7 @@ package org.springframework.pulsar.spring.cloud.stream.binder.properties;
|
||||
|
||||
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||
import org.springframework.boot.context.properties.NestedConfigurationProperty;
|
||||
import org.springframework.lang.Nullable;
|
||||
import org.springframework.pulsar.autoconfigure.ConsumerConfigProperties;
|
||||
import org.springframework.pulsar.autoconfigure.ProducerConfigProperties;
|
||||
|
||||
@@ -36,7 +37,8 @@ public class PulsarBinderConfigurationProperties {
|
||||
@NestedConfigurationProperty
|
||||
private final ProducerConfigProperties producer = new ProducerConfigProperties();
|
||||
|
||||
private int partitionCount = 1;
|
||||
@Nullable
|
||||
private Integer partitionCount;
|
||||
|
||||
public ConsumerConfigProperties getConsumer() {
|
||||
return this.consumer;
|
||||
@@ -46,11 +48,12 @@ public class PulsarBinderConfigurationProperties {
|
||||
return this.producer;
|
||||
}
|
||||
|
||||
public int partitionCount() {
|
||||
@Nullable
|
||||
public Integer getPartitionCount() {
|
||||
return this.partitionCount;
|
||||
}
|
||||
|
||||
public void setPartitionCount(int partitionCount) {
|
||||
public void setPartitionCount(Integer partitionCount) {
|
||||
this.partitionCount = partitionCount;
|
||||
}
|
||||
|
||||
|
||||
@@ -41,7 +41,8 @@ public class PulsarConsumerProperties extends ConsumerConfigProperties {
|
||||
@Nullable
|
||||
private Class<?> messageValueType;
|
||||
|
||||
private Integer partitionCount = 1;
|
||||
@Nullable
|
||||
private Integer partitionCount;
|
||||
|
||||
@Nullable
|
||||
public SchemaType getSchemaType() {
|
||||
@@ -79,6 +80,7 @@ public class PulsarConsumerProperties extends ConsumerConfigProperties {
|
||||
this.messageValueType = messageValueType;
|
||||
}
|
||||
|
||||
@Nullable
|
||||
public Integer getPartitionCount() {
|
||||
return this.partitionCount;
|
||||
}
|
||||
|
||||
@@ -41,6 +41,9 @@ public class PulsarProducerProperties extends ProducerConfigProperties {
|
||||
@Nullable
|
||||
private Class<?> messageValueType;
|
||||
|
||||
@Nullable
|
||||
private Integer partitionCount;
|
||||
|
||||
@Nullable
|
||||
public SchemaType getSchemaType() {
|
||||
return this.schemaType;
|
||||
@@ -77,4 +80,13 @@ public class PulsarProducerProperties extends ProducerConfigProperties {
|
||||
this.messageValueType = messageValueType;
|
||||
}
|
||||
|
||||
@Nullable
|
||||
public Integer getPartitionCount() {
|
||||
return this.partitionCount;
|
||||
}
|
||||
|
||||
public void setPartitionCount(Integer partitionCount) {
|
||||
this.partitionCount = partitionCount;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -50,28 +50,28 @@ public class PulsarTopicProvisioner implements
|
||||
public ProducerDestination provisionProducerDestination(String name,
|
||||
ExtendedProducerProperties<PulsarProducerProperties> pulsarProducerProperties)
|
||||
throws ProvisioningException {
|
||||
|
||||
int partitionCount = this.pulsarBinderConfigurationProperties.partitionCount();
|
||||
int partitionCountOnBinding = pulsarProducerProperties.getPartitionCount();
|
||||
if (partitionCountOnBinding > 1) {
|
||||
partitionCount = partitionCountOnBinding;
|
||||
}
|
||||
PulsarTopic pulsarTopic = PulsarTopic.builder(name).numberOfPartitions(partitionCount).build();
|
||||
Integer partitionCountFromBinding = pulsarProducerProperties.getExtension().getPartitionCount();
|
||||
var partitionCount = getPartitionCount(partitionCountFromBinding);
|
||||
var pulsarTopic = PulsarTopic.builder(name).numberOfPartitions(partitionCount).build();
|
||||
this.pulsarAdministration.createOrModifyTopics(pulsarTopic);
|
||||
return new PulsarDestination(pulsarTopic.topicName(), pulsarTopic.numberOfPartitions());
|
||||
}
|
||||
|
||||
private int getPartitionCount(Integer partitionCountConfig) {
|
||||
var partitionCount = this.pulsarBinderConfigurationProperties.getPartitionCount();
|
||||
if (partitionCountConfig != null && partitionCountConfig > 0) {
|
||||
partitionCount = partitionCountConfig;
|
||||
}
|
||||
return partitionCount == null ? 0 : partitionCount;
|
||||
}
|
||||
|
||||
@Override
|
||||
public ConsumerDestination provisionConsumerDestination(String name, String group,
|
||||
ExtendedConsumerProperties<PulsarConsumerProperties> pulsarConsumerProperties)
|
||||
throws ProvisioningException {
|
||||
int partitionCount = this.pulsarBinderConfigurationProperties.partitionCount();
|
||||
|
||||
int partitionCountOnBinding = pulsarConsumerProperties.getExtension().getPartitionCount();
|
||||
if (partitionCountOnBinding > 1) {
|
||||
partitionCount = partitionCountOnBinding;
|
||||
}
|
||||
PulsarTopic pulsarTopic = PulsarTopic.builder(name).numberOfPartitions(partitionCount).build();
|
||||
var partitionCountFromBinding = pulsarConsumerProperties.getExtension().getPartitionCount();
|
||||
var partitionCount = getPartitionCount(partitionCountFromBinding);
|
||||
var pulsarTopic = PulsarTopic.builder(name).numberOfPartitions(partitionCount).build();
|
||||
this.pulsarAdministration.createOrModifyTopics(pulsarTopic);
|
||||
return new PulsarDestination(pulsarTopic.topicName(), pulsarTopic.numberOfPartitions());
|
||||
}
|
||||
|
||||
@@ -52,9 +52,9 @@ public class PulsarBinderConfigurationPropertiesTests {
|
||||
|
||||
@Test
|
||||
void partitionCountProperty() {
|
||||
assertThat(properties.partitionCount()).isEqualTo(1);
|
||||
assertThat(properties.getPartitionCount()).isNull();
|
||||
bind(Map.of("spring.cloud.stream.pulsar.binder.partition-count", "5150"));
|
||||
assertThat(properties.partitionCount()).isEqualTo(5150);
|
||||
assertThat(properties.getPartitionCount()).isEqualTo(5150);
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -50,7 +50,7 @@ public class PulsarTopicProvisionerTests {
|
||||
new PulsarProducerProperties());
|
||||
ProducerDestination producerDestination = pulsarTopicProvisioner.provisionProducerDestination("foo",
|
||||
properties);
|
||||
verifyAndAssert(pulsarAdministration, producerDestination.getName(), "foo", 1);
|
||||
verifyAndAssert(pulsarAdministration, producerDestination.getName(), "foo", 0);
|
||||
}
|
||||
|
||||
private static void verifyAndAssert(PulsarAdministration pulsarAdministration, String actualProducerDestination,
|
||||
@@ -73,7 +73,7 @@ public class PulsarTopicProvisionerTests {
|
||||
new PulsarConsumerProperties());
|
||||
ConsumerDestination consumerDestination = pulsarTopicProvisioner.provisionConsumerDestination("bar", "",
|
||||
properties);
|
||||
verifyAndAssert(pulsarAdministration, consumerDestination.getName(), "bar", 1);
|
||||
verifyAndAssert(pulsarAdministration, consumerDestination.getName(), "bar", 0);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -98,7 +98,7 @@ public class PulsarTopicProvisionerTests {
|
||||
pulsarBinderConfigurationProperties);
|
||||
ExtendedProducerProperties<PulsarProducerProperties> properties = new ExtendedProducerProperties<>(
|
||||
new PulsarProducerProperties());
|
||||
properties.setPartitionCount(4);
|
||||
properties.getExtension().setPartitionCount(4);
|
||||
ProducerDestination producerDestination = pulsarTopicProvisioner.provisionProducerDestination("foo",
|
||||
properties);
|
||||
verifyAndAssert(pulsarAdministration, producerDestination.getName(), "foo", 4);
|
||||
|
||||
Reference in New Issue
Block a user