From 5b2d466c7d5ad71307cec7d12a33a5856000ffbb Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 13 Feb 2023 14:40:01 -0500 Subject: [PATCH] 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 --- .../PulsarBinderConfigurationProperties.java | 9 ++++-- .../properties/PulsarConsumerProperties.java | 4 ++- .../properties/PulsarProducerProperties.java | 12 ++++++++ .../provisioning/PulsarTopicProvisioner.java | 28 +++++++++---------- ...sarBinderConfigurationPropertiesTests.java | 4 +-- .../binder/PulsarTopicProvisionerTests.java | 6 ++-- 6 files changed, 40 insertions(+), 23 deletions(-) diff --git a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarBinderConfigurationProperties.java b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarBinderConfigurationProperties.java index 2646b7a3..b814f8e1 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarBinderConfigurationProperties.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarBinderConfigurationProperties.java @@ -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; } diff --git a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarConsumerProperties.java b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarConsumerProperties.java index f53fc553..2a6f509e 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarConsumerProperties.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarConsumerProperties.java @@ -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; } diff --git a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarProducerProperties.java b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarProducerProperties.java index 5d666e61..6e1e2dc8 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarProducerProperties.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/properties/PulsarProducerProperties.java @@ -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; + } + } diff --git a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/provisioning/PulsarTopicProvisioner.java b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/provisioning/PulsarTopicProvisioner.java index 820136c9..27c624c3 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/provisioning/PulsarTopicProvisioner.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/main/java/org/springframework/pulsar/spring/cloud/stream/binder/provisioning/PulsarTopicProvisioner.java @@ -50,28 +50,28 @@ public class PulsarTopicProvisioner implements public ProducerDestination provisionProducerDestination(String name, ExtendedProducerProperties 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) 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()); } diff --git a/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarBinderConfigurationPropertiesTests.java b/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarBinderConfigurationPropertiesTests.java index 28dc582b..a51750d4 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarBinderConfigurationPropertiesTests.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarBinderConfigurationPropertiesTests.java @@ -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 diff --git a/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarTopicProvisionerTests.java b/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarTopicProvisionerTests.java index dc7e7ec3..3ab8b5a9 100644 --- a/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarTopicProvisionerTests.java +++ b/spring-pulsar-spring-cloud-stream-binder/src/test/java/org/springframework/pulsar/spring/cloud/stream/binder/PulsarTopicProvisionerTests.java @@ -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 properties = new ExtendedProducerProperties<>( new PulsarProducerProperties()); - properties.setPartitionCount(4); + properties.getExtension().setPartitionCount(4); ProducerDestination producerDestination = pulsarTopicProvisioner.provisionProducerDestination("foo", properties); verifyAndAssert(pulsarAdministration, producerDestination.getName(), "foo", 4);