From a08e3c9d062ba50c5572c66c3f03e9ba899f9895 Mon Sep 17 00:00:00 2001 From: fxbing <94fxiaobing@gmail.com> Date: Thu, 15 Aug 2019 21:05:57 +0800 Subject: [PATCH 01/23] Add more config for auto-topic-creation --- conf/broker.conf | 6 ++ .../pulsar/broker/ServiceConfiguration.java | 13 ++- .../pulsar/broker/admin/AdminResource.java | 26 ++++++ .../BrokerServiceAutoTopicCreationTest.java | 83 +++++++++++++++++++ site2/docs/reference-configuration.md | 2 + 5 files changed, 129 insertions(+), 1 deletion(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java diff --git a/conf/broker.conf b/conf/broker.conf index b491336ca534c..5282850b34669 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -88,6 +88,12 @@ ttlDurationDefaultInSeconds=0 # Enable topic auto creation if new producer or consumer connected (disable auto creation with value false) allowAutoTopicCreation=true +# The type of topic that is allowed to be automatically created.(partition/non-partition) +allowAutoTopicCreationType=non-partition + +# The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partition. +allowAutoTopicCreationNumPartitions=1 + # Enable the deletion of inactive topics brokerDeleteInactiveTopicsEnabled=true diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 832f02759c681..b43bc794fe6ca 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -879,9 +879,20 @@ public class ServiceConfiguration implements PulsarConfiguration { private double managedLedgerDefaultMarkDeleteRateLimit = 1.0; @FieldContext( category = CATEGORY_STORAGE_ML, - doc = "Allow automated creation of non-partition topics if set to true (default value)." + doc = "Allow automated creation of topics if set to true (default value)." ) private boolean allowAutoTopicCreation = true; + @FieldContext( + category = CATEGORY_STORAGE_ML, + doc = "The type of topic that is allowed to be automatically created.(partition/non-partition)" + ) + private String allowAutoTopicCreationType = "non-partition"; + @FieldContext( + category = CATEGORY_STORAGE_ML, + doc = "The number of partitioned topics that is allowed to be automatically created" + + "if allowAutoTopicCreationType is partition." + ) + private int allowAutoTopicCreationNumPartitions = 1; @FieldContext( category = CATEGORY_STORAGE_ML, doc = "Number of threads to be used for managed ledger tasks dispatching" diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index 94409e1d071f6..06bdba59fc8b3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -18,7 +18,9 @@ */ package org.apache.pulsar.broker.admin; +import com.google.common.base.Joiner; import static com.google.common.base.Preconditions.checkArgument; +import org.apache.pulsar.broker.PulsarServerException; import static org.apache.pulsar.broker.cache.ConfigurationCacheService.POLICIES; import static org.apache.pulsar.common.util.Codec.decode; @@ -551,8 +553,26 @@ public PartitionedTopicMetadata deserialize(String key, byte[] content) throws E } }).thenAccept(metadata -> { // if the partitioned topic is not found in zk, then the topic is not partitioned + boolean allowAutoTopicCreation = pulsar.getConfiguration().isAllowAutoTopicCreation(); + String topicType = pulsar.getConfiguration().getAllowAutoTopicCreationType(); if (metadata.isPresent()) { metadataFuture.complete(metadata.get()); + } else if (allowAutoTopicCreation && "partition".equals(topicType)) { + String topicName = getTopicNameFromPath(path); + int configPartitions = pulsar.getConfiguration().getAllowAutoTopicCreationNumPartitions(); + try { + pulsar.getAdminClient().topics() + .createPartitionedTopicAsync(topicName, configPartitions) + .whenComplete((result, ex) -> { + if (ex == null) { + metadataFuture.complete(new PartitionedTopicMetadata(configPartitions)); + } else { + metadataFuture.completeExceptionally(ex); + } + }); + } catch (PulsarServerException e) { + metadataFuture.completeExceptionally(e); + } } else { metadataFuture.complete(new PartitionedTopicMetadata()); } @@ -627,4 +647,10 @@ protected List getPartitionedTopicList(TopicDomain topicDomain) { partitionedTopics.sort(null); return partitionedTopics; } + + private static String getTopicNameFromPath(String path) { + String formatPath = splitPath(path, 4); + String[] parts = formatPath.split("/"); + return Joiner.on("/").join(parts[2] + ":/", parts[0], parts[1], parts[3]); + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java new file mode 100644 index 0000000000000..29fc8570a09c6 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java @@ -0,0 +1,83 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.broker.service; + +import org.apache.pulsar.client.api.PulsarClientException; +import static org.testng.Assert.assertFalse; +import static org.testng.Assert.assertTrue; +import org.testng.annotations.AfterClass; +import org.testng.annotations.BeforeClass; +import org.testng.annotations.Test; + +public class BrokerServiceAutoTopicCreationTest extends BrokerTestBase{ + @BeforeClass + @Override + protected void setup() throws Exception { + super.baseSetup(); + } + + @AfterClass + @Override + protected void cleanup() throws Exception { + super.internalCleanup(); + } + + @Test + public void testAutoNonPartitionedTopicCreation() throws Exception{ + pulsar.getConfiguration().setAllowAutoTopicCreation(true); + pulsar.getConfiguration().setAllowAutoTopicCreationType("non-partition"); + + final String topicName = "persistent://prop/ns-abc/non-partitioned-topic"; + final String subscriptionName = "non-partitioned-topic-sub"; + pulsarClient.newConsumer().topic(topicName).subscriptionName(subscriptionName).subscribe(); + + assertTrue(admin.namespaces().getTopics("prop/ns-abc").contains(topicName)); + assertFalse(admin.topics().getPartitionedTopicList("prop/ns-abc").contains(topicName)); + } + + @Test + public void testAutoPartitionedTopicCreation() throws Exception{ + pulsar.getConfiguration().setAllowAutoTopicCreation(true); + pulsar.getConfiguration().setAllowAutoTopicCreationType("partition"); + pulsar.getConfiguration().setAllowAutoTopicCreationNumPartitions(3); + + final String topicName = "persistent://prop/ns-abc/partitioned-topic"; + final String subscriptionName = "partitioned-topic-sub"; + pulsarClient.newConsumer().topic(topicName).subscriptionName(subscriptionName).subscribe(); + + assertTrue(admin.topics().getPartitionedTopicList("prop/ns-abc").contains(topicName)); + for (int i = 0; i < 3; i++) { + assertTrue(admin.namespaces().getTopics("prop/ns-abc").contains(topicName + "-partition-" + i)); + } + } + + @Test + public void testAutoTopicCreationDisable() throws Exception{ + pulsar.getConfiguration().setAllowAutoTopicCreation(false); + + final String topicName = "persistent://prop/ns-abc/test-topic"; + final String subscriptionName = "test-topic-sub"; + try { + pulsarClient.newConsumer().topic(topicName).subscriptionName(subscriptionName).subscribe(); + } catch (Exception e) { + assertTrue(e instanceof PulsarClientException); + } + assertFalse(admin.namespaces().getTopics("prop/ns-abc").contains(topicName)); + } +} diff --git a/site2/docs/reference-configuration.md b/site2/docs/reference-configuration.md index bb549ae167d44..489057724606e 100644 --- a/site2/docs/reference-configuration.md +++ b/site2/docs/reference-configuration.md @@ -127,6 +127,8 @@ Pulsar brokers are responsible for handling incoming messages from producers, di |backlogQuotaCheckIntervalInSeconds| How often to check for topics that have reached the quota |60| |backlogQuotaDefaultLimitGB| Default per-topic backlog quota limit |10| |allowAutoTopicCreation| Enable topic auto creation if new producer or consumer connected |true| +|allowAutoTopicCreationType| The type of topic that is allowed to be automatically created.(partition/non-partition) |non-partition| +|allowAutoTopicCreationNumPartitions| The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partition |1| |brokerDeleteInactiveTopicsEnabled| Enable the deletion of inactive topics |true| |brokerDeleteInactiveTopicsFrequencySeconds| How often to check for inactive topics |60| |messageExpiryCheckIntervalInMinutes| How frequently to proactively check and purge expired messages |5| From 1dc9a9e1cc3aa039a408bc3d083742c0fa62aa74 Mon Sep 17 00:00:00 2001 From: fxbing <94fxiaobing@gmail.com> Date: Fri, 16 Aug 2019 19:52:21 +0800 Subject: [PATCH 02/23] Modify the location of the changes so that the changes do not affect other features. --- conf/broker.conf | 4 +- .../pulsar/broker/ServiceConfiguration.java | 4 +- .../pulsar/broker/admin/AdminResource.java | 92 +++++++++++++------ .../admin/impl/PersistentTopicsBase.java | 37 ++++---- .../broker/admin/v1/NonPersistentTopics.java | 2 +- .../broker/admin/v2/NonPersistentTopics.java | 2 +- .../BrokerServiceAutoTopicCreationTest.java | 4 +- site2/docs/reference-configuration.md | 2 +- 8 files changed, 94 insertions(+), 53 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index 5282850b34669..f43965d561f9e 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -88,8 +88,8 @@ ttlDurationDefaultInSeconds=0 # Enable topic auto creation if new producer or consumer connected (disable auto creation with value false) allowAutoTopicCreation=true -# The type of topic that is allowed to be automatically created.(partition/non-partition) -allowAutoTopicCreationType=non-partition +# The type of topic that is allowed to be automatically created.(partitioned/non-partitioned) +allowAutoTopicCreationType=non-partitioned # The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partition. allowAutoTopicCreationNumPartitions=1 diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index b43bc794fe6ca..3499b297b0021 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -884,9 +884,9 @@ public class ServiceConfiguration implements PulsarConfiguration { private boolean allowAutoTopicCreation = true; @FieldContext( category = CATEGORY_STORAGE_ML, - doc = "The type of topic that is allowed to be automatically created.(partition/non-partition)" + doc = "The type of topic that is allowed to be automatically created.(partitioned/non-partitioned)" ) - private String allowAutoTopicCreationType = "non-partition"; + private String allowAutoTopicCreationType = "non-partitioned"; @FieldContext( category = CATEGORY_STORAGE_ML, doc = "The number of partitioned topics that is allowed to be automatically created" diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index 06bdba59fc8b3..09f70bf8bec7c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -18,9 +18,8 @@ */ package org.apache.pulsar.broker.admin; -import com.google.common.base.Joiner; +import com.fasterxml.jackson.core.JsonProcessingException; import static com.google.common.base.Preconditions.checkArgument; -import org.apache.pulsar.broker.PulsarServerException; import static org.apache.pulsar.broker.cache.ConfigurationCacheService.POLICIES; import static org.apache.pulsar.common.util.Codec.decode; @@ -83,6 +82,7 @@ public abstract class AdminResource extends PulsarWebResource { private static final Logger log = LoggerFactory.getLogger(AdminResource.class); private static final String POLICIES_READONLY_FLAG_PATH = "/admin/flags/policies-readonly"; + private static final int PARTITIONED_TOPIC_WAIT_SYNC_TIME_MS = 1000; public static final String PARTITIONED_TOPIC_PATH_ZNODE = "partitioned-topics"; protected ZooKeeper globalZk() { @@ -502,7 +502,7 @@ protected ZooKeeperChildrenCache failureDomainListCache() { } protected PartitionedTopicMetadata getPartitionedTopicMetadata(TopicName topicName, - boolean authoritative) { + boolean authoritative, boolean checkAllowAutoCreation) { validateClusterOwnership(topicName.getCluster()); // validates global-namespace contains local/peer cluster: if peer/local cluster present then lookup can // serve/redirect request else fail partitioned-metadata-request so, client fails while creating @@ -521,7 +521,12 @@ protected PartitionedTopicMetadata getPartitionedTopicMetadata(TopicName topicNa } String path = path(PARTITIONED_TOPIC_PATH_ZNODE, namespaceName.toString(), domain(), topicName.getEncodedLocalName()); - PartitionedTopicMetadata partitionMetadata = fetchPartitionedTopicMetadata(pulsar(), path); + PartitionedTopicMetadata partitionMetadata; + if (checkAllowAutoCreation) { + partitionMetadata = fetchPartitionedTopicMetadataCheckAllowAutoCreation(pulsar(), path); + } else { + partitionMetadata = fetchPartitionedTopicMetadata(pulsar(), path); + } if (log.isDebugEnabled()) { log.debug("[{}] Total number of partitions for topic {} is {}", clientAppId(), topicName, @@ -553,26 +558,8 @@ public PartitionedTopicMetadata deserialize(String key, byte[] content) throws E } }).thenAccept(metadata -> { // if the partitioned topic is not found in zk, then the topic is not partitioned - boolean allowAutoTopicCreation = pulsar.getConfiguration().isAllowAutoTopicCreation(); - String topicType = pulsar.getConfiguration().getAllowAutoTopicCreationType(); if (metadata.isPresent()) { metadataFuture.complete(metadata.get()); - } else if (allowAutoTopicCreation && "partition".equals(topicType)) { - String topicName = getTopicNameFromPath(path); - int configPartitions = pulsar.getConfiguration().getAllowAutoTopicCreationNumPartitions(); - try { - pulsar.getAdminClient().topics() - .createPartitionedTopicAsync(topicName, configPartitions) - .whenComplete((result, ex) -> { - if (ex == null) { - metadataFuture.complete(new PartitionedTopicMetadata(configPartitions)); - } else { - metadataFuture.completeExceptionally(ex); - } - }); - } catch (PulsarServerException e) { - metadataFuture.completeExceptionally(e); - } } else { metadataFuture.complete(new PartitionedTopicMetadata()); } @@ -586,6 +573,51 @@ public PartitionedTopicMetadata deserialize(String key, byte[] content) throws E return metadataFuture; } + protected static PartitionedTopicMetadata fetchPartitionedTopicMetadataCheckAllowAutoCreation( + PulsarService pulsar, String path) { + try { + return fetchPartitionedTopicMetadataCheckAllowAutoCreationAsync(pulsar, path).get(); + } catch (Exception e) { + if (e.getCause() instanceof RestException) { + throw (RestException) e; + } + throw new RestException(e); + } + } + + protected static CompletableFuture fetchPartitionedTopicMetadataCheckAllowAutoCreationAsync( + PulsarService pulsar, String path) { + CompletableFuture metadataFuture = new CompletableFuture<>(); + try { + boolean allowAutoTopicCreation = pulsar.getConfiguration().isAllowAutoTopicCreation(); + String topicType = pulsar.getConfiguration().getAllowAutoTopicCreationType(); + fetchPartitionedTopicMetadataAsync(pulsar, path).whenComplete((metadata, ex) -> { + if (ex != null) { + metadataFuture.completeExceptionally(ex.getCause()); + } else if (metadata.partitions == 0 && allowAutoTopicCreation && + TopicType.PARTITIONED.toString().equals(topicType)) { + int configPartitions = pulsar.getConfiguration().getAllowAutoTopicCreationNumPartitions(); + try { + PartitionedTopicMetadata configMetadata = new PartitionedTopicMetadata(configPartitions); + byte[] content = jsonMapper().writeValueAsBytes(configMetadata); + ZkUtils.createFullPathOptimistic(pulsar.getGlobalZkCache().getZooKeeper(), path, content, + ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + // we wait for the data to be synced in all quorums and the observers + Thread.sleep(PARTITIONED_TOPIC_WAIT_SYNC_TIME_MS); + metadataFuture.complete(configMetadata); + } catch (JsonProcessingException | InterruptedException | KeeperException e) { + metadataFuture.completeExceptionally(e.getCause()); + } + } else { + metadataFuture.complete(metadata); + } + }); + } catch (Exception e) { + metadataFuture.completeExceptionally(e); + } + return metadataFuture; + } + protected void validateClusterExists(String cluster) { try { if (!clustersCache().get(path("clusters", cluster)).isPresent()) { @@ -648,9 +680,17 @@ protected List getPartitionedTopicList(TopicDomain topicDomain) { return partitionedTopics; } - private static String getTopicNameFromPath(String path) { - String formatPath = splitPath(path, 4); - String[] parts = formatPath.split("/"); - return Joiner.on("/").join(parts[2] + ":/", parts[0], parts[1], parts[3]); + enum TopicType { + PARTITIONED("partitioned"), + NON_PARTITIONED("non-partitioned"); + private String type; + + TopicType(String type) { + this.type = type; + } + + public String toString() { + return type; + } } } 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 9150f9965a882..f341fc45a8eac 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 @@ -445,7 +445,7 @@ protected void internalUpdatePartitionedTopic(int numPartitions) { } protected PartitionedTopicMetadata internalGetPartitionedMetadata(boolean authoritative) { - PartitionedTopicMetadata metadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata metadata = getPartitionedTopicMetadata(topicName, authoritative, true); if (metadata.partitions > 1) { validateClientVersion(); } @@ -457,7 +457,7 @@ protected void internalDeletePartitionedTopic(AsyncResponse asyncResponse, boole final CompletableFuture future = new CompletableFuture<>(); - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); final int numPartitions = partitionMetadata.partitions; if (numPartitions > 0) { final AtomicInteger count = new AtomicInteger(numPartitions); @@ -590,7 +590,7 @@ protected void internalGetSubscriptions(AsyncResponse asyncResponse, boolean aut final List subscriptions = Lists.newArrayList(); - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions > 0) { try { // get the subscriptions only from the 1st partition since all the other partitions will have the same @@ -685,7 +685,7 @@ public void getInfoFailed(ManagedLedgerException exception, Object ctx) { protected void internalGetPartitionedStats(AsyncResponse asyncResponse, boolean authoritative, boolean perPartition) { - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions == 0) { throw new RestException(Status.NOT_FOUND, "Partitioned Topic not found"); } @@ -743,7 +743,7 @@ protected void internalGetPartitionedStats(AsyncResponse asyncResponse, boolean } protected void internalGetPartitionedStatsInternal(AsyncResponse asyncResponse, boolean authoritative) { - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions == 0) { throw new RestException(Status.NOT_FOUND, "Partitioned Topic not found"); } @@ -786,7 +786,7 @@ protected void internalDeleteSubscription(AsyncResponse asyncResponse, String su if (topicName.isGlobal()) { validateGlobalNamespaceOwnership(namespaceName); } - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions > 0) { final List> futures = Lists.newArrayList(); @@ -855,7 +855,7 @@ protected void internalSkipAllMessages(AsyncResponse asyncResponse, String subNa if (topicName.isGlobal()) { validateGlobalNamespaceOwnership(namespaceName); } - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions > 0) { final List> futures = Lists.newArrayList(); @@ -920,7 +920,7 @@ protected void internalSkipMessages(String subName, int numMessages, boolean aut if (topicName.isGlobal()) { validateGlobalNamespaceOwnership(namespaceName); } - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions > 0) { throw new RestException(Status.METHOD_NOT_ALLOWED, "Skip messages on a partitioned topic is not allowed"); } @@ -952,7 +952,7 @@ protected void internalExpireMessagesForAllSubscriptions(AsyncResponse asyncResp if (topicName.isGlobal()) { validateGlobalNamespaceOwnership(namespaceName); } - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions > 0) { final List> futures = Lists.newArrayList(); @@ -1027,7 +1027,7 @@ protected void internalResetCursor(AsyncResponse asyncResponse, String subName, validateGlobalNamespaceOwnership(namespaceName); } - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); final int numPartitions = partitionMetadata.partitions; if (numPartitions > 0) { final CompletableFuture future = new CompletableFuture<>(); @@ -1141,7 +1141,7 @@ protected void internalCreateSubscription(AsyncResponse asyncResponse, String su log.info("[{}][{}] Creating subscription {} at message id {}", clientAppId(), topicName, subscriptionName, targetMessageId); - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); final int numPartitions = partitionMetadata.partitions; if (numPartitions > 0) { final CompletableFuture future = new CompletableFuture<>(); @@ -1249,7 +1249,7 @@ protected void internalResetCursorOnPosition(String subName, boolean authoritati log.info("[{}][{}] received reset cursor on subscription {} to position {}", clientAppId(), topicName, subName, messageId); - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions > 0) { log.warn("[{}] Not supported operation on partitioned-topic {} {}", clientAppId(), topicName, @@ -1288,7 +1288,7 @@ protected Response internalPeekNthMessage(String subName, int messagePosition, b if (topicName.isGlobal()) { validateGlobalNamespaceOwnership(namespaceName); } - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions > 0) { throw new RestException(Status.METHOD_NOT_ALLOWED, "Peek messages on a partitioned topic is not allowed"); } @@ -1413,7 +1413,7 @@ protected MessageId internalTerminate(boolean authoritative) { if (topicName.isGlobal()) { validateGlobalNamespaceOwnership(namespaceName); } - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions > 0) { throw new RestException(Status.METHOD_NOT_ALLOWED, "Termination of a partitioned topic is not allowed"); } @@ -1433,7 +1433,7 @@ protected void internalExpireMessages(AsyncResponse asyncResponse, String subNam validateGlobalNamespaceOwnership(namespaceName); } - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions > 0) { final List> futures = Lists.newArrayList(); @@ -1489,7 +1489,7 @@ private void internalExpireMessagesForSinglePartition(String subName, int expire validateGlobalNamespaceOwnership(namespaceName); } - PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative); + PartitionedTopicMetadata partitionMetadata = getPartitionedTopicMetadata(topicName, authoritative, false); if (partitionMetadata.partitions > 0) { String msg = "This method should not be called for partitioned topic"; log.error("[{}] {} {} {}", clientAppId(), msg, topicName, subName); @@ -1595,7 +1595,8 @@ public static CompletableFuture getPartitionedTopicMet // serve/redirect request else fail partitioned-metadata-request so, client fails while creating // producer/consumer checkLocalOrGetPeerReplicationCluster(pulsar, topicName.getNamespaceObject()) - .thenCompose(res -> fetchPartitionedTopicMetadataAsync(pulsar, path)).thenAccept(metadata -> { + .thenCompose(res -> fetchPartitionedTopicMetadataCheckAllowAutoCreationAsync(pulsar, path)) + .thenAccept(metadata -> { if (log.isDebugEnabled()) { log.debug("[{}] Total number of partitions for topic {} is {}", clientAppId, topicName, metadata.partitions); @@ -1632,7 +1633,7 @@ private RestException topicNotFoundReason(TopicName topicName) { } PartitionedTopicMetadata partitionedTopicMetadata = getPartitionedTopicMetadata( - TopicName.get(topicName.getPartitionedTopicName()), false); + TopicName.get(topicName.getPartitionedTopicName()), false, false); if (partitionedTopicMetadata == null || partitionedTopicMetadata.partitions == 0) { final String topicErrorType = partitionedTopicMetadata == null ? "has no metadata" : "has zero partitions"; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java index f1347e342ed4f..adca1bf21cae4 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java @@ -78,7 +78,7 @@ public PartitionedTopicMetadata getPartitionedMetadata(@PathParam("property") St @PathParam("topic") @Encoded String encodedTopic, @QueryParam("authoritative") @DefaultValue("false") boolean authoritative) { validateTopicName(property, cluster, namespace, encodedTopic); - return getPartitionedTopicMetadata(topicName, authoritative); + return getPartitionedTopicMetadata(topicName, authoritative, false); } @GET diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java index 8125f5b31b4fa..92c12b6fd0bfc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java @@ -88,7 +88,7 @@ public PartitionedTopicMetadata getPartitionedMetadata( @ApiParam(value = "Is authentication required to perform this operation") @QueryParam("authoritative") @DefaultValue("false") boolean authoritative) { validateTopicName(tenant, namespace, encodedTopic); - return getPartitionedTopicMetadata(topicName, authoritative); + return getPartitionedTopicMetadata(topicName, authoritative, false); } @GET diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java index 29fc8570a09c6..224be57d6c16f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java @@ -41,7 +41,7 @@ protected void cleanup() throws Exception { @Test public void testAutoNonPartitionedTopicCreation() throws Exception{ pulsar.getConfiguration().setAllowAutoTopicCreation(true); - pulsar.getConfiguration().setAllowAutoTopicCreationType("non-partition"); + pulsar.getConfiguration().setAllowAutoTopicCreationType("non-partitioned"); final String topicName = "persistent://prop/ns-abc/non-partitioned-topic"; final String subscriptionName = "non-partitioned-topic-sub"; @@ -54,7 +54,7 @@ public void testAutoNonPartitionedTopicCreation() throws Exception{ @Test public void testAutoPartitionedTopicCreation() throws Exception{ pulsar.getConfiguration().setAllowAutoTopicCreation(true); - pulsar.getConfiguration().setAllowAutoTopicCreationType("partition"); + pulsar.getConfiguration().setAllowAutoTopicCreationType("partitioned"); pulsar.getConfiguration().setAllowAutoTopicCreationNumPartitions(3); final String topicName = "persistent://prop/ns-abc/partitioned-topic"; diff --git a/site2/docs/reference-configuration.md b/site2/docs/reference-configuration.md index 489057724606e..aa4caf52eceda 100644 --- a/site2/docs/reference-configuration.md +++ b/site2/docs/reference-configuration.md @@ -127,7 +127,7 @@ Pulsar brokers are responsible for handling incoming messages from producers, di |backlogQuotaCheckIntervalInSeconds| How often to check for topics that have reached the quota |60| |backlogQuotaDefaultLimitGB| Default per-topic backlog quota limit |10| |allowAutoTopicCreation| Enable topic auto creation if new producer or consumer connected |true| -|allowAutoTopicCreationType| The type of topic that is allowed to be automatically created.(partition/non-partition) |non-partition| +|allowAutoTopicCreationType| The type of topic that is allowed to be automatically created.(partitioned/non-partitioned) |non-partitioned| |allowAutoTopicCreationNumPartitions| The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partition |1| |brokerDeleteInactiveTopicsEnabled| Enable the deletion of inactive topics |true| |brokerDeleteInactiveTopicsFrequencySeconds| How often to check for inactive topics |60| From 82dd8401578c71bb4a2e5d24913353f4ad34fb9f Mon Sep 17 00:00:00 2001 From: fxbing <94fxiaobing@gmail.com> Date: Sat, 17 Aug 2019 23:14:46 +0800 Subject: [PATCH 03/23] fix test bug about exception --- .../java/org/apache/pulsar/broker/admin/AdminResource.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index 09f70bf8bec7c..750408293d23c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -593,7 +593,7 @@ protected static CompletableFuture fetchPartitionedTop String topicType = pulsar.getConfiguration().getAllowAutoTopicCreationType(); fetchPartitionedTopicMetadataAsync(pulsar, path).whenComplete((metadata, ex) -> { if (ex != null) { - metadataFuture.completeExceptionally(ex.getCause()); + metadataFuture.completeExceptionally(ex); } else if (metadata.partitions == 0 && allowAutoTopicCreation && TopicType.PARTITIONED.toString().equals(topicType)) { int configPartitions = pulsar.getConfiguration().getAllowAutoTopicCreationNumPartitions(); @@ -606,7 +606,7 @@ protected static CompletableFuture fetchPartitionedTop Thread.sleep(PARTITIONED_TOPIC_WAIT_SYNC_TIME_MS); metadataFuture.complete(configMetadata); } catch (JsonProcessingException | InterruptedException | KeeperException e) { - metadataFuture.completeExceptionally(e.getCause()); + metadataFuture.completeExceptionally(e); } } else { metadataFuture.complete(metadata); From 55d4d692d0fa78acb925298b290178dff1a59ef8 Mon Sep 17 00:00:00 2001 From: fxbing <94fxiaobing@gmail.com> Date: Tue, 20 Aug 2019 13:48:42 +0800 Subject: [PATCH 04/23] Modify configuration parameter name and default value --- conf/broker.conf | 6 +++--- .../java/org/apache/pulsar/broker/ServiceConfiguration.java | 6 +++--- .../java/org/apache/pulsar/broker/admin/AdminResource.java | 2 +- .../broker/service/BrokerServiceAutoTopicCreationTest.java | 2 +- site2/docs/reference-configuration.md | 2 +- 5 files changed, 9 insertions(+), 9 deletions(-) diff --git a/conf/broker.conf b/conf/broker.conf index f43965d561f9e..dcd39f3e2e1e6 100644 --- a/conf/broker.conf +++ b/conf/broker.conf @@ -89,10 +89,10 @@ ttlDurationDefaultInSeconds=0 allowAutoTopicCreation=true # The type of topic that is allowed to be automatically created.(partitioned/non-partitioned) -allowAutoTopicCreationType=non-partitioned +allowAutoTopicCreationType=partitioned -# The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partition. -allowAutoTopicCreationNumPartitions=1 +# The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partitioned. +defaultNumPartitions=1 # Enable the deletion of inactive topics brokerDeleteInactiveTopicsEnabled=true diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index 3499b297b0021..4d7ef4133db1a 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -886,13 +886,13 @@ public class ServiceConfiguration implements PulsarConfiguration { category = CATEGORY_STORAGE_ML, doc = "The type of topic that is allowed to be automatically created.(partitioned/non-partitioned)" ) - private String allowAutoTopicCreationType = "non-partitioned"; + private String allowAutoTopicCreationType = "partitioned"; @FieldContext( category = CATEGORY_STORAGE_ML, doc = "The number of partitioned topics that is allowed to be automatically created" - + "if allowAutoTopicCreationType is partition." + + "if allowAutoTopicCreationType is partitioned." ) - private int allowAutoTopicCreationNumPartitions = 1; + private int defaultNumPartitions = 1; @FieldContext( category = CATEGORY_STORAGE_ML, doc = "Number of threads to be used for managed ledger tasks dispatching" diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index 750408293d23c..d6a91b4e7c439 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -596,7 +596,7 @@ protected static CompletableFuture fetchPartitionedTop metadataFuture.completeExceptionally(ex); } else if (metadata.partitions == 0 && allowAutoTopicCreation && TopicType.PARTITIONED.toString().equals(topicType)) { - int configPartitions = pulsar.getConfiguration().getAllowAutoTopicCreationNumPartitions(); + int configPartitions = pulsar.getConfiguration().getDefaultNumPartitions(); try { PartitionedTopicMetadata configMetadata = new PartitionedTopicMetadata(configPartitions); byte[] content = jsonMapper().writeValueAsBytes(configMetadata); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java index 224be57d6c16f..ecd9ee7bd8e26 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java @@ -55,7 +55,7 @@ public void testAutoNonPartitionedTopicCreation() throws Exception{ public void testAutoPartitionedTopicCreation() throws Exception{ pulsar.getConfiguration().setAllowAutoTopicCreation(true); pulsar.getConfiguration().setAllowAutoTopicCreationType("partitioned"); - pulsar.getConfiguration().setAllowAutoTopicCreationNumPartitions(3); + pulsar.getConfiguration().setDefaultNumPartitions(3); final String topicName = "persistent://prop/ns-abc/partitioned-topic"; final String subscriptionName = "partitioned-topic-sub"; diff --git a/site2/docs/reference-configuration.md b/site2/docs/reference-configuration.md index aa4caf52eceda..5261bae0c2651 100644 --- a/site2/docs/reference-configuration.md +++ b/site2/docs/reference-configuration.md @@ -128,7 +128,7 @@ Pulsar brokers are responsible for handling incoming messages from producers, di |backlogQuotaDefaultLimitGB| Default per-topic backlog quota limit |10| |allowAutoTopicCreation| Enable topic auto creation if new producer or consumer connected |true| |allowAutoTopicCreationType| The type of topic that is allowed to be automatically created.(partitioned/non-partitioned) |non-partitioned| -|allowAutoTopicCreationNumPartitions| The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partition |1| +|defaultNumPartitions| The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partitioned |1| |brokerDeleteInactiveTopicsEnabled| Enable the deletion of inactive topics |true| |brokerDeleteInactiveTopicsFrequencySeconds| How often to check for inactive topics |60| |messageExpiryCheckIntervalInMinutes| How frequently to proactively check and purge expired messages |5| From 42d6ea5cdddcbea41aa9249775bf3c1d7e505171 Mon Sep 17 00:00:00 2001 From: fxbing <94fxiaobing@gmail.com> Date: Tue, 20 Aug 2019 15:25:02 +0800 Subject: [PATCH 05/23] If topic is already exist, creating partitioned topic is not allowed. --- .../pulsar/broker/admin/AdminResource.java | 23 +++++++++++++++---- .../admin/impl/PersistentTopicsBase.java | 2 +- .../BrokerServiceAutoTopicCreationTest.java | 18 +++++++++++++++ 3 files changed, 37 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index d6a91b4e7c439..ebf3f23f3fecf 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -21,6 +21,7 @@ import com.fasterxml.jackson.core.JsonProcessingException; import static com.google.common.base.Preconditions.checkArgument; import static org.apache.pulsar.broker.cache.ConfigurationCacheService.POLICIES; +import org.apache.pulsar.common.api.proto.PulsarApi; import static org.apache.pulsar.common.util.Codec.decode; import java.net.MalformedURLException; @@ -523,7 +524,7 @@ protected PartitionedTopicMetadata getPartitionedTopicMetadata(TopicName topicNa String path = path(PARTITIONED_TOPIC_PATH_ZNODE, namespaceName.toString(), domain(), topicName.getEncodedLocalName()); PartitionedTopicMetadata partitionMetadata; if (checkAllowAutoCreation) { - partitionMetadata = fetchPartitionedTopicMetadataCheckAllowAutoCreation(pulsar(), path); + partitionMetadata = fetchPartitionedTopicMetadataCheckAllowAutoCreation(pulsar(), path, topicName); } else { partitionMetadata = fetchPartitionedTopicMetadata(pulsar(), path); } @@ -574,9 +575,10 @@ public PartitionedTopicMetadata deserialize(String key, byte[] content) throws E } protected static PartitionedTopicMetadata fetchPartitionedTopicMetadataCheckAllowAutoCreation( - PulsarService pulsar, String path) { + PulsarService pulsar, String path, TopicName topicName) { try { - return fetchPartitionedTopicMetadataCheckAllowAutoCreationAsync(pulsar, path).get(); + return fetchPartitionedTopicMetadataCheckAllowAutoCreationAsync(pulsar, path, topicName) + .get(); } catch (Exception e) { if (e.getCause() instanceof RestException) { throw (RestException) e; @@ -586,15 +588,26 @@ protected static PartitionedTopicMetadata fetchPartitionedTopicMetadataCheckAllo } protected static CompletableFuture fetchPartitionedTopicMetadataCheckAllowAutoCreationAsync( - PulsarService pulsar, String path) { + PulsarService pulsar, String path, TopicName topicName) { CompletableFuture metadataFuture = new CompletableFuture<>(); try { boolean allowAutoTopicCreation = pulsar.getConfiguration().isAllowAutoTopicCreation(); String topicType = pulsar.getConfiguration().getAllowAutoTopicCreationType(); + boolean topicExist; + try { + topicExist = pulsar.getNamespaceService() + .getListOfTopics(topicName.getNamespaceObject(), PulsarApi.CommandGetTopicsOfNamespace.Mode.ALL) + .contains(topicName.toString()); + } catch (Exception e) { + log.warn("Unexpected error while getting list of topics. topic={}. Error: {}", + topicName, e.getMessage(), e); + throw new RestException(e); + } fetchPartitionedTopicMetadataAsync(pulsar, path).whenComplete((metadata, ex) -> { if (ex != null) { metadataFuture.completeExceptionally(ex); - } else if (metadata.partitions == 0 && allowAutoTopicCreation && + // If topic is already exist, creating partitioned topic is not allowed. + } else if (metadata.partitions == 0 && !topicExist && allowAutoTopicCreation && TopicType.PARTITIONED.toString().equals(topicType)) { int configPartitions = pulsar.getConfiguration().getDefaultNumPartitions(); try { 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 f341fc45a8eac..cec6cc1f2590c 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 @@ -1595,7 +1595,7 @@ public static CompletableFuture getPartitionedTopicMet // serve/redirect request else fail partitioned-metadata-request so, client fails while creating // producer/consumer checkLocalOrGetPeerReplicationCluster(pulsar, topicName.getNamespaceObject()) - .thenCompose(res -> fetchPartitionedTopicMetadataCheckAllowAutoCreationAsync(pulsar, path)) + .thenCompose(res -> fetchPartitionedTopicMetadataCheckAllowAutoCreationAsync(pulsar, path, topicName)) .thenAccept(metadata -> { if (log.isDebugEnabled()) { log.debug("[{}] Total number of partitions for topic {} is {}", clientAppId, topicName, diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java index ecd9ee7bd8e26..17e6bfb81a606 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java @@ -80,4 +80,22 @@ public void testAutoTopicCreationDisable() throws Exception{ } assertFalse(admin.namespaces().getTopics("prop/ns-abc").contains(topicName)); } + + @Test + public void testAutoTopicCreationDisableIfNonPartitionedTopicAlreadyExist() throws Exception{ + pulsar.getConfiguration().setAllowAutoTopicCreation(true); + pulsar.getConfiguration().setAllowAutoTopicCreationType("partitioned"); + pulsar.getConfiguration().setDefaultNumPartitions(3); + + final String topicName = "persistent://prop/ns-abc/partitioned-topic"; + final String subscriptionName = "partitioned-topic-sub"; + admin.topics().createNonPartitionedTopic(topicName); + pulsarClient.newConsumer().topic(topicName).subscriptionName(subscriptionName).subscribe(); + + assertFalse(admin.topics().getPartitionedTopicList("prop/ns-abc").contains(topicName)); + for (int i = 0; i < 3; i++) { + assertFalse(admin.namespaces().getTopics("prop/ns-abc").contains(topicName + "-partition-" + i)); + } + assertTrue(admin.namespaces().getTopics("prop/ns-abc").contains(topicName)); + } } From f05cdebd3f5914e88b981ec557cadc617e39435e Mon Sep 17 00:00:00 2001 From: fxbing Date: Wed, 21 Aug 2019 14:41:33 +0800 Subject: [PATCH 06/23] If topic is already exist, creating partitioned topic is not allowed. --- conf/standalone.conf | 9 +++++++++ .../pulsar/broker/admin/AdminResource.java | 5 +++-- .../admin/impl/PersistentTopicsBase.java | 17 +++++++++++++++-- .../broker/admin/v1/NonPersistentTopics.java | 18 ++++++++++++++++-- .../broker/admin/v1/PersistentTopics.java | 5 +++-- .../broker/admin/v2/NonPersistentTopics.java | 19 +++++++++++++++++-- .../broker/admin/v2/PersistentTopics.java | 6 ++++-- .../broker/admin/PersistentTopicsTest.java | 6 +++--- .../configurations/pulsar_broker_test.conf | 2 ++ .../apache/pulsar/client/impl/HttpClient.java | 1 + 10 files changed, 73 insertions(+), 15 deletions(-) diff --git a/conf/standalone.conf b/conf/standalone.conf index 388c185851fa8..b25f21fdb643f 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -592,3 +592,12 @@ allowLoopback=true # will heart performance. It is better to give a higher number of gc # interval if there is enough disk capacity. gcWaitTime=300000 + +# Enable topic auto creation if new producer or consumer connected (disable auto creation with value false) +allowAutoTopicCreation=true + +# The type of topic that is allowed to be automatically created.(partitioned/non-partitioned) +allowAutoTopicCreationType=non-partitioned + +# The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partitioned. +defaultNumPartitions=1 \ No newline at end of file diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index ebf3f23f3fecf..58873c050fcef 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -609,9 +609,10 @@ protected static CompletableFuture fetchPartitionedTop // If topic is already exist, creating partitioned topic is not allowed. } else if (metadata.partitions == 0 && !topicExist && allowAutoTopicCreation && TopicType.PARTITIONED.toString().equals(topicType)) { - int configPartitions = pulsar.getConfiguration().getDefaultNumPartitions(); + int defaultNumPartitions = pulsar.getConfiguration().getDefaultNumPartitions(); + checkArgument(defaultNumPartitions > 0, "Default number of partitions should be more than 0"); try { - PartitionedTopicMetadata configMetadata = new PartitionedTopicMetadata(configPartitions); + PartitionedTopicMetadata configMetadata = new PartitionedTopicMetadata(defaultNumPartitions); byte[] content = jsonMapper().writeValueAsBytes(configMetadata); ZkUtils.createFullPathOptimistic(pulsar.getGlobalZkCache().getZooKeeper(), path, content, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); 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 cec6cc1f2590c..509b41de5c36b 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 @@ -20,6 +20,7 @@ import static com.google.common.base.Preconditions.checkNotNull; import static org.apache.pulsar.broker.cache.ConfigurationCacheService.POLICIES; +import org.apache.pulsar.common.api.proto.PulsarApi; import static org.apache.pulsar.common.util.Codec.decode; import com.github.zafarkhaja.semver.Version; @@ -374,6 +375,18 @@ protected void internalCreatePartitionedTopic(int numPartitions) { if (numPartitions <= 0) { throw new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be more than 0"); } + try { + boolean topicExist = pulsar().getNamespaceService() + .getListOfTopics(topicName.getNamespaceObject(), PulsarApi.CommandGetTopicsOfNamespace.Mode.ALL) + .contains(topicName.toString()); + if (topicExist) { + log.warn("[{}] Failed to create already existing topic {}", clientAppId(), topicName); + throw new RestException(Status.CONFLICT, "This topic already exists"); + } + } catch (Exception e) { + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); + throw new RestException(e); + } try { String path = ZkAdminPaths.partitionedTopicPath(topicName); byte[] data = jsonMapper().writeValueAsBytes(new PartitionedTopicMetadata(numPartitions)); @@ -444,8 +457,8 @@ protected void internalUpdatePartitionedTopic(int numPartitions) { } } - protected PartitionedTopicMetadata internalGetPartitionedMetadata(boolean authoritative) { - PartitionedTopicMetadata metadata = getPartitionedTopicMetadata(topicName, authoritative, true); + protected PartitionedTopicMetadata internalGetPartitionedMetadata(boolean authoritative, boolean checkAllowAutoCreation) { + PartitionedTopicMetadata metadata = getPartitionedTopicMetadata(topicName, authoritative, checkAllowAutoCreation); if (metadata.partitions > 1) { validateClientVersion(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java index adca1bf21cae4..96bc0841e5881 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/NonPersistentTopics.java @@ -47,6 +47,7 @@ import org.apache.pulsar.broker.service.Topic; import org.apache.pulsar.broker.service.nonpersistent.NonPersistentTopic; import org.apache.pulsar.broker.web.RestException; +import org.apache.pulsar.common.api.proto.PulsarApi; import org.apache.pulsar.common.naming.Constants; import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.NamespaceName; @@ -76,9 +77,10 @@ public class NonPersistentTopics extends PersistentTopics { public PartitionedTopicMetadata getPartitionedMetadata(@PathParam("property") String property, @PathParam("cluster") String cluster, @PathParam("namespace") String namespace, @PathParam("topic") @Encoded String encodedTopic, - @QueryParam("authoritative") @DefaultValue("false") boolean authoritative) { + @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, + @QueryParam("checkAllowAutoCreation") @DefaultValue("false") boolean checkAllowAutoCreation) { validateTopicName(property, cluster, namespace, encodedTopic); - return getPartitionedTopicMetadata(topicName, authoritative, false); + return getPartitionedTopicMetadata(topicName, authoritative, checkAllowAutoCreation); } @GET @@ -124,6 +126,18 @@ public void createPartitionedTopic(@PathParam("property") String property, @Path if (numPartitions <= 0) { throw new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be more than 0"); } + try { + boolean topicExist = pulsar().getNamespaceService() + .getListOfTopics(topicName.getNamespaceObject(), PulsarApi.CommandGetTopicsOfNamespace.Mode.ALL) + .contains(topicName.toString()); + if (topicExist) { + log.warn("[{}] Failed to create already existing topic {}", clientAppId(), topicName); + throw new RestException(Status.CONFLICT, "This topic already exists"); + } + } catch (Exception e) { + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); + throw new RestException(e); + } try { String path = path(PARTITIONED_TOPIC_PATH_ZNODE, namespaceName.toString(), domain(), topicName.getEncodedLocalName()); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java index f895ce184ff3f..ebece5e15455b 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/PersistentTopics.java @@ -181,9 +181,10 @@ public void updatePartitionedTopic(@PathParam("property") String property, @Path public PartitionedTopicMetadata getPartitionedMetadata(@PathParam("property") String property, @PathParam("cluster") String cluster, @PathParam("namespace") String namespace, @PathParam("topic") @Encoded String encodedTopic, - @QueryParam("authoritative") @DefaultValue("false") boolean authoritative) { + @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, + @QueryParam("checkAllowAutoCreation") @DefaultValue("false") boolean checkAllowAutoCreation) { validateTopicName(property, cluster, namespace, encodedTopic); - return internalGetPartitionedMetadata(authoritative); + return internalGetPartitionedMetadata(authoritative, checkAllowAutoCreation); } @DELETE diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java index 92c12b6fd0bfc..10dc5ee8d0485 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/NonPersistentTopics.java @@ -48,6 +48,7 @@ import org.apache.pulsar.broker.service.Topic; import org.apache.pulsar.broker.service.nonpersistent.NonPersistentTopic; import org.apache.pulsar.broker.web.RestException; +import org.apache.pulsar.common.api.proto.PulsarApi; import org.apache.pulsar.common.naming.NamespaceBundle; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.partition.PartitionedTopicMetadata; @@ -86,9 +87,11 @@ public PartitionedTopicMetadata getPartitionedMetadata( @ApiParam(value = "Specify topic name", required = true) @PathParam("topic") @Encoded String encodedTopic, @ApiParam(value = "Is authentication required to perform this operation") - @QueryParam("authoritative") @DefaultValue("false") boolean authoritative) { + @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, + @ApiParam(value = "Is check configuration required to automatically create topic") + @QueryParam("checkAllowAutoCreation") @DefaultValue("false") boolean checkAllowAutoCreation) { validateTopicName(tenant, namespace, encodedTopic); - return getPartitionedTopicMetadata(topicName, authoritative, false); + return getPartitionedTopicMetadata(topicName, authoritative, checkAllowAutoCreation); } @GET @@ -168,6 +171,18 @@ public void createPartitionedTopic( if (numPartitions <= 0) { throw new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be more than 0"); } + try { + boolean topicExist = pulsar().getNamespaceService() + .getListOfTopics(topicName.getNamespaceObject(), PulsarApi.CommandGetTopicsOfNamespace.Mode.ALL) + .contains(topicName.toString()); + if (topicExist) { + log.warn("[{}] Failed to create already existing topic {}", clientAppId(), topicName); + throw new RestException(Status.CONFLICT, "This topic already exists"); + } + } catch (Exception e) { + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); + throw new RestException(e); + } try { String path = path(PARTITIONED_TOPIC_PATH_ZNODE, namespaceName.toString(), domain(), topicName.getEncodedLocalName()); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java index e09c0121600b5..70e0624b12e02 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java @@ -288,9 +288,11 @@ public PartitionedTopicMetadata getPartitionedMetadata( @ApiParam(value = "Specify topic name", required = true) @PathParam("topic") @Encoded String encodedTopic, @ApiParam(value = "Is authentication required to perform this operation") - @QueryParam("authoritative") @DefaultValue("false") boolean authoritative) { + @QueryParam("authoritative") @DefaultValue("false") boolean authoritative, + @ApiParam(value = "Is check configuration required to automatically create topic") + @QueryParam("checkAllowAutoCreation") @DefaultValue("false") boolean checkAllowAutoCreation) { validateTopicName(tenant, namespace, encodedTopic); - return internalGetPartitionedMetadata(authoritative); + return internalGetPartitionedMetadata(authoritative, checkAllowAutoCreation); } @DELETE diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java index fc89acc5022bd..3d8a29a6b5feb 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/PersistentTopicsTest.java @@ -211,8 +211,8 @@ public void testNonPartitionedTopics() { final String nonPartitionTopic2 = "secondary-non-partitioned-topic"; persistentTopics.createNonPartitionedTopic(testTenant, testNamespace, nonPartitionTopic2, true); - Assert.assertEquals( - persistentTopics.getPartitionedMetadata(testTenant, testNamespace, nonPartitionTopic, true) .partitions, + Assert.assertEquals(persistentTopics + .getPartitionedMetadata(testTenant, testNamespace, nonPartitionTopic, true, false).partitions, 0); } @@ -221,7 +221,7 @@ public void testCreateNonPartitionedTopic() { final String topicName = "standard-topic"; persistentTopics.createNonPartitionedTopic(testTenant, testNamespace, topicName, true); PartitionedTopicMetadata pMetadata = persistentTopics.getPartitionedMetadata( - testTenant, testNamespace, topicName, true); + testTenant, testNamespace, topicName, true, false); Assert.assertEquals(pMetadata.partitions, 0); } diff --git a/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf b/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf index c55955fd34108..fe25a8ee87635 100644 --- a/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf +++ b/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf @@ -33,6 +33,8 @@ backlogQuotaCheckIntervalInSeconds=60 backlogQuotaDefaultLimitGB=50 brokerDeleteInactiveTopicsEnabled=true brokerDeleteInactiveTopicsFrequencySeconds=60 +allowAutoTopicCreation=true +allowAutoTopicCreationType=non-partitioned messageExpiryCheckIntervalInMinutes=5 clientLibraryVersionCheckEnabled=false clientLibraryVersionCheckAllowUnversioned=true diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpClient.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpClient.java index 96c62d2e9b6ef..a4211ff2b1e82 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpClient.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpClient.java @@ -152,6 +152,7 @@ public CompletableFuture get(String path, Class clazz) { // auth complete, use a new builder BoundRequestBuilder builder = httpClient.prepareGet(requestUrl) .setHeader("Accept", "application/json"); + builder.addQueryParam("checkAllowAutoCreation", "true"); if (authData.hasDataForHttp()) { Set> headers; From 476a4031e336be483c5f1f481055e8c5aa23b1e4 Mon Sep 17 00:00:00 2001 From: fxbing Date: Wed, 21 Aug 2019 19:39:02 +0800 Subject: [PATCH 07/23] add default config to original test and fix flaky test --- .../pulsar/broker/auth/MockedPulsarServiceBaseTest.java | 1 + .../pulsar/broker/service/BacklogQuotaManagerTest.java | 1 + .../pulsar/broker/service/BrokerBkEnsemblesTests.java | 1 + .../pulsar/broker/service/BrokerBookieIsolationTest.java | 4 ++++ .../service/BrokerServiceAutoTopicCreationTest.java | 2 +- .../apache/pulsar/broker/service/ReplicatorTestBase.java | 3 +++ .../apache/pulsar/client/api/NonPersistentTopicTest.java | 3 +++ .../pulsar/client/api/SimpleProducerConsumerTest.java | 4 ++-- .../pulsar/client/api/v1/V1_ProducerConsumerTest.java | 4 ++-- .../functions/worker/PulsarFunctionE2ESecurityTest.java | 1 + .../functions/worker/PulsarFunctionLocalRunTest.java | 1 + .../functions/worker/PulsarFunctionPublishTest.java | 1 + .../pulsar/functions/worker/PulsarFunctionStateTest.java | 1 + .../java/org/apache/pulsar/io/PulsarFunctionE2ETest.java | 1 + .../resources/configurations/pulsar_broker_test.conf | 1 + pulsar-client-cpp/tests/standalone.conf | 9 +++++++++ 16 files changed, 33 insertions(+), 5 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java index 341b96a69af8d..868d5916b7255 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/auth/MockedPulsarServiceBaseTest.java @@ -104,6 +104,7 @@ protected void resetConfig() { this.conf.setDefaultNumberOfNamespaceBundles(1); this.conf.setZookeeperServers("localhost:2181"); this.conf.setConfigurationStoreServers("localhost:3181"); + this.conf.setAllowAutoTopicCreationType("non-persistent"); } protected final void internalSetup() throws Exception { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java index c4a387cf2a62e..88a66036a528c 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BacklogQuotaManagerTest.java @@ -88,6 +88,7 @@ void setup() throws Exception { config.setBacklogQuotaCheckIntervalInSeconds(TIME_TO_CHECK_BACKLOG_QUOTA); config.setManagedLedgerMaxEntriesPerLedger(5); config.setManagedLedgerMinLedgerRolloverTimeMinutes(0); + config.setAllowAutoTopicCreationType("non-partitioned"); pulsar = new PulsarService(config); pulsar.start(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java index 940002d13b8e1..f8807349f6e91 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBkEnsemblesTests.java @@ -104,6 +104,7 @@ protected void setup() throws Exception { config.setManagedLedgerMaxEntriesPerLedger(5); config.setManagedLedgerMinLedgerRolloverTimeMinutes(0); config.setAdvertisedAddress("127.0.0.1"); + config.setAllowAutoTopicCreationType("non-partitioned"); pulsar = new PulsarService(config); pulsar.start(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBookieIsolationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBookieIsolationTest.java index 3b9edc9088eb7..68f72b911203a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBookieIsolationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBookieIsolationTest.java @@ -158,6 +158,8 @@ public void testBookieIsolation() throws Exception { config.setManagedLedgerDefaultWriteQuorum(2); config.setManagedLedgerDefaultAckQuorum(2); + config.setAllowAutoTopicCreationType("non-partitioned"); + int totalEntriesPerLedger = 20; int totalLedgers = totalPublish / totalEntriesPerLedger; config.setManagedLedgerMaxEntriesPerLedger(totalEntriesPerLedger); @@ -288,6 +290,7 @@ public void testBookieIsilationWithSecondaryGroup() throws Exception { config.setManagedLedgerDefaultEnsembleSize(2); config.setManagedLedgerDefaultWriteQuorum(2); config.setManagedLedgerDefaultAckQuorum(2); + config.setAllowAutoTopicCreationType("non-partitioned"); int totalEntriesPerLedger = 20; int totalLedgers = totalPublish / totalEntriesPerLedger; @@ -410,6 +413,7 @@ public void testDeleteIsolationGroup() throws Exception { config.setManagedLedgerDefaultEnsembleSize(2); config.setManagedLedgerDefaultWriteQuorum(2); config.setManagedLedgerDefaultAckQuorum(2); + config.setAllowAutoTopicCreationType("non-partitioned"); config.setManagedLedgerMinLedgerRolloverTimeMinutes(0); pulsarService = new PulsarService(config); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java index 17e6bfb81a606..4b62ce91f281a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java @@ -87,7 +87,7 @@ public void testAutoTopicCreationDisableIfNonPartitionedTopicAlreadyExist() thro pulsar.getConfiguration().setAllowAutoTopicCreationType("partitioned"); pulsar.getConfiguration().setDefaultNumPartitions(3); - final String topicName = "persistent://prop/ns-abc/partitioned-topic"; + final String topicName = "persistent://prop/ns-abc/test-topic-2"; final String subscriptionName = "partitioned-topic-sub"; admin.topics().createNonPartitionedTopic(topicName); pulsarClient.newConsumer().topic(topicName).subscriptionName(subscriptionName).subscribe(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTestBase.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTestBase.java index ca5ba6864cecf..5864f48e789e2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTestBase.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ReplicatorTestBase.java @@ -132,6 +132,7 @@ void setup() throws Exception { config1.setTlsTrustCertsFilePath(TLS_SERVER_CERT_FILE_PATH); config1.setBacklogQuotaCheckIntervalInSeconds(TIME_TO_CHECK_BACKLOG_QUOTA); config1.setDefaultNumberOfNamespaceBundles(1); + config1.setAllowAutoTopicCreationType("non-partitioned"); pulsar1 = new PulsarService(config1); pulsar1.start(); ns1 = pulsar1.getBrokerService(); @@ -165,6 +166,7 @@ void setup() throws Exception { config2.setTlsTrustCertsFilePath(TLS_SERVER_CERT_FILE_PATH); config2.setBacklogQuotaCheckIntervalInSeconds(TIME_TO_CHECK_BACKLOG_QUOTA); config2.setDefaultNumberOfNamespaceBundles(1); + config2.setAllowAutoTopicCreationType("non-partitioned"); pulsar2 = new PulsarService(config2); pulsar2.start(); ns2 = pulsar2.getBrokerService(); @@ -198,6 +200,7 @@ void setup() throws Exception { config3.setTlsKeyFilePath(TLS_SERVER_KEY_FILE_PATH); config3.setTlsTrustCertsFilePath(TLS_SERVER_CERT_FILE_PATH); config3.setDefaultNumberOfNamespaceBundles(1); + config3.setAllowAutoTopicCreationType("non-partitioned"); pulsar3 = new PulsarService(config3); pulsar3.start(); ns3 = pulsar3.getBrokerService(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/NonPersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/NonPersistentTopicTest.java index f70e6a1150e45..b54265a615753 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/NonPersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/NonPersistentTopicTest.java @@ -891,6 +891,7 @@ void setupReplicationCluster() throws Exception { inSec(getBrokerServicePurgeInactiveFrequency(), TimeUnit.SECONDS)); config1.setBrokerServicePort(Optional.ofNullable(PortManager.nextFreePort())); config1.setBacklogQuotaCheckIntervalInSeconds(TIME_TO_CHECK_BACKLOG_QUOTA); + config1.setAllowAutoTopicCreationType("non-partitioned"); pulsar1 = new PulsarService(config1); pulsar1.start(); ns1 = pulsar1.getBrokerService(); @@ -917,6 +918,7 @@ void setupReplicationCluster() throws Exception { inSec(getBrokerServicePurgeInactiveFrequency(), TimeUnit.SECONDS)); config2.setBrokerServicePort(Optional.ofNullable(PortManager.nextFreePort())); config2.setBacklogQuotaCheckIntervalInSeconds(TIME_TO_CHECK_BACKLOG_QUOTA); + config2.setAllowAutoTopicCreationType("non-partitioned"); pulsar2 = new PulsarService(config2); pulsar2.start(); ns2 = pulsar2.getBrokerService(); @@ -942,6 +944,7 @@ void setupReplicationCluster() throws Exception { config3.setBrokerServicePurgeInactiveFrequencyInSeconds( inSec(getBrokerServicePurgeInactiveFrequency(), TimeUnit.SECONDS)); config3.setBrokerServicePort(Optional.ofNullable(PortManager.nextFreePort())); + config3.setAllowAutoTopicCreationType("non-partitioned"); pulsar3 = new PulsarService(config3); pulsar3.start(); ns3 = pulsar3.getBrokerService(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java index 35d4c77d82bbb..c2fad94f0b3a2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java @@ -2309,7 +2309,7 @@ public void testFailReceiveAsyncOnConsumerClose() throws Exception { // (1) simple consumers Consumer consumer = pulsarClient.newConsumer() - .topic("persistent://my-property/my-ns/failAsyncReceive").subscriptionName("my-subscriber-name") + .topic("persistent://my-property/my-ns/failAsyncReceive-1").subscriptionName("my-subscriber-name") .subscribe(); consumer.close(); // receive messages @@ -2322,7 +2322,7 @@ public void testFailReceiveAsyncOnConsumerClose() throws Exception { // (2) Partitioned-consumer int numPartitions = 4; - TopicName topicName = TopicName.get("persistent://my-property/my-ns/failAsyncReceive"); + TopicName topicName = TopicName.get("persistent://my-property/my-ns/failAsyncReceive-2"); admin.topics().createPartitionedTopic(topicName.toString(), numPartitions); Consumer partitionedConsumer = pulsarClient.newConsumer().topic(topicName.toString()) .subscriptionName("my-partitioned-subscriber").subscribe(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java index 24fb5a7562942..7d9a8cdc66b79 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/v1/V1_ProducerConsumerTest.java @@ -2040,7 +2040,7 @@ public void testFailReceiveAsyncOnConsumerClose() throws Exception { // (1) simple consumers Consumer consumer = pulsarClient.newConsumer() - .topic("persistent://my-property/use/my-ns/failAsyncReceive") + .topic("persistent://my-property/use/my-ns/failAsyncReceive-1") .subscriptionName("my-subscriber-name") .subscribe(); consumer.close(); @@ -2054,7 +2054,7 @@ public void testFailReceiveAsyncOnConsumerClose() throws Exception { // (2) Partitioned-consumer int numPartitions = 4; - TopicName topicName = TopicName.get("persistent://my-property/use/my-ns/failAsyncReceive"); + TopicName topicName = TopicName.get("persistent://my-property/use/my-ns/failAsyncReceive-2"); admin.topics().createPartitionedTopic(topicName.toString(), numPartitions); Consumer partitionedConsumer = pulsarClient.newConsumer().topic(topicName.toString()) .subscriptionName("my-partitioned-subscriber") diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java index 54f05d8125ff6..afe79f493e09d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java @@ -132,6 +132,7 @@ void setup(Method method) throws Exception { config.setBrokerServicePort(Optional.of(brokerServicePort)); config.setLoadManagerClassName(SimpleLoadManagerImpl.class.getName()); config.setAdvertisedAddress("localhost"); + config.setAllowAutoTopicCreationType("non-partitioned"); Set providers = new HashSet<>(); providers.add(AuthenticationProviderToken.class.getName()); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java index 97cabc335c2b6..80047953acf5e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java @@ -180,6 +180,7 @@ void setup(Method method) throws Exception { "tlsCertFile:" + TLS_CLIENT_CERT_FILE_PATH + "," + "tlsKeyFile:" + TLS_CLIENT_KEY_FILE_PATH); config.setBrokerClientTrustCertsFilePath(TLS_TRUST_CERT_FILE_PATH); config.setBrokerClientTlsEnabled(true); + config.setAllowAutoTopicCreationType("non-partitioned"); functionsWorkerService = createPulsarFunctionWorker(config); urlTls = new URL(brokerServiceUrl); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java index 9c2b7b7688c5b..0a0ea6257b9d5 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java @@ -169,6 +169,7 @@ void setup(Method method) throws Exception { "tlsCertFile:" + TLS_CLIENT_CERT_FILE_PATH + "," + "tlsKeyFile:" + TLS_CLIENT_KEY_FILE_PATH); config.setBrokerClientTrustCertsFilePath(TLS_TRUST_CERT_FILE_PATH); config.setBrokerClientTlsEnabled(true); + config.setAllowAutoTopicCreationType("non-partitioned"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionStateTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionStateTest.java index 809fca79c4634..a7de6c9962ffc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionStateTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionStateTest.java @@ -168,6 +168,7 @@ void setup(Method method) throws Exception { "tlsCertFile:" + TLS_CLIENT_CERT_FILE_PATH + "," + "tlsKeyFile:" + TLS_CLIENT_KEY_FILE_PATH); config.setBrokerClientTrustCertsFilePath(TLS_TRUST_CERT_FILE_PATH); config.setBrokerClientTlsEnabled(true); + config.setAllowAutoTopicCreationType("non-partitioned"); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java index 847958fc223b0..addd2c6928a5a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java @@ -194,6 +194,7 @@ void setup(Method method) throws Exception { "tlsCertFile:" + TLS_CLIENT_CERT_FILE_PATH + "," + "tlsKeyFile:" + TLS_CLIENT_KEY_FILE_PATH); config.setBrokerClientTrustCertsFilePath(TLS_TRUST_CERT_FILE_PATH); config.setBrokerClientTlsEnabled(true); + config.setAllowAutoTopicCreationType("non-partitioned"); diff --git a/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf b/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf index fe25a8ee87635..75b6d03311092 100644 --- a/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf +++ b/pulsar-broker/src/test/resources/configurations/pulsar_broker_test.conf @@ -35,6 +35,7 @@ brokerDeleteInactiveTopicsEnabled=true brokerDeleteInactiveTopicsFrequencySeconds=60 allowAutoTopicCreation=true allowAutoTopicCreationType=non-partitioned +defaultNumPartitions=1 messageExpiryCheckIntervalInMinutes=5 clientLibraryVersionCheckEnabled=false clientLibraryVersionCheckAllowUnversioned=true diff --git a/pulsar-client-cpp/tests/standalone.conf b/pulsar-client-cpp/tests/standalone.conf index 8a016426d96da..857285ab3d570 100644 --- a/pulsar-client-cpp/tests/standalone.conf +++ b/pulsar-client-cpp/tests/standalone.conf @@ -266,3 +266,12 @@ keepAliveIntervalSeconds=30 # How often broker checks for inactive topics to be deleted (topics with no subscriptions and no one connected) brokerServicePurgeInactiveFrequencyInSeconds=60 + +# Enable topic auto creation if new producer or consumer connected (disable auto creation with value false) +allowAutoTopicCreation=true + +# The type of topic that is allowed to be automatically created.(partitioned/non-partitioned) +allowAutoTopicCreationType=non-partitioned + +# The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partitioned. +defaultNumPartitions=1 From 300c35985e62c57825a9d6bf8ff6d99af0674c0e Mon Sep 17 00:00:00 2001 From: fxbing Date: Thu, 22 Aug 2019 09:58:56 +0800 Subject: [PATCH 08/23] add config to cpp test --- pulsar-client-cpp/test-conf/standalone.conf | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/pulsar-client-cpp/test-conf/standalone.conf b/pulsar-client-cpp/test-conf/standalone.conf index 6f799c171cc27..2de6a376af5bb 100644 --- a/pulsar-client-cpp/test-conf/standalone.conf +++ b/pulsar-client-cpp/test-conf/standalone.conf @@ -263,3 +263,12 @@ keepAliveIntervalSeconds=30 # How often broker checks for inactive topics to be deleted (topics with no subscriptions and no one connected) brokerServicePurgeInactiveFrequencyInSeconds=60 + +# Enable topic auto creation if new producer or consumer connected (disable auto creation with value false) +allowAutoTopicCreation=true + +# The type of topic that is allowed to be automatically created.(partitioned/non-partitioned) +allowAutoTopicCreationType=non-partitioned + +# The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partitioned. +defaultNumPartitions=1 \ No newline at end of file From 0fe2d3aa6ac0f8c6b908590f59522b7b0bfd0d08 Mon Sep 17 00:00:00 2001 From: fxbing Date: Thu, 22 Aug 2019 11:31:38 +0800 Subject: [PATCH 09/23] add config to cpp test --- pulsar-client-cpp/test-conf/standalone-ssl.conf | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/pulsar-client-cpp/test-conf/standalone-ssl.conf b/pulsar-client-cpp/test-conf/standalone-ssl.conf index 426ba437df8d0..6ab440683782c 100644 --- a/pulsar-client-cpp/test-conf/standalone-ssl.conf +++ b/pulsar-client-cpp/test-conf/standalone-ssl.conf @@ -278,3 +278,12 @@ keepAliveIntervalSeconds=30 # How often broker checks for inactive topics to be deleted (topics with no subscriptions and no one connected) brokerServicePurgeInactiveFrequencyInSeconds=60 + +# Enable topic auto creation if new producer or consumer connected (disable auto creation with value false) +allowAutoTopicCreation=true + +# The type of topic that is allowed to be automatically created.(partitioned/non-partitioned) +allowAutoTopicCreationType=non-partitioned + +# The number of partitioned topics that is allowed to be automatically created if allowAutoTopicCreationType is partitioned. +defaultNumPartitions=1 \ No newline at end of file From 39f9ee1841093b5ae06aa57b95fc7f94915f096c Mon Sep 17 00:00:00 2001 From: fxbing Date: Thu, 22 Aug 2019 17:50:57 +0800 Subject: [PATCH 10/23] trigger test --- .../broker/service/BrokerServiceAutoTopicCreationTest.java | 1 + 1 file changed, 1 insertion(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java index 4b62ce91f281a..cf08c0e4ea892 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java @@ -26,6 +26,7 @@ import org.testng.annotations.Test; public class BrokerServiceAutoTopicCreationTest extends BrokerTestBase{ + @BeforeClass @Override protected void setup() throws Exception { From f08c9a766d7b92704d3e71c4a1774fee679ec302 Mon Sep 17 00:00:00 2001 From: fxbing Date: Tue, 27 Aug 2019 20:30:19 +0800 Subject: [PATCH 11/23] fix integration tests --- .../org/apache/pulsar/broker/admin/impl/BrokersBase.java | 1 + .../org/apache/pulsar/tests/integration/cli/CLITest.java | 6 ++++++ 2 files changed, 7 insertions(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/BrokersBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/BrokersBase.java index 63e392a48c7a6..d07b519a99e6f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/BrokersBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/BrokersBase.java @@ -262,6 +262,7 @@ public void healthcheck(@Suspended AsyncResponse asyncResponse) throws Exception PulsarClient client = pulsar().getClient(); String messageStr = UUID.randomUUID().toString(); + pulsar().getAdminClient().topics().createNonPartitionedTopic(topic); CompletableFuture> producerFuture = client.newProducer(Schema.STRING).topic(topic).createAsync(); CompletableFuture> readerFuture = client.newReader(Schema.STRING) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/CLITest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/CLITest.java index 3af34d7f60dea..ca39cd46ab1db 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/CLITest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/CLITest.java @@ -93,6 +93,12 @@ public void testCreateSubscriptionCommand() throws Exception { public void testTopicTerminationOnTopicsWithoutConnectedConsumers() throws Exception { String topicName = "persistent://public/default/test-topic-termination"; BrokerContainer container = pulsarCluster.getAnyBroker(); + container.execCmd( + PulsarCluster.ADMIN_SCRIPT, + "topics", + "create", + topicName); + ContainerExecResult result = container.execCmd( PulsarCluster.CLIENT_SCRIPT, "produce", From 356e99b9c800cbcdf89799c5e62cf9e44b3d7c3e Mon Sep 17 00:00:00 2001 From: fxbing Date: Wed, 28 Aug 2019 15:12:35 +0800 Subject: [PATCH 12/23] fix integration tests --- .../org/apache/pulsar/broker/admin/impl/BrokersBase.java | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/BrokersBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/BrokersBase.java index d07b519a99e6f..fc529e9e575d0 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/BrokersBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/BrokersBase.java @@ -262,7 +262,13 @@ public void healthcheck(@Suspended AsyncResponse asyncResponse) throws Exception PulsarClient client = pulsar().getClient(); String messageStr = UUID.randomUUID().toString(); - pulsar().getAdminClient().topics().createNonPartitionedTopic(topic); + // create non-partitioned topic manually + try { + pulsar().getBrokerService().getTopic(topic, true).get(); + } catch (Exception e) { + asyncResponse.resume(new RestException(e)); + return; + } CompletableFuture> producerFuture = client.newProducer(Schema.STRING).topic(topic).createAsync(); CompletableFuture> readerFuture = client.newReader(Schema.STRING) From f47d5a95ad428368c4bbbbcc9bfa023ba2ee4695 Mon Sep 17 00:00:00 2001 From: fxbing Date: Sun, 1 Sep 2019 17:04:18 +0800 Subject: [PATCH 13/23] create non-partitioned topic manually for function --- .../java/org/apache/pulsar/functions/worker/WorkerService.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerService.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerService.java index 9ec9688b1a0ac..0d739156b02a4 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerService.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerService.java @@ -148,6 +148,9 @@ public void start(URI dlogUri, } log.info("Created Pulsar client"); + brokerAdmin.topics().createNonPartitionedTopic(workerConfig.getFunctionAssignmentTopic()); + brokerAdmin.topics().createNonPartitionedTopic(workerConfig.getClusterCoordinationTopic()); + brokerAdmin.topics().createNonPartitionedTopic(workerConfig.getFunctionMetadataTopic()); //create scheduler manager this.schedulerManager = new SchedulerManager(this.workerConfig, this.client, this.brokerAdmin, this.executor); From 41fa7cafb36016094cdac940203b806846971336 Mon Sep 17 00:00:00 2001 From: fxbing Date: Mon, 2 Sep 2019 23:37:30 +0800 Subject: [PATCH 14/23] fix test --- .../pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java | 2 +- .../pulsar/functions/worker/PulsarFunctionLocalRunTest.java | 2 +- .../pulsar/functions/worker/PulsarFunctionPublishTest.java | 2 +- .../apache/pulsar/functions/worker/PulsarFunctionStateTest.java | 2 +- .../pulsar/functions/worker/PulsarWorkerAssignmentTest.java | 2 +- .../test/java/org/apache/pulsar/io/PulsarFunctionAdminTest.java | 2 +- .../test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java | 2 +- .../pulsar/tests/integration/functions/PulsarFunctionsTest.java | 2 ++ 8 files changed, 9 insertions(+), 7 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java index afe79f493e09d..3782286fab7fe 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionE2ESecurityTest.java @@ -89,7 +89,7 @@ public class PulsarFunctionE2ESecurityTest { final String TENANT2 = "tenant2"; final String NAMESPACE = "test-ns"; - String pulsarFunctionsNamespace = TENANT + "/use/pulsar-function-admin"; + String pulsarFunctionsNamespace = TENANT + "/pulsar-function-admin"; String primaryHost; String workerId; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java index 80047953acf5e..8c2525e6c048a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionLocalRunTest.java @@ -102,7 +102,7 @@ public class PulsarFunctionLocalRunTest { BrokerStats brokerStatsClient; WorkerService functionsWorkerService; final String tenant = "external-repl-prop"; - String pulsarFunctionsNamespace = tenant + "/" + CLUSTER + "/pulsar-function-admin"; + String pulsarFunctionsNamespace = tenant + "/pulsar-function-admin"; String primaryHost; String workerId; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java index 0a0ea6257b9d5..efa06957111dc 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionPublishTest.java @@ -91,7 +91,7 @@ public class PulsarFunctionPublishTest { BrokerStats brokerStatsClient; WorkerService functionsWorkerService; final String tenant = "external-repl-prop"; - String pulsarFunctionsNamespace = tenant + "/use/pulsar-function-admin"; + String pulsarFunctionsNamespace = tenant + "/pulsar-function-admin"; String primaryHost; String workerId; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionStateTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionStateTest.java index a7de6c9962ffc..d0b85c3a87fb8 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionStateTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarFunctionStateTest.java @@ -90,7 +90,7 @@ public class PulsarFunctionStateTest { BrokerStats brokerStatsClient; WorkerService functionsWorkerService; final String tenant = "external-repl-prop"; - String pulsarFunctionsNamespace = tenant + "/use/pulsar-function-admin"; + String pulsarFunctionsNamespace = tenant + "/pulsar-function-admin"; String primaryHost; String workerId; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarWorkerAssignmentTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarWorkerAssignmentTest.java index 55e4a70dcc812..27a7d6981f07f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarWorkerAssignmentTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/functions/worker/PulsarWorkerAssignmentTest.java @@ -70,7 +70,7 @@ public class PulsarWorkerAssignmentTest { BrokerStats brokerStatsClient; WorkerService functionsWorkerService; final String tenant = "external-repl-prop"; - final String pulsarFunctionsNamespace = tenant + "/use/pulsar-function-admin"; + final String pulsarFunctionsNamespace = tenant + "/pulsar-function-admin"; String primaryHost; String workerId; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionAdminTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionAdminTest.java index ecfce1c738541..0c5db1922f010 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionAdminTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionAdminTest.java @@ -72,7 +72,7 @@ public class PulsarFunctionAdminTest { WorkerServer functionsWorkerServer; WorkerService functionsWorkerService; final String tenant = "external-repl-prop"; - String pulsarFunctionsNamespace = tenant + "/use/pulsar-function-admin"; + String pulsarFunctionsNamespace = tenant + "/pulsar-function-admin"; String primaryHost; private final int ZOOKEEPER_PORT = PortManager.nextFreePort(); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java index addd2c6928a5a..d8f9cecd07634 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/io/PulsarFunctionE2ETest.java @@ -118,7 +118,7 @@ public class PulsarFunctionE2ETest { BrokerStats brokerStatsClient; WorkerService functionsWorkerService; final String tenant = "external-repl-prop"; - String pulsarFunctionsNamespace = tenant + "/use/pulsar-function-admin"; + String pulsarFunctionsNamespace = tenant + "/pulsar-function-admin"; String primaryHost; String workerId; diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java index 7fda45be2a51f..915a6f8e74910 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java @@ -869,6 +869,8 @@ public void testFunctionLocalRun(Runtime runtime) throws Exception { try (PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl(pulsarCluster.getHttpServiceUrl()).build()) { + admin.topics().createNonPartitionedTopic(inputTopicName); + admin.topics().createNonPartitionedTopic(outputTopicName); retryStrategically((test) -> { try { return admin.topics().getStats(inputTopicName).subscriptions.size() == 1; From 87dd93888a244e984f01e6acdc0c7ed75c353084 Mon Sep 17 00:00:00 2001 From: fxbing Date: Thu, 5 Sep 2019 15:07:32 +0800 Subject: [PATCH 15/23] create async method for default partitioned topic creation --- .../pulsar/broker/admin/AdminResource.java | 42 ++++++++++++------- 1 file changed, 28 insertions(+), 14 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index 58873c050fcef..b3e52f362d6de 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -606,22 +606,16 @@ protected static CompletableFuture fetchPartitionedTop fetchPartitionedTopicMetadataAsync(pulsar, path).whenComplete((metadata, ex) -> { if (ex != null) { metadataFuture.completeExceptionally(ex); - // If topic is already exist, creating partitioned topic is not allowed. + // If topic is already exist, creating partitioned topic is not allowed. } else if (metadata.partitions == 0 && !topicExist && allowAutoTopicCreation && TopicType.PARTITIONED.toString().equals(topicType)) { - int defaultNumPartitions = pulsar.getConfiguration().getDefaultNumPartitions(); - checkArgument(defaultNumPartitions > 0, "Default number of partitions should be more than 0"); - try { - PartitionedTopicMetadata configMetadata = new PartitionedTopicMetadata(defaultNumPartitions); - byte[] content = jsonMapper().writeValueAsBytes(configMetadata); - ZkUtils.createFullPathOptimistic(pulsar.getGlobalZkCache().getZooKeeper(), path, content, - ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); - // we wait for the data to be synced in all quorums and the observers - Thread.sleep(PARTITIONED_TOPIC_WAIT_SYNC_TIME_MS); - metadataFuture.complete(configMetadata); - } catch (JsonProcessingException | InterruptedException | KeeperException e) { - metadataFuture.completeExceptionally(e); - } + createDefaultPartitionedTopicAsync(pulsar, path).whenComplete((defaultMetadata, e) -> { + if (e == null) { + metadataFuture.complete(defaultMetadata); + } else { + metadataFuture.completeExceptionally(e); + } + }); } else { metadataFuture.complete(metadata); } @@ -632,6 +626,26 @@ protected static CompletableFuture fetchPartitionedTop return metadataFuture; } + protected static CompletableFuture createDefaultPartitionedTopicAsync( + PulsarService pulsar, String path) { + int defaultNumPartitions = pulsar.getConfiguration().getDefaultNumPartitions(); + checkArgument(defaultNumPartitions > 0, "Default number of partitions should be more than 0"); + PartitionedTopicMetadata configMetadata = new PartitionedTopicMetadata(defaultNumPartitions); + CompletableFuture partitionedTopicFuture = new CompletableFuture<>(); + try { + byte[] content = jsonMapper().writeValueAsBytes(configMetadata); + ZkUtils.createFullPathOptimistic(pulsar.getGlobalZkCache().getZooKeeper(), path, content, + ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + // we wait for the data to be synced in all quorums and the observers + Thread.sleep(PARTITIONED_TOPIC_WAIT_SYNC_TIME_MS); + partitionedTopicFuture.complete(configMetadata); + } catch (JsonProcessingException | InterruptedException | KeeperException e) { + log.error("Failed to create default partitioned topic.", e); + partitionedTopicFuture.completeExceptionally(e); + } + return partitionedTopicFuture; + } + protected void validateClusterExists(String cluster) { try { if (!clustersCache().get(path("clusters", cluster)).isPresent()) { From a38efb3300f64c5859f598e718b527e662776fd4 Mon Sep 17 00:00:00 2001 From: fxbing Date: Fri, 6 Sep 2019 14:16:31 +0800 Subject: [PATCH 16/23] fix comment --- .../pulsar/broker/admin/AdminResource.java | 23 ++++++++++++++----- 1 file changed, 17 insertions(+), 6 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index b3e52f362d6de..69d6988cea882 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -634,12 +634,23 @@ protected static CompletableFuture createDefaultPartit CompletableFuture partitionedTopicFuture = new CompletableFuture<>(); try { byte[] content = jsonMapper().writeValueAsBytes(configMetadata); - ZkUtils.createFullPathOptimistic(pulsar.getGlobalZkCache().getZooKeeper(), path, content, - ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); - // we wait for the data to be synced in all quorums and the observers - Thread.sleep(PARTITIONED_TOPIC_WAIT_SYNC_TIME_MS); - partitionedTopicFuture.complete(configMetadata); - } catch (JsonProcessingException | InterruptedException | KeeperException e) { + ZkUtils.asyncCreateFullPathOptimistic(pulsar.getGlobalZkCache().getZooKeeper(), path, content, + ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT, (rc, path1, ctx, name) -> { + if (rc == KeeperException.Code.OK.intValue()) { + // we wait for the data to be synced in all quorums and the observers + try { + Thread.sleep(PARTITIONED_TOPIC_WAIT_SYNC_TIME_MS); + } catch (InterruptedException exc) { + partitionedTopicFuture.completeExceptionally(exc); + return; + } + partitionedTopicFuture.complete(configMetadata); + } else { + partitionedTopicFuture.completeExceptionally( + KeeperException.create(KeeperException.Code.get(rc))); + } + }, null); + } catch (JsonProcessingException e) { log.error("Failed to create default partitioned topic.", e); partitionedTopicFuture.completeExceptionally(e); } From 30913de305b20b66348f06f92e1dcba818f3ad1a Mon Sep 17 00:00:00 2001 From: fxbing Date: Fri, 6 Sep 2019 15:04:44 +0800 Subject: [PATCH 17/23] fix comment --- .../main/java/org/apache/pulsar/broker/admin/AdminResource.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index 69d6988cea882..9de3b8e9301d9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -603,7 +603,7 @@ protected static CompletableFuture fetchPartitionedTop topicName, e.getMessage(), e); throw new RestException(e); } - fetchPartitionedTopicMetadataAsync(pulsar, path).whenComplete((metadata, ex) -> { + fetchPartitionedTopicMetadataAsync(pulsar, path).whenCompleteAsync((metadata, ex) -> { if (ex != null) { metadataFuture.completeExceptionally(ex); // If topic is already exist, creating partitioned topic is not allowed. From 62f23134cc5246e4331115bb4b8980b76343ea15 Mon Sep 17 00:00:00 2001 From: fxbing Date: Fri, 6 Sep 2019 15:26:50 +0800 Subject: [PATCH 18/23] fix comment --- .../pulsar/broker/admin/AdminResource.java | 23 +++++-------------- 1 file changed, 6 insertions(+), 17 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index 9de3b8e9301d9..86bd1fc96e618 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -634,23 +634,12 @@ protected static CompletableFuture createDefaultPartit CompletableFuture partitionedTopicFuture = new CompletableFuture<>(); try { byte[] content = jsonMapper().writeValueAsBytes(configMetadata); - ZkUtils.asyncCreateFullPathOptimistic(pulsar.getGlobalZkCache().getZooKeeper(), path, content, - ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT, (rc, path1, ctx, name) -> { - if (rc == KeeperException.Code.OK.intValue()) { - // we wait for the data to be synced in all quorums and the observers - try { - Thread.sleep(PARTITIONED_TOPIC_WAIT_SYNC_TIME_MS); - } catch (InterruptedException exc) { - partitionedTopicFuture.completeExceptionally(exc); - return; - } - partitionedTopicFuture.complete(configMetadata); - } else { - partitionedTopicFuture.completeExceptionally( - KeeperException.create(KeeperException.Code.get(rc))); - } - }, null); - } catch (JsonProcessingException e) { + ZkUtils.createFullPathOptimistic(pulsar.getGlobalZkCache().getZooKeeper(), path, content, + ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT); + // we wait for the data to be synced in all quorums and the observers + Thread.sleep(PARTITIONED_TOPIC_WAIT_SYNC_TIME_MS); + partitionedTopicFuture.complete(configMetadata); + } catch (JsonProcessingException | KeeperException | InterruptedException e) { log.error("Failed to create default partitioned topic.", e); partitionedTopicFuture.completeExceptionally(e); } From 65978c63f8da320209bd3cd95cd136ada9a3e5ee Mon Sep 17 00:00:00 2001 From: fxbing Date: Fri, 6 Sep 2019 20:38:03 +0800 Subject: [PATCH 19/23] fix comment --- .../pulsar/broker/admin/AdminResource.java | 18 ++++++++++++++++++ .../apache/pulsar/client/impl/HttpClient.java | 1 - .../pulsar/client/impl/HttpLookupService.java | 3 ++- 3 files changed, 20 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java index 86bd1fc96e618..ed497e80c01bf 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/AdminResource.java @@ -612,6 +612,24 @@ protected static CompletableFuture fetchPartitionedTop createDefaultPartitionedTopicAsync(pulsar, path).whenComplete((defaultMetadata, e) -> { if (e == null) { metadataFuture.complete(defaultMetadata); + } else if (e instanceof KeeperException) { + try { + Thread.sleep(PARTITIONED_TOPIC_WAIT_SYNC_TIME_MS); + if (!pulsar.getGlobalZkCache().exists(path)){ + metadataFuture.completeExceptionally(e); + return; + } + } catch (InterruptedException | KeeperException exc) { + metadataFuture.completeExceptionally(exc); + return; + } + fetchPartitionedTopicMetadataAsync(pulsar, path).whenComplete((metadata2, ex2) -> { + if (ex2 != null) { + metadataFuture.completeExceptionally(ex2); + } else { + metadataFuture.complete(metadata2); + } + }); } else { metadataFuture.completeExceptionally(e); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpClient.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpClient.java index a4211ff2b1e82..96c62d2e9b6ef 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpClient.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpClient.java @@ -152,7 +152,6 @@ public CompletableFuture get(String path, Class clazz) { // auth complete, use a new builder BoundRequestBuilder builder = httpClient.prepareGet(requestUrl) .setHeader("Accept", "application/json"); - builder.addQueryParam("checkAllowAutoCreation", "true"); if (authData.hasDataForHttp()) { Set> headers; diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpLookupService.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpLookupService.java index 4fe030dd0647e..e5334dfd1b6a0 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpLookupService.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpLookupService.java @@ -103,7 +103,8 @@ public CompletableFuture> getBroker(T public CompletableFuture getPartitionedTopicMetadata(TopicName topicName) { String format = topicName.isV2() ? "admin/v2/%s/partitions" : "admin/%s/partitions"; - return httpClient.get(String.format(format, topicName.getLookupName()), PartitionedTopicMetadata.class); + return httpClient.get(String.format(format, topicName.getLookupName()) + "?checkAllowAutoCreation=true", + PartitionedTopicMetadata.class); } public String getServiceUrl() { From 84bd59b47b52635bcd390ee7e0233afa81919e1a Mon Sep 17 00:00:00 2001 From: fxbing Date: Fri, 6 Sep 2019 21:09:17 +0800 Subject: [PATCH 20/23] add test --- .../BrokerServiceAutoTopicCreationTest.java | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java index cf08c0e4ea892..cd1aec9922a9f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerServiceAutoTopicCreationTest.java @@ -19,6 +19,7 @@ package org.apache.pulsar.broker.service; import org.apache.pulsar.client.api.PulsarClientException; +import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertTrue; import org.testng.annotations.AfterClass; @@ -99,4 +100,21 @@ public void testAutoTopicCreationDisableIfNonPartitionedTopicAlreadyExist() thro } assertTrue(admin.namespaces().getTopics("prop/ns-abc").contains(topicName)); } + + /** + * CheckAllowAutoCreation's default value is false. + * So using getPartitionedTopicMetadata() directly will not produce partitioned topic + * even if the option to automatically create partitioned topic is configured + */ + @Test + public void testGetPartitionedMetadataWithoutCheckAllowAutoCreation() throws Exception{ + pulsar.getConfiguration().setAllowAutoTopicCreation(true); + pulsar.getConfiguration().setAllowAutoTopicCreationType("partitioned"); + pulsar.getConfiguration().setDefaultNumPartitions(3); + + final String topicName = "persistent://prop/ns-abc/test-topic-3"; + int partitions = admin.topics().getPartitionedTopicMetadata(topicName).partitions; + assertEquals(partitions, 0); + assertFalse(admin.namespaces().getTopics("prop/ns-abc").contains(topicName)); + } } From ee69a44e1b03e3c372c57e621077a64e6375b91f Mon Sep 17 00:00:00 2001 From: fxbing Date: Mon, 9 Sep 2019 14:53:05 +0800 Subject: [PATCH 21/23] fix test --- .../integration/functions/PulsarFunctionsTest.java | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java index 915a6f8e74910..40c5144e090fc 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java @@ -1116,6 +1116,10 @@ private void testFunctionNegAck(Runtime runtime) throws Exception { String inputTopicName = "persistent://public/default/test-neg-ack-" + runtime + "-input-" + randomName(8); String outputTopicName = "test-neg-ack-" + runtime + "-output-" + randomName(8); + try (PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl(pulsarCluster.getHttpServiceUrl()).build()) { + admin.topics().createNonPartitionedTopic(inputTopicName); + admin.topics().createNonPartitionedTopic(outputTopicName); + } String functionName = "test-neg-ack-fn-" + randomName(8); final int numMessages = 20; @@ -1292,6 +1296,10 @@ private void testPublishFunction(Runtime runtime) throws Exception { String inputTopicName = "persistent://public/default/test-publish-" + runtime + "-input-" + randomName(8); String outputTopicName = "test-publish-" + runtime + "-output-" + randomName(8); + try (PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl(pulsarCluster.getHttpServiceUrl()).build()) { + admin.topics().createNonPartitionedTopic(inputTopicName); + admin.topics().createNonPartitionedTopic(outputTopicName); + } String functionName = "test-publish-fn-" + randomName(8); final int numMessages = 10; @@ -1418,6 +1426,10 @@ private void testExclamationFunction(Runtime runtime, String inputTopicName = "persistent://public/default/test-exclamation-" + runtime + "-input-" + randomName(8); String outputTopicName = "test-exclamation-" + runtime + "-output-" + randomName(8); + try (PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl(pulsarCluster.getHttpServiceUrl()).build()) { + admin.topics().createNonPartitionedTopic(inputTopicName); + admin.topics().createNonPartitionedTopic(outputTopicName); + } if (isTopicPattern) { @Cleanup PulsarClient client = PulsarClient.builder() .serviceUrl(pulsarCluster.getPlainTextServiceUrl()) From 0408e34fab8359121b30cfc66b31e309c66b5c79 Mon Sep 17 00:00:00 2001 From: fxbing Date: Mon, 9 Sep 2019 22:08:12 +0800 Subject: [PATCH 22/23] fix test --- .../tests/integration/functions/PulsarFunctionsTest.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java index 40c5144e090fc..1f0355356763f 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java @@ -959,7 +959,10 @@ public void testWindowFunction(String type, String[] expectedResults) throws Exc String inputTopicName = "test-" + type + "-count-window-" + functionRuntimeType + "-input-" + randomName(8); String outputTopicName = "test-" + type + "-count-window-" + functionRuntimeType + "-output-" + randomName(8); - + try (PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl(pulsarCluster.getHttpServiceUrl()).build()) { + admin.topics().createNonPartitionedTopic(inputTopicName); + admin.topics().createNonPartitionedTopic(outputTopicName); + } CommandGenerator generator = CommandGenerator.createDefaultGenerator( inputTopicName, From fd5da63bc88a032d3907b7333a01d0435de3b587 Mon Sep 17 00:00:00 2001 From: fxbing Date: Tue, 10 Sep 2019 10:46:28 +0800 Subject: [PATCH 23/23] fix test --- .../tests/integration/functions/PulsarFunctionsTest.java | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java index 1f0355356763f..22d4c6b288ed7 100644 --- a/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java +++ b/tests/integration/src/test/java/org/apache/pulsar/tests/integration/functions/PulsarFunctionsTest.java @@ -546,6 +546,9 @@ private void testSource(SourceTester tester) throws Exception { .serviceUrl(pulsarCluster.getPlainTextServiceUrl()) .build(); + @Cleanup + PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl(pulsarCluster.getHttpServiceUrl()).build(); + admin.topics().createNonPartitionedTopic(outputTopicName); @Cleanup Consumer consumer = client.newConsumer(Schema.STRING) .topic(outputTopicName) @@ -1944,6 +1947,11 @@ private void testDebeziumMySqlConnect() .serviceUrl(pulsarCluster.getPlainTextServiceUrl()) .build(); + @Cleanup + PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl(pulsarCluster.getHttpServiceUrl()).build(); + admin.topics().createNonPartitionedTopic(consumeTopicName); + admin.topics().createNonPartitionedTopic(outputTopicName); + @Cleanup Consumer> consumer = client.newConsumer(KeyValueSchema.kvBytes()) .topic(consumeTopicName)