From cbc7981b67da634f25f9fe6ede041d35600950b4 Mon Sep 17 00:00:00 2001 From: Chris Bono Date: Fri, 17 Feb 2023 09:07:31 -0600 Subject: [PATCH] Improve topic validation (#351) * Improve topic validation * Check if topic exists already for both partitioned and unpartitioned cases prior to creating. Fixes #266 * Fix test failure --- .../pulsar/core/PulsarAdministration.java | 60 ++++++++++++------- .../core/PulsarAdministrationTests.java | 45 +++++++++++--- 2 files changed, 74 insertions(+), 31 deletions(-) diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarAdministration.java b/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarAdministration.java index a9b0e5ec..777096ce 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarAdministration.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/core/PulsarAdministration.java @@ -22,6 +22,7 @@ import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.Set; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; @@ -37,6 +38,7 @@ import org.springframework.beans.factory.SmartInitializingSingleton; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; import org.springframework.core.log.LogAccessor; +import org.springframework.lang.Nullable; import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; @@ -55,6 +57,7 @@ public class PulsarAdministration private final PulsarAdminBuilder adminBuilder; + @Nullable private ApplicationContext applicationContext; /** @@ -87,7 +90,7 @@ public class PulsarAdministration } private void loadConf(PulsarAdminBuilder builder, Map adminConfig) { - Map conf = new HashMap<>(adminConfig); + var conf = new HashMap<>(adminConfig); // Workaround the fact that the PulsarAdminImpl does not attempt to construct the // timeout settings from the config props @@ -107,8 +110,8 @@ public class PulsarAdministration // Workaround the fact that the PulsarAdminImpl does not attempt to construct the // authentication from the config props - String authPluginClassName = (String) conf.get("authPluginClassName"); - String authParams = (String) conf.get("authParams"); + var authPluginClassName = (String) conf.get("authPluginClassName"); + var authParams = (String) conf.get("authParams"); if (StringUtils.hasText(authPluginClassName) && StringUtils.hasText(authParams)) { try { builder.authentication(authPluginClassName, authParams); @@ -120,8 +123,8 @@ public class PulsarAdministration } private void initialize() { - Collection topics = this.applicationContext.getBeansOfType(PulsarTopic.class, false, false) - .values(); + var topics = Objects.requireNonNull(this.applicationContext, "Application context was not set") + .getBeansOfType(PulsarTopic.class, false, false).values(); createOrModifyTopicsIfNeeded(topics); } @@ -161,39 +164,52 @@ public class PulsarAdministration } private void doCreateOrModifyTopicsIfNeeded(PulsarAdmin admin, Collection topics) { - Map> topicsPerNamespace = getTopicsPerNamespace(topics); + var topicsPerNamespace = getTopicsPerNamespace(topics); topicsPerNamespace.forEach((namespace, requestedTopics) -> { - Set topicsToCreate = new HashSet<>(); - Set topicsToModify = new HashSet<>(); + var topicsToCreate = new HashSet(); + var topicsToModify = new HashSet(); try { - List existingTopicsInNamespace = admin.topics().getList(namespace); + var existingTopicsInNamespace = admin.topics().getList(namespace); - for (PulsarTopic topic : requestedTopics) { + for (var topic : requestedTopics) { + var topicName = topic.getFullyQualifiedTopicName(); if (topic.isPartitioned()) { - List matchingPartitions = getMatchingTopicPartitions(topic, existingTopicsInNamespace); + if (existingTopicsInNamespace.contains(topicName)) { + throw new IllegalStateException( + "Topic '%s' already exists un-partitioned - needs to be deleted first" + .formatted(topicName)); + } + var matchingPartitions = getMatchingTopicPartitions(topic, existingTopicsInNamespace); if (matchingPartitions.isEmpty()) { - this.logger.debug(() -> "Topic " + topic.getFullyQualifiedTopicName() + " does not exist."); + this.logger.debug(() -> "Topic '%s' does not yet exist - will add".formatted(topicName)); topicsToCreate.add(topic); } else { - int numberOfExistingPartitions = matchingPartitions.size(); + var numberOfExistingPartitions = matchingPartitions.size(); if (numberOfExistingPartitions < topic.numberOfPartitions()) { - this.logger.debug(() -> "Topic " + topic.getFullyQualifiedTopicName() + " found with " - + numberOfExistingPartitions + " partitions."); + this.logger.debug(() -> "Topic '%s' found with %d partitions - will update to %d" + .formatted(topicName, numberOfExistingPartitions, topic.numberOfPartitions())); topicsToModify.add(topic); } else if (numberOfExistingPartitions > topic.numberOfPartitions()) { - throw new IllegalStateException("Topic " + topic.getFullyQualifiedTopicName() - + " found with " + numberOfExistingPartitions - + " partitions. Needs to be deleted first."); + throw new IllegalStateException( + "Topic '%s' found w/ %d partitions but can't shrink to %d - needs to be deleted first" + .formatted(topicName, numberOfExistingPartitions, + topic.numberOfPartitions())); } } } else { - if (!existingTopicsInNamespace.contains(topic.getFullyQualifiedTopicName())) { - this.logger.debug(() -> "Topic " + topic.getFullyQualifiedTopicName() + " does not exist."); + var matchingPartitions = getMatchingTopicPartitions(topic, existingTopicsInNamespace); + if (!matchingPartitions.isEmpty()) { + throw new IllegalStateException( + "Topic '%s' already exists partitioned - needs to be deleted first" + .formatted(topicName)); + } + if (!existingTopicsInNamespace.contains(topicName)) { + this.logger.debug(() -> "Topic '%s' does not yet exist - will add".formatted(topicName)); topicsToCreate.add(topic); } } @@ -211,7 +227,7 @@ public class PulsarAdministration private void createTopics(PulsarAdmin admin, Set topicsToCreate) throws PulsarAdminException { this.logger.debug(() -> "Creating topics: " + topicsToCreate.stream() .map(PulsarTopic::getFullyQualifiedTopicName).collect(Collectors.joining(","))); - for (PulsarTopic topic : topicsToCreate) { + for (var topic : topicsToCreate) { if (topic.isPartitioned()) { admin.topics().createPartitionedTopic(topic.topicName(), topic.numberOfPartitions()); } @@ -224,7 +240,7 @@ public class PulsarAdministration private void modifyTopics(PulsarAdmin admin, Set topicsToModify) throws PulsarAdminException { this.logger.debug(() -> "Modifying topics: " + topicsToModify.stream() .map(PulsarTopic::getFullyQualifiedTopicName).collect(Collectors.joining(","))); - for (PulsarTopic topic : topicsToModify) { + for (var topic : topicsToModify) { admin.topics().updatePartitionedTopic(topic.topicName(), topic.numberOfPartitions()); } } diff --git a/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarAdministrationTests.java b/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarAdministrationTests.java index 894d7de2..0126f5bf 100644 --- a/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarAdministrationTests.java +++ b/spring-pulsar/src/test/java/org/springframework/pulsar/core/PulsarAdministrationTests.java @@ -52,6 +52,7 @@ import org.springframework.test.context.junit.jupiter.SpringExtension; */ @ExtendWith(SpringExtension.class) @ContextConfiguration +@SuppressWarnings("JUnitMalformedDeclaration") public class PulsarAdministrationTests implements PulsarTestContainerSupport { private static final String NAMESPACE = "public/default"; @@ -112,8 +113,8 @@ public class PulsarAdministrationTests implements PulsarTestContainerSupport { } @Nested - @ContextConfiguration(classes = CreateMissingTopicsTest.CreateMissingTopicsConfig.class) - class CreateMissingTopicsTest { + @ContextConfiguration(classes = CreateMissingTopicsTests.CreateMissingTopicsConfig.class) + class CreateMissingTopicsTests { @Test void topicsExist(@Autowired ObjectProvider expectedTopics) throws Exception { @@ -143,8 +144,8 @@ public class PulsarAdministrationTests implements PulsarTestContainerSupport { } @Nested - @ContextConfiguration(classes = CreateMissingTopicsInSeparateNamespacesTest.CreateMissingTopicsConfig.class) - class CreateMissingTopicsInSeparateNamespacesTest { + @ContextConfiguration(classes = CreateMissingTopicsInSeparateNamespacesTests.CreateMissingTopicsConfig.class) + class CreateMissingTopicsInSeparateNamespacesTests { @Test void topicsExist(@Autowired PulsarTopic partitionedGreenTopic, @Autowired PulsarTopic partitionedBlueTopic) @@ -190,8 +191,8 @@ public class PulsarAdministrationTests implements PulsarTestContainerSupport { } @Nested - @ContextConfiguration(classes = IncrementPartitionCountTest.IncrementPartitionCountConfig.class) - class IncrementPartitionCountTest { + @ContextConfiguration(classes = IncrementPartitionCountTests.IncrementPartitionCountConfig.class) + class IncrementPartitionCountTests { @Test void topicsExist(@Autowired ObjectProvider expectedTopics) throws Exception { @@ -214,8 +215,8 @@ public class PulsarAdministrationTests implements PulsarTestContainerSupport { } @Nested - @ContextConfiguration(classes = DecrementPartitionCountTest.DecrementPartitionCountConfig.class) - class DecrementPartitionCountTest { + @ContextConfiguration(classes = DecrementPartitionCountTests.DecrementPartitionCountConfig.class) + class DecrementPartitionCountTests { @Test void topicModificationThrows(@Autowired ObjectProvider expectedTopics) throws Exception { @@ -223,7 +224,7 @@ public class PulsarAdministrationTests implements PulsarTestContainerSupport { PulsarTopic smallerTopic = PulsarTopic.builder("dpc-partitioned-1").numberOfPartitions(4).build(); assertThatIllegalStateException().isThrownBy(() -> pulsarAdministration.createOrModifyTopics(smallerTopic)) .withMessage( - "Topic persistent://public/default/dpc-partitioned-1 found with 8 partitions. Needs to be deleted first."); + "Topic 'persistent://public/default/dpc-partitioned-1' found w/ 8 partitions but can't shrink to 4 - needs to be deleted first"); } @@ -239,4 +240,30 @@ public class PulsarAdministrationTests implements PulsarTestContainerSupport { } + @Nested + @ContextConfiguration + class ConflictingTopicsTests { + + @Test + void unpartitionedTopicAlreadyExists() { + var unpartitionedTopic = PulsarTopic.builder("ctt-foo").numberOfPartitions(0).build(); + var partitionedTopic = PulsarTopic.builder("ctt-foo").numberOfPartitions(3).build(); + pulsarAdministration.createOrModifyTopics(unpartitionedTopic); + assertThatIllegalStateException() + .isThrownBy(() -> pulsarAdministration.createOrModifyTopics(partitionedTopic)).withMessage( + "Topic 'persistent://public/default/ctt-foo' already exists un-partitioned - needs to be deleted first"); + } + + @Test + void partitionedTopicAlreadyExists() { + var unpartitionedTopic = PulsarTopic.builder("ctt-bar").numberOfPartitions(0).build(); + var partitionedTopic = PulsarTopic.builder("ctt-bar").numberOfPartitions(3).build(); + pulsarAdministration.createOrModifyTopics(partitionedTopic); + assertThatIllegalStateException() + .isThrownBy(() -> pulsarAdministration.createOrModifyTopics(unpartitionedTopic)).withMessage( + "Topic 'persistent://public/default/ctt-bar' already exists partitioned - needs to be deleted first"); + } + + } + }