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 ba6eb1a9..fe08e554 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 @@ -163,10 +163,10 @@ public class PulsarAdministration private void doCreateOrModifyTopicsIfNeeded(PulsarAdmin admin, Collection topics) { Map> topicsPerNamespace = getTopicsPerNamespace(topics); - Set topicsToCreate = new HashSet<>(); - Set topicsToModify = new HashSet<>(); - topicsPerNamespace.forEach((namespace, requestedTopics) -> { + Set topicsToCreate = new HashSet<>(); + Set topicsToModify = new HashSet<>(); + try { List existingTopicsInNamespace = admin.topics().getList(namespace); 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 dac82560..b0c24252 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 @@ -24,6 +24,7 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Set; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminBuilder; @@ -61,20 +62,27 @@ public class PulsarAdministrationTests implements PulsarTestContainerSupport { @Autowired private PulsarAdministration pulsarAdministration; - private void assertThatTopicsExist(List expected) throws PulsarAdminException { - List expectedTopics = expected.stream().mapMulti((topic, consumer) -> { - if (topic.isPartitioned()) { - for (int i = 0; i < topic.numberOfPartitions(); i++) { - consumer.accept(topic.getFullyQualifiedTopicName() + "-partition-" + i); - } - } - else { - consumer.accept(topic.getFullyQualifiedTopicName()); - } + private void assertThatTopicsExist(List expected) throws PulsarAdminException { + assertThatTopicsExistIn(expected, NAMESPACE); + } - }).toList(); - assertThat(pulsarAdminClient.topics().getList(NAMESPACE)).containsAll(expectedTopics); - } + private void assertThatTopicsExistIn(List expected, String namespace) throws PulsarAdminException { + List expectedTopics = expectedTopics(expected); + assertThat(pulsarAdminClient.topics().getList(namespace)).containsAll(expectedTopics); + } + + static List expectedTopics(List expected) { + return expected.stream().mapMulti((topic, consumer) -> { + if (topic.isPartitioned()) { + for (int i = 0; i < topic.numberOfPartitions(); i++) { + consumer.accept(topic.getFullyQualifiedTopicName() + "-partition-" + i); + } + } else { + consumer.accept(topic.getFullyQualifiedTopicName()); + } + + }).toList(); + } @Test void constructorRespectsAuthenticationProps() { @@ -134,6 +142,56 @@ public class PulsarAdministrationTests implements PulsarTestContainerSupport { } + } + + @Nested + @ContextConfiguration(classes = CreateMissingTopicsInSeparateNamespacesTest.CreateMissingTopicsConfig.class) + class CreateMissingTopicsInSeparateNamespacesTest { + + @Test + void topicsExist(@Autowired ObjectProvider expectedTopics) throws Exception { + assertThatTopicsExistIn(expectedTopics.stream().filter(topic -> topic.topicName().contains(CreateMissingTopicsConfig.PUBLIC_BLUE_NAMESPACE)).toList(), CreateMissingTopicsConfig.PUBLIC_BLUE_NAMESPACE); + assertThatTopicsExistIn(expectedTopics.stream().filter(topic -> topic.topicName().contains(CreateMissingTopicsConfig.PUBLIC_GREEN_NAMESPACE)).toList(), CreateMissingTopicsConfig.PUBLIC_GREEN_NAMESPACE); + } + + + @Configuration(proxyBeanMethods = false) + static class CreateMissingTopicsConfig { + + public static final String PUBLIC_GREEN_NAMESPACE = "public/green"; + + public static final String PUBLIC_BLUE_NAMESPACE = "public/blue"; + + static { + var pulsarAdminBuilder = PulsarAdmin.builder().serviceHttpUrl(PulsarTestContainerSupport.getHttpServiceUrl()); + try (var pulsarAdmin = pulsarAdminBuilder.build()) { + Set.of(PUBLIC_GREEN_NAMESPACE, PUBLIC_BLUE_NAMESPACE) + .forEach(ns -> { + try { + pulsarAdmin.namespaces().createNamespace(ns); + } catch (Exception e) { + throw new RuntimeException(e); + } + }); + } catch (Exception e) { + throw new RuntimeException(e); + } + } + + @Bean + PulsarTopic partitionedGreenTopic() { + return PulsarTopic.builder("persistent://%s/partitioned-1".formatted(PUBLIC_GREEN_NAMESPACE)) + .numberOfPartitions(2).build(); + } + + @Bean + PulsarTopic partitionedBlueTopic() { + return PulsarTopic.builder("persistent://%s/partitioned-1".formatted(PUBLIC_BLUE_NAMESPACE)) + .numberOfPartitions(2).build(); + } + + } + } @Nested