From 73547fbd0295f0758919afde7a7e776eb7b7be6a Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 28 Dec 2022 11:46:16 +0800 Subject: [PATCH 01/20] [fix][broker] Reject create non existent partitions. --- .../pulsar/broker/service/BrokerService.java | 28 +++++++-------- .../persistent/PersistentTopicTest.java | 36 +++++++++++++++---- 2 files changed, 43 insertions(+), 21 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index bcec8351733a9..03f803e127879 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -1016,27 +1016,27 @@ public CompletableFuture> getTopic(final TopicName topicName, bo } } final boolean isPersistentTopic = topicName.getDomain().equals(TopicDomain.persistent); - if (isPersistentTopic) { - return topics.computeIfAbsent(topicName.toString(), (k) -> { - return this.loadOrCreatePersistentTopic(k, createIfMissing, properties); - }); - } else { - return topics.computeIfAbsent(topicName.toString(), (name) -> { + return topics.computeIfAbsent(topicName.toString(), (name) -> { + // partitioned topic if (topicName.isPartitioned()) { final TopicName partitionedTopicName = TopicName.get(topicName.getPartitionedTopicName()); - return this.fetchPartitionedTopicMetadataAsync(partitionedTopicName).thenCompose((metadata) -> { + return fetchPartitionedTopicMetadataAsync(partitionedTopicName).thenCompose((metadata) -> { if (topicName.getPartitionIndex() < metadata.partitions) { - return createNonPersistentTopic(name); + return isPersistentTopic ? + loadOrCreatePersistentTopic(name, createIfMissing, properties) : + createNonPersistentTopic(name); } return CompletableFuture.completedFuture(Optional.empty()); }); - } else if (createIfMissing) { - return createNonPersistentTopic(name); - } else { - return CompletableFuture.completedFuture(Optional.empty()); } - }); - } + // non-partitioned topic + return isPersistentTopic ? + // persistent topic + loadOrCreatePersistentTopic(name, createIfMissing, properties) : + // non-persistent topic + (createIfMissing ? createNonPersistentTopic(name) : + CompletableFuture.completedFuture(Optional.empty())); + }); } catch (IllegalArgumentException e) { log.warn("[{}] Illegalargument exception when loading topic", topicName, e); return FutureUtil.failedFuture(e); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java index aa05624a5b0c9..4c6294308d96e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java @@ -29,6 +29,7 @@ import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertNull; +import static org.testng.Assert.fail; import java.io.ByteArrayOutputStream; import java.lang.reflect.Field; import java.nio.charset.StandardCharsets; @@ -38,6 +39,7 @@ import java.util.UUID; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; import com.google.common.collect.Multimap; import com.google.common.collect.Sets; import lombok.Cleanup; @@ -46,6 +48,7 @@ import org.apache.pulsar.broker.service.BrokerTestBase; import org.apache.pulsar.broker.stats.PrometheusMetricsTest; import org.apache.pulsar.broker.stats.prometheus.PrometheusMetricsGenerator; +import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.*; import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.TopicName; @@ -128,8 +131,9 @@ public void testUnblockStuckSubscription() throws Exception { PersistentDispatcherMultipleConsumers sharedDispatcher = (PersistentDispatcherMultipleConsumers) sharedSub .getDispatcher(); - PersistentDispatcherSingleActiveConsumer failOverDispatcher = (PersistentDispatcherSingleActiveConsumer) failOverSub - .getDispatcher(); + PersistentDispatcherSingleActiveConsumer failOverDispatcher = + (PersistentDispatcherSingleActiveConsumer) failOverSub + .getDispatcher(); // build backlog consumer1.close(); @@ -279,14 +283,15 @@ public void testPersistentPartitionedTopicUnload() throws Exception { @DataProvider(name = "topicAndMetricsLevel") public Object[][] indexPatternTestData() { return new Object[][]{ - new Object[] {"persistent://prop/autoNs/test_delayed_message_metric", true}, - new Object[] {"persistent://prop/autoNs/test_delayed_message_metric", false}, + new Object[]{"persistent://prop/autoNs/test_delayed_message_metric", true}, + new Object[]{"persistent://prop/autoNs/test_delayed_message_metric", false}, }; } @Test(dataProvider = "topicAndMetricsLevel") - public void testDelayedDeliveryTrackerMemoryUsageMetric(String topic, boolean exposeTopicLevelMetrics) throws Exception { + public void testDelayedDeliveryTrackerMemoryUsageMetric(String topic, boolean exposeTopicLevelMetrics) + throws Exception { PulsarClient client = pulsar.getClient(); String namespace = TopicName.get(topic).getNamespace(); admin.namespaces().createNamespace(namespace); @@ -365,8 +370,8 @@ public void testUpdateCursorLastActive() throws Exception { Producer producer = pulsarClient.newProducer(Schema.STRING).topic(topicName).create(); PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); - PersistentSubscription persistentSubscription = topic.getSubscription(sharedSubName); - PersistentSubscription persistentSubscription2 = topic.getSubscription(failoverSubName); + PersistentSubscription persistentSubscription = topic.getSubscription(sharedSubName); + PersistentSubscription persistentSubscription2 = topic.getSubscription(failoverSubName); // `addConsumer` should update last active assertTrue(persistentSubscription.getCursor().getLastActive() > beforeAddConsumerTimestamp); @@ -402,4 +407,21 @@ public void testUpdateCursorLastActive() throws Exception { assertTrue(persistentSubscription.getCursor().getLastActive() > beforeRemoveConsumerTimestamp); assertTrue(persistentSubscription2.getCursor().getLastActive() > beforeRemoveConsumerTimestamp); } + + + @Test + public void testCreateNonExistentPartitions() throws PulsarAdminException { + final String topicName = "non-persistent://prop/ns-abc/testCreateNonExistentPartitions"; + admin.topics().createPartitionedTopic(topicName, 4); + TopicName partition = TopicName.get(topicName).getPartition(4); + try { + @Cleanup + Producer producer = pulsarClient.newProducer() + .topic(partition.toString()) + .create(); + fail("unexpected behaviour"); + } catch (PulsarClientException ex) { + Assert.assertTrue(ex instanceof PulsarClientException.TopicDoesNotExistException); + } + } } From a2350f1a4502aa0a4838b8de02a0a212b189821e Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 28 Dec 2022 12:03:57 +0800 Subject: [PATCH 02/20] Revert format change --- .../service/persistent/PersistentTopicTest.java | 16 +++++++--------- 1 file changed, 7 insertions(+), 9 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java index 4c6294308d96e..24a232eac86b1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java @@ -131,9 +131,8 @@ public void testUnblockStuckSubscription() throws Exception { PersistentDispatcherMultipleConsumers sharedDispatcher = (PersistentDispatcherMultipleConsumers) sharedSub .getDispatcher(); - PersistentDispatcherSingleActiveConsumer failOverDispatcher = - (PersistentDispatcherSingleActiveConsumer) failOverSub - .getDispatcher(); + PersistentDispatcherSingleActiveConsumer failOverDispatcher = (PersistentDispatcherSingleActiveConsumer) failOverSub + .getDispatcher(); // build backlog consumer1.close(); @@ -283,15 +282,14 @@ public void testPersistentPartitionedTopicUnload() throws Exception { @DataProvider(name = "topicAndMetricsLevel") public Object[][] indexPatternTestData() { return new Object[][]{ - new Object[]{"persistent://prop/autoNs/test_delayed_message_metric", true}, - new Object[]{"persistent://prop/autoNs/test_delayed_message_metric", false}, + new Object[] {"persistent://prop/autoNs/test_delayed_message_metric", true}, + new Object[] {"persistent://prop/autoNs/test_delayed_message_metric", false}, }; } @Test(dataProvider = "topicAndMetricsLevel") - public void testDelayedDeliveryTrackerMemoryUsageMetric(String topic, boolean exposeTopicLevelMetrics) - throws Exception { + public void testDelayedDeliveryTrackerMemoryUsageMetric(String topic, boolean exposeTopicLevelMetrics) throws Exception { PulsarClient client = pulsar.getClient(); String namespace = TopicName.get(topic).getNamespace(); admin.namespaces().createNamespace(namespace); @@ -370,8 +368,8 @@ public void testUpdateCursorLastActive() throws Exception { Producer producer = pulsarClient.newProducer(Schema.STRING).topic(topicName).create(); PersistentTopic topic = (PersistentTopic) pulsar.getBrokerService().getTopicReference(topicName).get(); - PersistentSubscription persistentSubscription = topic.getSubscription(sharedSubName); - PersistentSubscription persistentSubscription2 = topic.getSubscription(failoverSubName); + PersistentSubscription persistentSubscription = topic.getSubscription(sharedSubName); + PersistentSubscription persistentSubscription2 = topic.getSubscription(failoverSubName); // `addConsumer` should update last active assertTrue(persistentSubscription.getCursor().getLastActive() > beforeAddConsumerTimestamp); From 3c7b9e95d953049e69018429757d290bcc7534b4 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 28 Dec 2022 12:09:51 +0800 Subject: [PATCH 03/20] Add partition number test --- .../pulsar/broker/service/persistent/PersistentTopicTest.java | 1 + 1 file changed, 1 insertion(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java index 24a232eac86b1..df0a5290e0d6e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java @@ -421,5 +421,6 @@ public void testCreateNonExistentPartitions() throws PulsarAdminException { } catch (PulsarClientException ex) { Assert.assertTrue(ex instanceof PulsarClientException.TopicDoesNotExistException); } + Assert.assertEquals(admin.topics().getPartitionedTopicMetadata(topicName).partitions, 4); } } From e3763c48e69c205fe20e97323178aeef39aa5c9e Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 28 Dec 2022 21:13:30 +0800 Subject: [PATCH 04/20] Fix test --- .../pulsar/broker/service/persistent/PersistentTopicTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java index df0a5290e0d6e..cf0a0c3678adb 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java @@ -409,7 +409,7 @@ public void testUpdateCursorLastActive() throws Exception { @Test public void testCreateNonExistentPartitions() throws PulsarAdminException { - final String topicName = "non-persistent://prop/ns-abc/testCreateNonExistentPartitions"; + final String topicName = "persistent://prop/ns-abc/testCreateNonExistentPartitions"; admin.topics().createPartitionedTopic(topicName, 4); TopicName partition = TopicName.get(topicName).getPartition(4); try { From 75e88b03126c003876674cce8bf93c626d093404 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 28 Dec 2022 21:15:50 +0800 Subject: [PATCH 05/20] Apply comment --- .../broker/service/persistent/PersistentTopicTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java index cf0a0c3678adb..cdf7fb4eb011f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java @@ -408,7 +408,7 @@ public void testUpdateCursorLastActive() throws Exception { @Test - public void testCreateNonExistentPartitions() throws PulsarAdminException { + public void testCreateNonExistentPartitions() throws PulsarAdminException, PulsarClientException { final String topicName = "persistent://prop/ns-abc/testCreateNonExistentPartitions"; admin.topics().createPartitionedTopic(topicName, 4); TopicName partition = TopicName.get(topicName).getPartition(4); @@ -418,8 +418,8 @@ public void testCreateNonExistentPartitions() throws PulsarAdminException { .topic(partition.toString()) .create(); fail("unexpected behaviour"); - } catch (PulsarClientException ex) { - Assert.assertTrue(ex instanceof PulsarClientException.TopicDoesNotExistException); + } catch (PulsarClientException.TopicDoesNotExistException ignored) { + } Assert.assertEquals(admin.topics().getPartitionedTopicMetadata(topicName).partitions, 4); } From 5d0668847eb12e60594af9ec73de65688745a262 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 28 Dec 2022 21:24:08 +0800 Subject: [PATCH 06/20] Fix checkstyle --- .../org/apache/pulsar/broker/service/BrokerService.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 03f803e127879..aea1c12c4d391 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -1022,17 +1022,17 @@ public CompletableFuture> getTopic(final TopicName topicName, bo final TopicName partitionedTopicName = TopicName.get(topicName.getPartitionedTopicName()); return fetchPartitionedTopicMetadataAsync(partitionedTopicName).thenCompose((metadata) -> { if (topicName.getPartitionIndex() < metadata.partitions) { - return isPersistentTopic ? - loadOrCreatePersistentTopic(name, createIfMissing, properties) : + return isPersistentTopic + ? loadOrCreatePersistentTopic(name, createIfMissing, properties) : createNonPersistentTopic(name); } return CompletableFuture.completedFuture(Optional.empty()); }); } // non-partitioned topic - return isPersistentTopic ? + return isPersistentTopic // persistent topic - loadOrCreatePersistentTopic(name, createIfMissing, properties) : + ? loadOrCreatePersistentTopic(name, createIfMissing, properties) : // non-persistent topic (createIfMissing ? createNonPersistentTopic(name) : CompletableFuture.completedFuture(Optional.empty())); From 390f12e35a8d91c6a50ed7f87241f4e71a817dcd Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Wed, 28 Dec 2022 21:29:33 +0800 Subject: [PATCH 07/20] Fix checkstyle --- .../pulsar/broker/service/persistent/PersistentTopicTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java index cdf7fb4eb011f..19c5bd5c9aa1d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java @@ -39,7 +39,6 @@ import java.util.UUID; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; -import java.util.stream.Collectors; import com.google.common.collect.Multimap; import com.google.common.collect.Sets; import lombok.Cleanup; From 2d68680f7d051369390f8aa79c1fe0c29c212a28 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Thu, 29 Dec 2022 10:06:43 +0800 Subject: [PATCH 08/20] Fix unexpected behaviour --- .../admin/impl/PersistentTopicsBase.java | 209 ++++++++---------- 1 file changed, 98 insertions(+), 111 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 3fb551967b94e..a7a71cc179cd3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -67,10 +67,12 @@ import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.PulsarService; +import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.admin.AdminResource; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; import org.apache.pulsar.broker.authorization.AuthorizationService; import org.apache.pulsar.broker.service.AnalyzeBacklogResult; +import org.apache.pulsar.broker.service.BrokerService; import org.apache.pulsar.broker.service.BrokerServiceException.AlreadyRunningException; import org.apache.pulsar.broker.service.BrokerServiceException.SubscriptionBusyException; import org.apache.pulsar.broker.service.BrokerServiceException.SubscriptionInvalidCursorPosition; @@ -411,34 +413,35 @@ protected CompletableFuture internalCreateNonPartitionedTopicAsync(boolean * recreate them at application so, newly created producers and consumers can connect to newly added partitions as * well. Therefore, it can violate partition ordering at producers until all producers are restarted at application. * - * @param numPartitions + * @param expectPartitions * @param updateLocalTopicOnly * @param authoritative * @param force */ - protected CompletableFuture internalUpdatePartitionedTopicAsync(int numPartitions, + protected CompletableFuture internalUpdatePartitionedTopicAsync(int expectPartitions, boolean updateLocalTopicOnly, boolean authoritative, boolean force) { - if (numPartitions <= 0) { - return FutureUtil.failedFuture(new RestException(Status.NOT_ACCEPTABLE, - "Number of partitions should be more than 0")); + if (expectPartitions <= 0) { + return FutureUtil.failedFuture(new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be more than 0")); } + final BrokerService brokerService = pulsar().getBrokerService(); + ServiceConfiguration configuration = pulsar().getConfiguration(); return validateTopicOwnershipAsync(topicName, authoritative) - .thenCompose(__ -> validateTopicPolicyOperationAsync(topicName, PolicyName.PARTITION, - PolicyOperation.WRITE)) + .thenCompose(__ -> validateTopicPolicyOperationAsync(topicName, PolicyName.PARTITION, PolicyOperation.WRITE)) .thenCompose(__ -> { if (!updateLocalTopicOnly && !force) { - return validatePartitionTopicUpdateAsync(topicName.getLocalName(), numPartitions); + return validatePartitionTopicUpdateAsync(topicName.getLocalName(), expectPartitions); } else { return CompletableFuture.completedFuture(null); } - }).thenCompose(__ -> pulsar().getBrokerService().fetchPartitionedTopicMetadataAsync(topicName)) + }).thenCompose(__ -> brokerService.fetchPartitionedTopicMetadataAsync(topicName)) .thenCompose(topicMetadata -> { - final int maxPartitions = pulsar().getConfig().getMaxNumPartitionsPerPartitionedTopic(); - if (maxPartitions > 0 && numPartitions > maxPartitions) { + final int maxPartitions = configuration.getMaxNumPartitionsPerPartitionedTopic(); + if (maxPartitions > 0 && expectPartitions > maxPartitions) { throw new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be less than or equal to " + maxPartitions); } + final int previousPartitions = topicMetadata.partitions; // Only do the validation if it's the first hop. if (topicName.isGlobal() && isNamespaceReplicated(topicName.getNamespaceObject())) { return getNamespaceReplicatedClustersAsync(topicName.getNamespaceObject()) @@ -452,25 +455,24 @@ protected CompletableFuture internalUpdatePartitionedTopicAsync(int numPar return clusters; }) .thenCompose(clusters -> - tryCreateExtendedPartitionsAsync(topicMetadata.partitions, numPartitions) + tryCreateExtendedPartitionsAsync(topicMetadata.partitions, expectPartitions) .thenApply(ignore -> clusters)) - .thenCompose(clusters -> createSubscriptions(topicName, numPartitions, force).thenApply( - ignore -> clusters)) .thenCompose(clusters -> { if (!updateLocalTopicOnly) { - return updatePartitionInOtherCluster(numPartitions, clusters) + return updatePartitionInOtherCluster(expectPartitions, clusters) .thenCompose(v -> namespaceResources().getPartitionedTopicResources() .updatePartitionedTopicAsync(topicName, p -> - new PartitionedTopicMetadata(numPartitions, + new PartitionedTopicMetadata(expectPartitions, p.properties) )); } else { return CompletableFuture.completedFuture(null); } - }); + }).thenCompose(clusters -> createSubscriptions(topicName, previousPartitions, + expectPartitions, force)); } else { - return tryCreateExtendedPartitionsAsync(topicMetadata.partitions, numPartitions) - .thenCompose(ignore -> updatePartitionedTopic(topicName, numPartitions, force)); + return tryCreateExtendedPartitionsAsync(previousPartitions, expectPartitions) + .thenCompose(ignore -> updatePartitionedTopic(topicName, previousPartitions, expectPartitions, force)); } }); } @@ -4363,124 +4365,109 @@ private PersistentReplicator getReplicatorReference(String replName, PersistentT } } - private CompletableFuture updatePartitionedTopic(TopicName topicName, int numPartitions, boolean force) { - CompletableFuture result = new CompletableFuture<>(); - createSubscriptions(topicName, numPartitions, force).thenCompose(__ -> { - CompletableFuture future = namespaceResources().getPartitionedTopicResources() - .updatePartitionedTopicAsync(topicName, p -> - new PartitionedTopicMetadata(numPartitions, p.properties)); - future.exceptionally(ex -> { - // If the update operation fails, clean up the partitions that were created - getPartitionedTopicMetadataAsync(topicName, false, false).thenAccept(metadata -> { - int oldPartition = metadata.partitions; - for (int i = oldPartition; i < numPartitions; i++) { - topicResources().deletePersistentTopicAsync(topicName.getPartition(i)).exceptionally(ex1 -> { - log.warn("[{}] Failed to clean up managedLedger {}", clientAppId(), topicName, - ex1.getCause()); - return null; - }); - } - }).exceptionally(e -> { - log.warn("[{}] Failed to clean up managedLedger", topicName, e); - return null; - }); + private CompletableFuture updatePartitionedTopic(TopicName topicName,int previousPartitions, + int expectPartitions, boolean force) { + CompletableFuture future = namespaceResources().getPartitionedTopicResources() + .updatePartitionedTopicAsync(topicName, p -> new PartitionedTopicMetadata(expectPartitions, p.properties)); + future.exceptionally(ex -> { + // If the update operation fails, clean up the partitions that were created + getPartitionedTopicMetadataAsync(topicName, false, false) + .thenAccept(metadata -> { + int oldPartition = metadata.partitions; + for (int i = oldPartition; i < expectPartitions; i++) { + topicResources().deletePersistentTopicAsync(topicName.getPartition(i)).exceptionally(ex1 -> { + log.warn("[{}] Failed to clean up managedLedger {}", clientAppId(), topicName, + ex1.getCause()); + return null; + }); + } + }).exceptionally(e -> { + log.warn("[{}] Failed to clean up managedLedger", topicName, e); return null; }); - return future; - }).thenAccept(__ -> result.complete(null)).exceptionally(ex -> { - result.completeExceptionally(ex); return null; }); - return result; + return future.thenCompose(__ -> createSubscriptions(topicName, previousPartitions, expectPartitions, force)); } /** * It creates subscriptions for new partitions of existing partitioned-topics. * * @param topicName : topic-name: persistent://prop/cluster/ns/topic - * @param numPartitions : number partitions for the topics + * @param previousPartitions number of previous partitions + * @param expectPartitions : number of expected partitions * @param ignoreConflictException : If true, ignore ConflictException: subscription already exists for topic * */ - private CompletableFuture createSubscriptions(TopicName topicName, int numPartitions, - boolean ignoreConflictException) { + private CompletableFuture createSubscriptions(TopicName topicName,int previousPartitions, + int expectPartitions, boolean ignoreConflictException) { CompletableFuture result = new CompletableFuture<>(); - pulsar().getBrokerService().fetchPartitionedTopicMetadataAsync(topicName).thenAccept(partitionMetadata -> { - if (partitionMetadata.partitions < 1) { - result.completeExceptionally(new RestException(Status.CONFLICT, "Topic is not partitioned topic")); - return; - } - - if (partitionMetadata.partitions >= numPartitions) { - result.completeExceptionally(new RestException(Status.CONFLICT, - "number of partitions must be more than existing " + partitionMetadata.partitions)); - return; - } + if (previousPartitions < 1) { + return FutureUtil.failedFuture(new RestException(Status.CONFLICT, "Topic is not partitioned topic")); + } - PulsarAdmin admin; - try { - admin = pulsar().getAdminClient(); - } catch (PulsarServerException e1) { - result.completeExceptionally(e1); - return; - } + if (previousPartitions >= expectPartitions) { + return FutureUtil.failedFuture(new RestException(Status.CONFLICT, + "number of partitions must be more than existing " + previousPartitions)); + } - admin.topics().getStatsAsync(topicName.getPartition(0).toString()).thenAccept(stats -> { - List> subscriptionFutures = new ArrayList<>(); + PulsarAdmin admin; + try { + admin = pulsar().getAdminClient(); + } catch (PulsarServerException e1) { + return FutureUtil.failedFuture(e1); + } - stats.getSubscriptions().entrySet().forEach(e -> { - String subscription = e.getKey(); - SubscriptionStats ss = e.getValue(); - if (!ss.isDurable()) { - // We must not re-create non-durable subscriptions on the new partitions - return; - } - boolean replicated = ss.isReplicated(); + admin.topics().getStatsAsync(topicName.getPartition(0).toString()).thenAccept(stats -> { + List> subscriptionFutures = new ArrayList<>(); - for (int i = partitionMetadata.partitions; i < numPartitions; i++) { - final String topicNamePartition = topicName.getPartition(i).toString(); - CompletableFuture future = new CompletableFuture<>(); - admin.topics().createSubscriptionAsync(topicNamePartition, - subscription, MessageId.latest, replicated).whenComplete((__, ex) -> { - if (ex == null) { + stats.getSubscriptions().entrySet().forEach(e -> { + String subscription = e.getKey(); + SubscriptionStats ss = e.getValue(); + if (!ss.isDurable()) { + // We must not re-create non-durable subscriptions on the new partitions + return; + } + boolean replicated = ss.isReplicated(); + + for (int i = previousPartitions; i < expectPartitions; i++) { + final String topicNamePartition = topicName.getPartition(i).toString(); + CompletableFuture future = new CompletableFuture<>(); + admin.topics().createSubscriptionAsync(topicNamePartition, + subscription, MessageId.latest, replicated).whenComplete((__, ex) -> { + if (ex == null) { + future.complete(null); + } else { + if (ignoreConflictException + && ex instanceof PulsarAdminException.ConflictException) { future.complete(null); } else { - if (ignoreConflictException - && ex instanceof PulsarAdminException.ConflictException) { - future.complete(null); - } else { - future.completeExceptionally(ex); - } + future.completeExceptionally(ex); } - }); - subscriptionFutures.add(future); - } - }); + } + }); + subscriptionFutures.add(future); + } + }); - FutureUtil.waitForAll(subscriptionFutures).thenRun(() -> { - log.info("[{}] Successfully created subscriptions on new partitions {}", clientAppId(), topicName); - result.complete(null); - }).exceptionally(ex -> { - log.warn("[{}] Failed to create subscriptions on new partitions for {}", - clientAppId(), topicName, ex); - result.completeExceptionally(ex); - return null; - }); + FutureUtil.waitForAll(subscriptionFutures).thenRun(() -> { + log.info("[{}] Successfully created subscriptions on new partitions {}", clientAppId(), topicName); + result.complete(null); }).exceptionally(ex -> { - if (ex.getCause() instanceof PulsarAdminException.NotFoundException) { - // The first partition doesn't exist, so there are currently to subscriptions to recreate - result.complete(null); - } else { - log.warn("[{}] Failed to get list of subscriptions of {}", - clientAppId(), topicName.getPartition(0), ex); - result.completeExceptionally(ex); - } + log.warn("[{}] Failed to create subscriptions on new partitions for {}", + clientAppId(), topicName, ex); + result.completeExceptionally(ex); return null; }); }).exceptionally(ex -> { - log.warn("[{}] Failed to get partition metadata for {}", - clientAppId(), topicName.toString()); - result.completeExceptionally(ex); + if (ex.getCause() instanceof PulsarAdminException.NotFoundException) { + // The first partition doesn't exist, so there are currently to subscriptions to recreate + result.complete(null); + } else { + log.warn("[{}] Failed to get list of subscriptions of {}", + clientAppId(), topicName.getPartition(0), ex); + result.completeExceptionally(ex); + } return null; }); return result; From 8266e890c7d7ef64a6bd52900b0bd994f5f78dd4 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Thu, 29 Dec 2022 10:45:20 +0800 Subject: [PATCH 09/20] Fix checkstyle --- .../admin/impl/PersistentTopicsBase.java | 19 ++++++++++++------- 1 file changed, 12 insertions(+), 7 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index a7a71cc179cd3..3bad9ff19cb2e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -422,12 +422,14 @@ protected CompletableFuture internalUpdatePartitionedTopicAsync(int expect boolean updateLocalTopicOnly, boolean authoritative, boolean force) { if (expectPartitions <= 0) { - return FutureUtil.failedFuture(new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be more than 0")); + return FutureUtil.failedFuture( + new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be more than 0")); } final BrokerService brokerService = pulsar().getBrokerService(); ServiceConfiguration configuration = pulsar().getConfiguration(); return validateTopicOwnershipAsync(topicName, authoritative) - .thenCompose(__ -> validateTopicPolicyOperationAsync(topicName, PolicyName.PARTITION, PolicyOperation.WRITE)) + .thenCompose(__ -> + validateTopicPolicyOperationAsync(topicName, PolicyName.PARTITION, PolicyOperation.WRITE)) .thenCompose(__ -> { if (!updateLocalTopicOnly && !force) { return validatePartitionTopicUpdateAsync(topicName.getLocalName(), expectPartitions); @@ -472,13 +474,15 @@ protected CompletableFuture internalUpdatePartitionedTopicAsync(int expect expectPartitions, force)); } else { return tryCreateExtendedPartitionsAsync(previousPartitions, expectPartitions) - .thenCompose(ignore -> updatePartitionedTopic(topicName, previousPartitions, expectPartitions, force)); + .thenCompose(ignore -> + updatePartitionedTopic(topicName, previousPartitions, expectPartitions, force)); } }); } protected void internalCreateMissedPartitions(AsyncResponse asyncResponse) { - getPartitionedTopicMetadataAsync(topicName, false, false).thenAccept(metadata -> { + getPartitionedTopicMetadataAsync(topicName, false, false) + .thenAccept(metadata -> { if (metadata != null) { tryCreatePartitionsAsync(metadata.partitions).thenAccept(v -> { asyncResponse.resume(Response.noContent().build()); @@ -4365,10 +4369,11 @@ private PersistentReplicator getReplicatorReference(String replName, PersistentT } } - private CompletableFuture updatePartitionedTopic(TopicName topicName,int previousPartitions, + private CompletableFuture updatePartitionedTopic(TopicName topicName, int previousPartitions, int expectPartitions, boolean force) { CompletableFuture future = namespaceResources().getPartitionedTopicResources() - .updatePartitionedTopicAsync(topicName, p -> new PartitionedTopicMetadata(expectPartitions, p.properties)); + .updatePartitionedTopicAsync(topicName, p -> + new PartitionedTopicMetadata(expectPartitions, p.properties)); future.exceptionally(ex -> { // If the update operation fails, clean up the partitions that were created getPartitionedTopicMetadataAsync(topicName, false, false) @@ -4399,7 +4404,7 @@ private CompletableFuture updatePartitionedTopic(TopicName topicName,int p * @param ignoreConflictException : If true, ignore ConflictException: subscription already exists for topic * */ - private CompletableFuture createSubscriptions(TopicName topicName,int previousPartitions, + private CompletableFuture createSubscriptions(TopicName topicName, int previousPartitions, int expectPartitions, boolean ignoreConflictException) { CompletableFuture result = new CompletableFuture<>(); if (previousPartitions < 1) { From 3c5e3aa7e526cabb4eb71ef7df36e64439e9e484 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Thu, 29 Dec 2022 11:58:16 +0800 Subject: [PATCH 10/20] Fix some failed test --- .../org/apache/pulsar/broker/admin/AdminApi2Test.java | 11 ++++------- .../broker/systopic/PartitionedSystemTopicTest.java | 3 +++ 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java index 11c84d990f68d..ca9bb6950b8c1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java @@ -2676,15 +2676,12 @@ public void testFailedUpdatePartitionedTopic() throws Exception { assertEquals(admin.topics().getPartitionedTopicMetadata(partitionedTopicName).partitions, startPartitions); // create a subscription for few new partition which can fail - admin.topics().createSubscription(partitionedTopicName + "-partition-" + startPartitions, subName1, - MessageId.earliest); - try { - admin.topics().updatePartitionedTopic(partitionedTopicName, newPartitions, false, false); - } catch (PulsarAdminException.PreconditionFailedException e) { - // Ok + admin.topics().createSubscription(partitionedTopicName + "-partition-" + startPartitions, subName1, + MessageId.earliest); + } catch (PulsarAdminException.PreconditionFailedException ex) { + // OK } - assertEquals(admin.topics().getPartitionedTopicMetadata(partitionedTopicName).partitions, startPartitions); admin.topics().updatePartitionedTopic(partitionedTopicName, newPartitions, false, true); // validate subscription is created for new partition. diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java index 008c2143a3566..39bcd0d99a85f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java @@ -52,6 +52,7 @@ import org.apache.pulsar.common.naming.SystemTopicNames; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.naming.TopicVersion; +import org.apache.pulsar.common.partition.PartitionedTopicMetadata; import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.TenantInfo; import org.apache.pulsar.common.policies.data.TenantInfoImpl; @@ -194,6 +195,8 @@ public void testHeartbeatTopicNotAllowedToSendEvent() throws Exception { NamespaceName namespaceName = NamespaceService.getHeartbeatNamespaceV2(pulsar.getAdvertisedAddress(), pulsar.getConfig()); TopicName topicName = TopicName.get("persistent", namespaceName, SystemTopicNames.NAMESPACE_EVENTS_LOCAL_NAME); + pulsar.getPulsarResources().getNamespaceResources().getPartitionedTopicResources() + .createPartitionedTopic(topicName, new PartitionedTopicMetadata(PARTITIONS)); for (int partition = 0; partition < PARTITIONS; partition ++) { pulsar.getBrokerService() .getTopic(topicName.getPartition(partition).toString(), true).join(); From 9f49db0fc1b01896d87e2340bb6be843a60eae17 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 30 Dec 2022 09:59:13 +0800 Subject: [PATCH 11/20] Rollback some logic --- .../pulsar/broker/service/BrokerService.java | 35 ++++++++++++------- 1 file changed, 23 insertions(+), 12 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index aea1c12c4d391..93ea33ff6069b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -1016,27 +1016,38 @@ public CompletableFuture> getTopic(final TopicName topicName, bo } } final boolean isPersistentTopic = topicName.getDomain().equals(TopicDomain.persistent); + if (isPersistentTopic) { + return topics.computeIfAbsent(topicName.toString(), (tpName) -> { + if (topicName.isPartitioned()) { + return this.fetchPartitionedTopicMetadataAsync(TopicName.get(topicName.getPartitionedTopicName())) + .thenCompose((metadata) -> { + // Allow crate non-partitioned persistent topic that name includes `partition` + if (metadata.partitions == 0 || + topicName.getPartitionIndex() < metadata.partitions) { + return loadOrCreatePersistentTopic(tpName, createIfMissing, properties); + } + return CompletableFuture.completedFuture(Optional.empty()); + }); + } + return loadOrCreatePersistentTopic(tpName, createIfMissing, properties); + }); + } else { return topics.computeIfAbsent(topicName.toString(), (name) -> { - // partitioned topic if (topicName.isPartitioned()) { final TopicName partitionedTopicName = TopicName.get(topicName.getPartitionedTopicName()); - return fetchPartitionedTopicMetadataAsync(partitionedTopicName).thenCompose((metadata) -> { + return this.fetchPartitionedTopicMetadataAsync(partitionedTopicName).thenCompose((metadata) -> { if (topicName.getPartitionIndex() < metadata.partitions) { - return isPersistentTopic - ? loadOrCreatePersistentTopic(name, createIfMissing, properties) : - createNonPersistentTopic(name); + return createNonPersistentTopic(name); } return CompletableFuture.completedFuture(Optional.empty()); }); + } else if (createIfMissing) { + return createNonPersistentTopic(name); + } else { + return CompletableFuture.completedFuture(Optional.empty()); } - // non-partitioned topic - return isPersistentTopic - // persistent topic - ? loadOrCreatePersistentTopic(name, createIfMissing, properties) : - // non-persistent topic - (createIfMissing ? createNonPersistentTopic(name) : - CompletableFuture.completedFuture(Optional.empty())); }); + } } catch (IllegalArgumentException e) { log.warn("[{}] Illegalargument exception when loading topic", topicName, e); return FutureUtil.failedFuture(e); From c77a7744e299243a4efb9c06538b55d7006ddb18 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 30 Dec 2022 10:03:32 +0800 Subject: [PATCH 12/20] Add non-persistent topic test --- .../nonpersistent/NonPersistentTopicTest.java | 23 +++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopicTest.java index 71caa1edb527c..73a1084f30f2a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentTopicTest.java @@ -18,13 +18,18 @@ */ package org.apache.pulsar.broker.service.nonpersistent; +import lombok.Cleanup; import org.apache.pulsar.broker.service.BrokerTestBase; +import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.Producer; +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.naming.TopicName; import org.apache.pulsar.common.policies.data.TopicStats; +import org.junit.Assert; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; @@ -32,6 +37,7 @@ import static org.testng.Assert.assertTrue; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotNull; +import static org.testng.Assert.fail; @Test(groups = "broker") public class NonPersistentTopicTest extends BrokerTestBase { @@ -96,4 +102,21 @@ public void testAccumulativeStats() throws Exception { assertEquals(statsAfterUnsubscribe.getBytesOutCounter(), statsBeforeUnsubscribe.getBytesOutCounter()); assertEquals(statsAfterUnsubscribe.getMsgOutCounter(), statsBeforeUnsubscribe.getMsgOutCounter()); } + + @Test + public void testCreateNonExistentPartitions() throws PulsarAdminException, PulsarClientException { + final String topicName = "non-persistent://prop/ns-abc/testCreateNonExistentPartitions"; + admin.topics().createPartitionedTopic(topicName, 4); + TopicName partition = TopicName.get(topicName).getPartition(4); + try { + @Cleanup + Producer producer = pulsarClient.newProducer() + .topic(partition.toString()) + .create(); + fail("unexpected behaviour"); + } catch (PulsarClientException.TopicDoesNotExistException ignored) { + + } + Assert.assertEquals(admin.topics().getPartitionedTopicMetadata(topicName).partitions, 4); + } } From 6fce4df6993c1d70a60f6709c4631c76ae8a3ddb Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 30 Dec 2022 10:15:25 +0800 Subject: [PATCH 13/20] Fix checkstyle --- .../org/apache/pulsar/broker/service/BrokerService.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 93ea33ff6069b..77f26e4e17a7f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -1019,11 +1019,11 @@ public CompletableFuture> getTopic(final TopicName topicName, bo if (isPersistentTopic) { return topics.computeIfAbsent(topicName.toString(), (tpName) -> { if (topicName.isPartitioned()) { - return this.fetchPartitionedTopicMetadataAsync(TopicName.get(topicName.getPartitionedTopicName())) + return fetchPartitionedTopicMetadataAsync(TopicName.get(topicName.getPartitionedTopicName())) .thenCompose((metadata) -> { // Allow crate non-partitioned persistent topic that name includes `partition` - if (metadata.partitions == 0 || - topicName.getPartitionIndex() < metadata.partitions) { + if (metadata.partitions == 0 + || topicName.getPartitionIndex() < metadata.partitions) { return loadOrCreatePersistentTopic(tpName, createIfMissing, properties); } return CompletableFuture.completedFuture(Optional.empty()); From 655a6abb05b6bb538d3c5ef2b2d1b3e4e98e3890 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 30 Dec 2022 10:54:02 +0800 Subject: [PATCH 14/20] Revert test fix --- .../pulsar/broker/systopic/PartitionedSystemTopicTest.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java index 39bcd0d99a85f..88f10df842285 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java @@ -195,8 +195,6 @@ public void testHeartbeatTopicNotAllowedToSendEvent() throws Exception { NamespaceName namespaceName = NamespaceService.getHeartbeatNamespaceV2(pulsar.getAdvertisedAddress(), pulsar.getConfig()); TopicName topicName = TopicName.get("persistent", namespaceName, SystemTopicNames.NAMESPACE_EVENTS_LOCAL_NAME); - pulsar.getPulsarResources().getNamespaceResources().getPartitionedTopicResources() - .createPartitionedTopic(topicName, new PartitionedTopicMetadata(PARTITIONS)); for (int partition = 0; partition < PARTITIONS; partition ++) { pulsar.getBrokerService() .getTopic(topicName.getPartition(partition).toString(), true).join(); From 5b460cbcb3f899c517f200ecdfbf3e96a6093460 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 30 Dec 2022 10:56:37 +0800 Subject: [PATCH 15/20] Revert some logic to help review --- .../pulsar/broker/admin/impl/PersistentTopicsBase.java | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 3bad9ff19cb2e..52bb1364008c8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -425,8 +425,6 @@ protected CompletableFuture internalUpdatePartitionedTopicAsync(int expect return FutureUtil.failedFuture( new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be more than 0")); } - final BrokerService brokerService = pulsar().getBrokerService(); - ServiceConfiguration configuration = pulsar().getConfiguration(); return validateTopicOwnershipAsync(topicName, authoritative) .thenCompose(__ -> validateTopicPolicyOperationAsync(topicName, PolicyName.PARTITION, PolicyOperation.WRITE)) @@ -436,9 +434,9 @@ protected CompletableFuture internalUpdatePartitionedTopicAsync(int expect } else { return CompletableFuture.completedFuture(null); } - }).thenCompose(__ -> brokerService.fetchPartitionedTopicMetadataAsync(topicName)) + }).thenCompose(__ -> pulsar().getBrokerService().fetchPartitionedTopicMetadataAsync(topicName)) .thenCompose(topicMetadata -> { - final int maxPartitions = configuration.getMaxNumPartitionsPerPartitionedTopic(); + final int maxPartitions = pulsar().getConfig().getMaxNumPartitionsPerPartitionedTopic(); if (maxPartitions > 0 && expectPartitions > maxPartitions) { throw new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be less than or equal to " + maxPartitions); @@ -481,8 +479,7 @@ protected CompletableFuture internalUpdatePartitionedTopicAsync(int expect } protected void internalCreateMissedPartitions(AsyncResponse asyncResponse) { - getPartitionedTopicMetadataAsync(topicName, false, false) - .thenAccept(metadata -> { + getPartitionedTopicMetadataAsync(topicName, false, false).thenAccept(metadata -> { if (metadata != null) { tryCreatePartitionsAsync(metadata.partitions).thenAccept(v -> { asyncResponse.resume(Response.noContent().build()); From 3d59543903bc10b3fd90336cd3e72c4677a02fb0 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 30 Dec 2022 10:57:19 +0800 Subject: [PATCH 16/20] Remove unused import --- .../pulsar/broker/systopic/PartitionedSystemTopicTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java index 88f10df842285..008c2143a3566 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/systopic/PartitionedSystemTopicTest.java @@ -52,7 +52,6 @@ import org.apache.pulsar.common.naming.SystemTopicNames; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.naming.TopicVersion; -import org.apache.pulsar.common.partition.PartitionedTopicMetadata; import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.TenantInfo; import org.apache.pulsar.common.policies.data.TenantInfoImpl; From bf674734f2fba299127b420c3083964dbcd8c80c Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 30 Dec 2022 11:06:55 +0800 Subject: [PATCH 17/20] Fix checkstyle --- .../apache/pulsar/broker/admin/impl/PersistentTopicsBase.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 52bb1364008c8..ecc2ab3cf4e88 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -67,12 +67,10 @@ import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.PulsarService; -import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.admin.AdminResource; import org.apache.pulsar.broker.authentication.AuthenticationDataSource; import org.apache.pulsar.broker.authorization.AuthorizationService; import org.apache.pulsar.broker.service.AnalyzeBacklogResult; -import org.apache.pulsar.broker.service.BrokerService; import org.apache.pulsar.broker.service.BrokerServiceException.AlreadyRunningException; import org.apache.pulsar.broker.service.BrokerServiceException.SubscriptionBusyException; import org.apache.pulsar.broker.service.BrokerServiceException.SubscriptionInvalidCursorPosition; From 68f2bd55e310ae38da07575fc1a4fb9cbc6b684e Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 30 Dec 2022 16:23:12 +0800 Subject: [PATCH 18/20] Fix failed test --- .../admin/impl/PersistentTopicsBase.java | 47 +++++++------------ 1 file changed, 16 insertions(+), 31 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index ecc2ab3cf4e88..b5cd83ccf8219 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -439,7 +439,6 @@ protected CompletableFuture internalUpdatePartitionedTopicAsync(int expect throw new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be less than or equal to " + maxPartitions); } - final int previousPartitions = topicMetadata.partitions; // Only do the validation if it's the first hop. if (topicName.isGlobal() && isNamespaceReplicated(topicName.getNamespaceObject())) { return getNamespaceReplicatedClustersAsync(topicName.getNamespaceObject()) @@ -453,25 +452,22 @@ protected CompletableFuture internalUpdatePartitionedTopicAsync(int expect return clusters; }) .thenCompose(clusters -> - tryCreateExtendedPartitionsAsync(topicMetadata.partitions, expectPartitions) - .thenApply(ignore -> clusters)) + tryCreatePartitionsAsync(expectPartitions).thenApply(ignore -> clusters)) .thenCompose(clusters -> { if (!updateLocalTopicOnly) { - return updatePartitionInOtherCluster(expectPartitions, clusters) - .thenCompose(v -> namespaceResources().getPartitionedTopicResources() - .updatePartitionedTopicAsync(topicName, p -> - new PartitionedTopicMetadata(expectPartitions, - p.properties) - )); + return namespaceResources().getPartitionedTopicResources() + .updatePartitionedTopicAsync(topicName, p -> + new PartitionedTopicMetadata(expectPartitions, p.properties)) + .thenCompose(__ -> + updatePartitionInOtherCluster(expectPartitions, clusters)); } else { return CompletableFuture.completedFuture(null); } - }).thenCompose(clusters -> createSubscriptions(topicName, previousPartitions, - expectPartitions, force)); + }).thenCompose(clusters -> createSubscriptions(topicName, + expectPartitions)); } else { - return tryCreateExtendedPartitionsAsync(previousPartitions, expectPartitions) - .thenCompose(ignore -> - updatePartitionedTopic(topicName, previousPartitions, expectPartitions, force)); + return tryCreatePartitionsAsync(expectPartitions) + .thenCompose(ignore -> updatePartitionedTopic(topicName, expectPartitions)); } }); } @@ -4364,8 +4360,7 @@ private PersistentReplicator getReplicatorReference(String replName, PersistentT } } - private CompletableFuture updatePartitionedTopic(TopicName topicName, int previousPartitions, - int expectPartitions, boolean force) { + private CompletableFuture updatePartitionedTopic(TopicName topicName, int expectPartitions) { CompletableFuture future = namespaceResources().getPartitionedTopicResources() .updatePartitionedTopicAsync(topicName, p -> new PartitionedTopicMetadata(expectPartitions, p.properties)); @@ -4387,30 +4382,21 @@ private CompletableFuture updatePartitionedTopic(TopicName topicName, int }); return null; }); - return future.thenCompose(__ -> createSubscriptions(topicName, previousPartitions, expectPartitions, force)); + return future.thenCompose(__ -> createSubscriptions(topicName, expectPartitions)); } /** * It creates subscriptions for new partitions of existing partitioned-topics. * * @param topicName : topic-name: persistent://prop/cluster/ns/topic - * @param previousPartitions number of previous partitions * @param expectPartitions : number of expected partitions - * @param ignoreConflictException : If true, ignore ConflictException: subscription already exists for topic * */ - private CompletableFuture createSubscriptions(TopicName topicName, int previousPartitions, - int expectPartitions, boolean ignoreConflictException) { + private CompletableFuture createSubscriptions(TopicName topicName, int expectPartitions) { CompletableFuture result = new CompletableFuture<>(); - if (previousPartitions < 1) { + if (expectPartitions < 1) { return FutureUtil.failedFuture(new RestException(Status.CONFLICT, "Topic is not partitioned topic")); } - - if (previousPartitions >= expectPartitions) { - return FutureUtil.failedFuture(new RestException(Status.CONFLICT, - "number of partitions must be more than existing " + previousPartitions)); - } - PulsarAdmin admin; try { admin = pulsar().getAdminClient(); @@ -4430,7 +4416,7 @@ private CompletableFuture createSubscriptions(TopicName topicName, int pre } boolean replicated = ss.isReplicated(); - for (int i = previousPartitions; i < expectPartitions; i++) { + for (int i = 0; i < expectPartitions; i++) { final String topicNamePartition = topicName.getPartition(i).toString(); CompletableFuture future = new CompletableFuture<>(); admin.topics().createSubscriptionAsync(topicNamePartition, @@ -4438,8 +4424,7 @@ private CompletableFuture createSubscriptions(TopicName topicName, int pre if (ex == null) { future.complete(null); } else { - if (ignoreConflictException - && ex instanceof PulsarAdminException.ConflictException) { + if (ex instanceof PulsarAdminException.ConflictException) { future.complete(null); } else { future.completeExceptionally(ex); From 58f516f67bc3c26591462a701b5582d3e8ea06e3 Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Fri, 6 Jan 2023 20:27:38 +0800 Subject: [PATCH 19/20] Fix test --- .../admin/impl/PersistentTopicsBase.java | 83 ++++++++++++------- 1 file changed, 52 insertions(+), 31 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index b5cd83ccf8219..f9891dd054b40 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -439,36 +439,56 @@ protected CompletableFuture internalUpdatePartitionedTopicAsync(int expect throw new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be less than or equal to " + maxPartitions); } - // Only do the validation if it's the first hop. - if (topicName.isGlobal() && isNamespaceReplicated(topicName.getNamespaceObject())) { - return getNamespaceReplicatedClustersAsync(topicName.getNamespaceObject()) - .thenApply(clusters -> { - if (!clusters.contains(pulsar().getConfig().getClusterName())) { - log.error("[{}] local cluster is not part of replicated cluster for namespace {}", - clientAppId(), topicName); - throw new RestException(Status.FORBIDDEN, "Local cluster is not part of replicate" - + " cluster list"); - } - return clusters; - }) - .thenCompose(clusters -> - tryCreatePartitionsAsync(expectPartitions).thenApply(ignore -> clusters)) - .thenCompose(clusters -> { - if (!updateLocalTopicOnly) { - return namespaceResources().getPartitionedTopicResources() - .updatePartitionedTopicAsync(topicName, p -> - new PartitionedTopicMetadata(expectPartitions, p.properties)) - .thenCompose(__ -> - updatePartitionInOtherCluster(expectPartitions, clusters)); - } else { - return CompletableFuture.completedFuture(null); - } - }).thenCompose(clusters -> createSubscriptions(topicName, - expectPartitions)); - } else { - return tryCreatePartitionsAsync(expectPartitions) - .thenCompose(ignore -> updatePartitionedTopic(topicName, expectPartitions)); + final PulsarAdmin adminClient; + try { + adminClient = pulsar().getAdminClient(); + } catch (PulsarServerException e) { + throw new RuntimeException(e); } + return adminClient.topics().getListAsync(topicName.getNamespace()) + .thenCompose(topics -> { + long existPartitions = topics.stream() + .filter(t -> t.startsWith(topicName.getPartitionedTopicName())) + .count(); + if (existPartitions >= expectPartitions) { + throw new RestException(Status.CONFLICT, + "Number of new partitions must be greater than existing number of partitions"); + } + // Only do the validation if it's the first hop. + if (topicName.isGlobal() && isNamespaceReplicated(topicName.getNamespaceObject())) { + return getNamespaceReplicatedClustersAsync(topicName.getNamespaceObject()) + .thenApply(clusters -> { + if (!clusters.contains(pulsar().getConfig().getClusterName())) { + log.error("[{}] local cluster is not part of" + + " replicated cluster for namespace {}", + clientAppId(), topicName); + throw new RestException(Status.FORBIDDEN, + "Local cluster is not part of replicate cluster list"); + } + return clusters; + }) + .thenCompose(clusters -> + tryCreatePartitionsAsync(expectPartitions) + .thenApply(ignore -> clusters)) + .thenCompose(clusters -> { + if (!updateLocalTopicOnly) { + return namespaceResources().getPartitionedTopicResources() + .updatePartitionedTopicAsync(topicName, p -> + new PartitionedTopicMetadata(expectPartitions, + p.properties)) + .thenCompose(__ -> + updatePartitionInOtherCluster(expectPartitions, + clusters)); + } else { + return CompletableFuture.completedFuture(null); + } + }).thenCompose(clusters -> createSubscriptions(topicName, + expectPartitions)); + } else { + return tryCreatePartitionsAsync(expectPartitions) + .thenCompose(ignore -> updatePartitionedTopic(topicName, expectPartitions)); + } + }); }); } @@ -4424,10 +4444,11 @@ private CompletableFuture createSubscriptions(TopicName topicName, int exp if (ex == null) { future.complete(null); } else { - if (ex instanceof PulsarAdminException.ConflictException) { + Throwable realCause = FutureUtil.unwrapCompletionException(ex); + if (realCause instanceof PulsarAdminException.ConflictException) { future.complete(null); } else { - future.completeExceptionally(ex); + future.completeExceptionally(realCause); } } }); From 5e2c937deab065e9350c8b71e1daf174922d12aa Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Mon, 9 Jan 2023 16:00:06 +0800 Subject: [PATCH 20/20] Apply comment --- .../apache/pulsar/broker/admin/impl/PersistentTopicsBase.java | 3 ++- .../java/org/apache/pulsar/broker/admin/AdminApi2Test.java | 1 + 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index f9891dd054b40..81c9638632eb7 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java @@ -448,7 +448,8 @@ protected CompletableFuture internalUpdatePartitionedTopicAsync(int expect return adminClient.topics().getListAsync(topicName.getNamespace()) .thenCompose(topics -> { long existPartitions = topics.stream() - .filter(t -> t.startsWith(topicName.getPartitionedTopicName())) + .filter(t -> TopicName.get(t).getPartitionedTopicName() + .equals(topicName.getPartitionedTopicName())) .count(); if (existPartitions >= expectPartitions) { throw new RestException(Status.CONFLICT, diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java index ca9bb6950b8c1..e35e9311b9fd2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApi2Test.java @@ -2679,6 +2679,7 @@ public void testFailedUpdatePartitionedTopic() throws Exception { try { admin.topics().createSubscription(partitionedTopicName + "-partition-" + startPartitions, subName1, MessageId.earliest); + fail("Unexpected behaviour"); } catch (PulsarAdminException.PreconditionFailedException ex) { // OK }