From 830ac9b94bcef2a5e341ec2af2d89f2fcea1fc6f Mon Sep 17 00:00:00 2001 From: penghui Date: Thu, 5 Mar 2020 11:55:56 +0800 Subject: [PATCH 1/7] Fix create partitioned topic with a substring of an existing topic name. --- .../pulsar/broker/admin/AdminResource.java | 77 +++++++++++++++++++ .../admin/impl/PersistentTopicsBase.java | 76 +----------------- .../broker/admin/v1/NonPersistentTopics.java | 35 +-------- .../broker/admin/v1/PersistentTopics.java | 10 ++- .../broker/admin/v2/NonPersistentTopics.java | 38 ++------- .../broker/admin/v2/PersistentTopics.java | 12 ++- .../pulsar/broker/admin/AdminApiTest.java | 4 + .../apache/pulsar/broker/admin/AdminTest.java | 4 +- .../broker/admin/PersistentTopicsTest.java | 34 +++++--- 9 files changed, 136 insertions(+), 154 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 722da4f49214c..092c645ce6906 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 @@ -36,6 +36,7 @@ import javax.servlet.ServletContext; import javax.ws.rs.WebApplicationException; +import javax.ws.rs.container.AsyncResponse; import javax.ws.rs.core.Response; import javax.ws.rs.core.Response.Status; import javax.ws.rs.core.UriBuilder; @@ -46,6 +47,7 @@ import org.apache.pulsar.broker.cache.LocalZooKeeperCacheService; import org.apache.pulsar.broker.web.PulsarWebResource; 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.NamespaceBundleFactory; @@ -707,4 +709,79 @@ protected List getPartitionedTopicList(TopicDomain topicDomain) { partitionedTopics.sort(null); return partitionedTopics; } + + protected void internalCreatePartitionedTopic(AsyncResponse asyncResponse, int numPartitions) { + validateAdminAccessForTenant(topicName.getTenant()); + if (numPartitions <= 0) { + asyncResponse.resume(new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be more than 0")); + return; + } + checkTopicExistsAsync(topicName).thenAccept(exists -> { + if (exists) { + log.warn("[{}] Failed to create already existing topic {}", clientAppId(), topicName); + asyncResponse.resume(new RestException(Status.CONFLICT, "This topic already exists")); + } else { + try { + String path = ZkAdminPaths.partitionedTopicPath(topicName); + byte[] data = jsonMapper().writeValueAsBytes(new PartitionedTopicMetadata(numPartitions)); + zkCreateOptimisticAsync(globalZk(), path, data, (rc, s, o, s1) -> { + if (KeeperException.Code.OK.intValue() == rc) { + if (topicName.isPersistent()) { + tryCreatePartitionsAsync(numPartitions); + } + globalZk().sync(path, (rc2, s2, ctx) -> { + if (KeeperException.Code.OK.intValue() == rc2) { + log.info("[{}] Successfully created partitioned topic {}", clientAppId(), topicName); + asyncResponse.resume(Response.noContent().build()); + } else { + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, KeeperException.create(KeeperException.Code.get(rc))); + asyncResponse.resume(new RestException(KeeperException.create(KeeperException.Code.get(rc)))); + } + }, null); + } else if (KeeperException.Code.NODEEXISTS.intValue() == rc) { + log.warn("[{}] Failed to create already existing partitioned topic {}", clientAppId(), topicName); + asyncResponse.resume(new RestException(Status.CONFLICT, "Partitioned topic already exists")); + } else if (KeeperException.Code.BADVERSION.intValue() == rc) { + log.warn("[{}] Failed to create partitioned topic {}: concurrent modification", clientAppId(), + topicName); + asyncResponse.resume(new RestException(Status.CONFLICT, "Concurrent modification")); + } else { + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, KeeperException.create(KeeperException.Code.get(rc))); + asyncResponse.resume(new RestException(KeeperException.create(KeeperException.Code.get(rc)))); + } + }); + } catch (Exception e) { + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); + asyncResponse.resume(new RestException(e)); + } + } + }).exceptionally(ex -> { + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, ex); + asyncResponse.resume(new RestException(ex)); + return null; + }); + } + + /** + * Check the exists topics contains the given topic. + * Since there are topic partitions and non-partitioned topics in Pulsar, must ensure both partitions + * and non-partitioned topics are not duplicated. So, if compare with a partition name, we should compare + * to the partitioned name of this partition. + * + * @param topicName given topic name + */ + protected CompletableFuture checkTopicExistsAsync(TopicName topicName) { + return pulsar().getNamespaceService().getListOfTopics(topicName.getNamespaceObject(), + PulsarApi.CommandGetTopicsOfNamespace.Mode.ALL) + .thenCompose(topics -> { + boolean exists = false; + for (String topic : topics) { + if (topicName.getPartitionedTopicName().equals(TopicName.get(topic).getPartitionedTopicName())) { + exists = true; + break; + } + } + return CompletableFuture.completedFuture(exists); + }); + } } 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 503b21b430d58..d265ee8c20965 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,7 +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; @@ -390,46 +390,6 @@ protected void internalRevokePermissionsOnTopic(String role) { revokePermissions(topicName.toString(), role); } - protected void internalCreatePartitionedTopic(int numPartitions) { - validateAdminAccessForTenant(topicName.getTenant()); - if (numPartitions <= 0) { - throw new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be more than 0"); - } - validatePartitionTopicName(topicName.getLocalName()); - try { - boolean topicExist = pulsar().getNamespaceService() - .getListOfTopics(topicName.getNamespaceObject(), PulsarApi.CommandGetTopicsOfNamespace.Mode.ALL) - .join() - .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)); - zkCreateOptimistic(path, data); - tryCreatePartitionsAsync(numPartitions); - // Sync data to all quorums and the observers - zkSync(path); - log.info("[{}] Successfully created partitioned topic {}", clientAppId(), topicName); - } catch (KeeperException.NodeExistsException e) { - log.warn("[{}] Failed to create already existing partitioned topic {}", clientAppId(), topicName); - throw new RestException(Status.CONFLICT, "Partitioned topic already exists"); - } catch (KeeperException.BadVersionException e) { - log.warn("[{}] Failed to create partitioned topic {}: concurrent modification", clientAppId(), - topicName); - throw new RestException(Status.CONFLICT, "Concurrent modification"); - } catch (Exception e) { - log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); - throw new RestException(e); - } - } - protected void internalCreateNonPartitionedTopic(boolean authoritative) { validateAdminAccessForTenant(topicName.getTenant()); validateNonPartitionTopicName(topicName.getLocalName()); @@ -2071,40 +2031,6 @@ private void validatePartitionTopicUpdate(String topicName, int numberOfPartitio } } - /** - * Validate partitioned topic name. - * Validation will fail and throw RestException if - * 1) There's already a partitioned topic with same topic name and have some of its partition created. - * 2) There's already non partition topic with same name and contains partition suffix "-partition-" - * followed by numeric value. In this case internal created partition of partitioned topic could override - * the existing non partition topic. - * - * @param topicName - */ - private void validatePartitionTopicName(String topicName) { - List existingTopicList = internalGetList(); - String prefix = topicName + TopicName.PARTITIONED_TOPIC_SUFFIX; - for (String existingTopicName : existingTopicList) { - if (existingTopicName.contains(prefix)) { - try { - Long.parseLong(existingTopicName.substring( - existingTopicName.indexOf(TopicName.PARTITIONED_TOPIC_SUFFIX) - + TopicName.PARTITIONED_TOPIC_SUFFIX.length())); - log.warn("[{}] Already have topic {} which contains partition " + - "suffix '-partition-' and end with numeric value. Creation of partitioned topic {}" - + "could cause conflict.", clientAppId(), existingTopicName, topicName); - throw new RestException(Status.PRECONDITION_FAILED, - "Already have topic " + existingTopicName + " which contains partition suffix '-partition-' " + - "and end with numeric value, Creation of partitioned topic " + topicName + - " could cause conflict."); - } catch (NumberFormatException e) { - // Do nothing, if value after partition suffix is not pure numeric value, - // as it can't conflict with internal created partitioned topic's name. - } - } - } - } - /** * Validate non partition topic name, * Validation will fail and throw RestException if 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 4bc0ddfd154d9..bcbde9606d4bb 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 @@ -124,41 +124,14 @@ public PersistentTopicInternalStats getInternalStats(@PathParam("property") Stri @ApiOperation(hidden = true, value = "Create a partitioned topic.", notes = "It needs to be called before creating a producer on a partitioned topic.") @ApiResponses(value = { @ApiResponse(code = 403, message = "Don't have admin permission"), @ApiResponse(code = 409, message = "Partitioned topic already exist") }) - public void createPartitionedTopic(@PathParam("property") String property, @PathParam("cluster") String cluster, + public void createPartitionedTopic(@Suspended final AsyncResponse asyncResponse, @PathParam("property") String property, @PathParam("cluster") String cluster, @PathParam("namespace") String namespace, @PathParam("topic") @Encoded String encodedTopic, int numPartitions) { - validateTopicName(property, cluster, namespace, encodedTopic); - validateAdminAccessForTenant(topicName.getTenant()); - 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) - .join() - .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()); - byte[] data = jsonMapper().writeValueAsBytes(new PartitionedTopicMetadata(numPartitions)); - zkCreateOptimistic(path, data); - // Sync data to all quorums and the observers - zkSync(path); - log.info("[{}] Successfully created partitioned topic {}", clientAppId(), topicName); - } catch (KeeperException.NodeExistsException e) { - log.warn("[{}] Failed to create already existing partitioned topic {}", clientAppId(), topicName); - throw new RestException(Status.CONFLICT, "Partitioned topic already exist"); + validateTopicName(property, cluster, namespace, encodedTopic); + internalCreatePartitionedTopic(asyncResponse, numPartitions); } catch (Exception e) { - log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); - throw new RestException(e); + asyncResponse.resume(e); } } 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 836ca14282781..7c814085f85d2 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 @@ -147,11 +147,15 @@ public void revokePermissionsOnTopic(@PathParam("property") String property, @ApiResponse(code = 307, message = "Current broker doesn't serve the namespace of this topic"), @ApiResponse(code = 403, message = "Don't have admin permission"), @ApiResponse(code = 409, message = "Partitioned topic already exist") }) - public void createPartitionedTopic(@PathParam("property") String property, @PathParam("cluster") String cluster, + public void createPartitionedTopic(@Suspended final AsyncResponse asyncResponse, @PathParam("property") String property, @PathParam("cluster") String cluster, @PathParam("namespace") String namespace, @PathParam("topic") @Encoded String encodedTopic, int numPartitions) { - validateTopicName(property, cluster, namespace, encodedTopic); - internalCreatePartitionedTopic(numPartitions); + try { + validateTopicName(property, cluster, namespace, encodedTopic); + internalCreatePartitionedTopic(asyncResponse, numPartitions); + } catch (Exception e) { + asyncResponse.resume(e); + } } /** 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 7e88eed2fe5ab..6493c35942845 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 @@ -164,6 +164,7 @@ public PersistentTopicInternalStats getInternalStats( @ApiResponse(code = 503, message = "Failed to validate global cluster configuration"), }) public void createPartitionedTopic( + @Suspended final AsyncResponse asyncResponse, @ApiParam(value = "Specify the tenant", required = true) @PathParam("tenant") String tenant, @ApiParam(value = "Specify the namespace", required = true) @@ -172,39 +173,14 @@ public void createPartitionedTopic( @PathParam("topic") @Encoded String encodedTopic, @ApiParam(value = "The number of partitions for the topic", required = true, type = "int", defaultValue = "0") int numPartitions) { - validateGlobalNamespaceOwnership(tenant,namespace); - validateTopicName(tenant, namespace, encodedTopic); - validateAdminAccessForTenant(topicName.getTenant()); - 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) - .join() - .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()); - byte[] data = jsonMapper().writeValueAsBytes(new PartitionedTopicMetadata(numPartitions)); - zkCreateOptimistic(path, data); - // Sync data to all quorums and the observers - zkSync(path); - log.info("[{}] Successfully created partitioned topic {}", clientAppId(), topicName); - } catch (KeeperException.NodeExistsException e) { - log.warn("[{}] Failed to create already existing partitioned topic {}", clientAppId(), topicName); - throw new RestException(Status.CONFLICT, "Partitioned topic already exists"); + validateGlobalNamespaceOwnership(tenant,namespace); + validateTopicName(tenant, namespace, encodedTopic); + validateAdminAccessForTenant(topicName.getTenant()); + internalCreatePartitionedTopic(asyncResponse, numPartitions); } catch (Exception e) { - log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); - throw new RestException(e); + asyncResponse.resume(e); } } 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 8c59fa5b6b079..f3e5a2cd35fd7 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 @@ -192,6 +192,7 @@ public void revokePermissionsOnTopic( @ApiResponse(code = 503, message = "Failed to validate global cluster configuration") }) public void createPartitionedTopic( + @Suspended final AsyncResponse asyncResponse, @ApiParam(value = "Specify the tenant", required = true) @PathParam("tenant") String tenant, @ApiParam(value = "Specify the namespace", required = true) @@ -200,9 +201,14 @@ public void createPartitionedTopic( @PathParam("topic") @Encoded String encodedTopic, @ApiParam(value = "The number of partitions for the topic", required = true, type = "int", defaultValue = "0") int numPartitions) { - validateGlobalNamespaceOwnership(tenant,namespace); - validatePartitionedTopicName(tenant, namespace, encodedTopic); - internalCreatePartitionedTopic(numPartitions); + try { + validateGlobalNamespaceOwnership(tenant,namespace); + validatePartitionedTopicName(tenant, namespace, encodedTopic); + validateAdminAccessForTenant(topicName.getTenant()); + internalCreatePartitionedTopic(asyncResponse, numPartitions); + } catch (Exception e) { + asyncResponse.resume(e); + } } @PUT diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java index f0ee4f60e985c..5a86d9b714f06 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java @@ -2010,6 +2010,10 @@ public void testPersistentTopicCreation() throws Exception { } catch (PulsarAdminException e) { assertTrue(e instanceof ConflictException); } + + // Check create partitioned topic with substring topic name + admin.topics().createPartitionedTopic("persistent://prop-xyz/ns1/create_substring_topic", 1); + admin.topics().createPartitionedTopic("persistent://prop-xyz/ns1/substring_topic", 1); } /** diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java index 00d1a31b1ef53..b839e6827477e 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java @@ -677,7 +677,9 @@ void persistentTopics() throws Exception { verify(response, times(1)).resume(Lists.newArrayList()); // create topic assertEquals(persistentTopics.getPartitionedTopicList(property, cluster, namespace), Lists.newArrayList()); - persistentTopics.createPartitionedTopic(property, cluster, namespace, topic, 5); + response = mock(AsyncResponse.class); + persistentTopics.createPartitionedTopic(response, property, cluster, namespace, topic, 5); + verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); assertEquals(persistentTopics.getPartitionedTopicList(property, cluster, namespace), Lists .newArrayList(String.format("persistent://%s/%s/%s/%s", property, cluster, namespace, topic))); 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 1825d313124e7..164f3e893e558 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 @@ -157,7 +157,9 @@ public void testGetSubscriptions() { "Partitioned Topic not found: persistent://my-tenant/my-namespace/topic-not-found-partition-0 has zero partitions"); // 3) Create the partitioned topic - persistentTopics.createPartitionedTopic(testTenant, testNamespace, testLocalTopicName, 3); + response = mock(AsyncResponse.class); + persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, testLocalTopicName, 3); + verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); // 4) Create a subscription response = mock(AsyncResponse.class); @@ -250,7 +252,9 @@ public void testCreatePartitionedTopicHavingNonPartitionTopicWithPartitionSuffix doReturn(mockLocalZooKeeperCacheService).when(pulsar).getLocalZkCacheService(); doReturn(mockZooKeeperChildrenCache).when(mockLocalZooKeeperCacheService).managedLedgerListCache(); doReturn(ImmutableSet.of(nonPartitionTopicName1, nonPartitionTopicName2)).when(mockZooKeeperChildrenCache).get(anyString()); - persistentTopics.createPartitionedTopic(testTenant, testNamespace, partitionedTopicName, 5); + AsyncResponse response = mock(AsyncResponse.class); + persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, partitionedTopicName, 5); + verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); } @Test(expectedExceptions = RestException.class) @@ -269,7 +273,9 @@ public void testUpdatePartitionedTopicHavingNonPartitionTopicWithPartitionSuffix return null; }).when(persistentTopics).validatePartitionedTopicName(any(), any(), any()); doNothing().when(persistentTopics).validateAdminAccessForTenant(anyString()); - persistentTopics.createPartitionedTopic(testTenant, testNamespace, partitionedTopicName, 5); + AsyncResponse response = mock(AsyncResponse.class); + persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, partitionedTopicName, 5); + verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); persistentTopics.updatePartitionedTopic(testTenant, testNamespace, partitionedTopicName, true, 10); } @@ -295,7 +301,9 @@ public void testUnloadTopic() { // 3) create partitioned topic and unload response = mock(AsyncResponse.class); - persistentTopics.createPartitionedTopic(testTenant, testNamespace, partitionTopicName, 6); + persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, partitionTopicName, 6); + verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); + response = mock(AsyncResponse.class); persistentTopics.unloadTopic(response, testTenant, testNamespace, partitionTopicName, true); responseCaptor = ArgumentCaptor.forClass(Response.class); verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); @@ -320,10 +328,13 @@ public void testUnloadTopicShallThrowNotFoundWhenTopicNotExist() { @Test public void testGetPartitionedTopicsList() throws KeeperException, InterruptedException, PulsarAdminException { + AsyncResponse response = mock(AsyncResponse.class); + persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, "test-topic1", 3); + verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); - persistentTopics.createPartitionedTopic(testTenant, testNamespace, "test-topic1", 3); - - nonPersistentTopic.createPartitionedTopic(testTenant, testNamespace, "test-topic2", 3); + response = mock(AsyncResponse.class); + nonPersistentTopic.createPartitionedTopic(response, testTenant, testNamespace, "test-topic2", 3); + verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); List persistentPartitionedTopics = persistentTopics.getPartitionedTopicList(testTenant, testNamespace); @@ -351,7 +362,9 @@ public void testGrantNonPartitionedTopic() { public void testGrantPartitionedTopic() { final String partitionedTopicName = "partitioned-topic"; final int numPartitions = 5; - persistentTopics.createPartitionedTopic(testTenant, testNamespace, partitionedTopicName, numPartitions); + AsyncResponse response = mock(AsyncResponse.class); + persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, partitionedTopicName, numPartitions); + verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); String role = "role"; Set expectActions = new HashSet<>(); @@ -387,8 +400,9 @@ public void testRevokeNonPartitionedTopic() { public void testRevokePartitionedTopic() { final String partitionedTopicName = "partitioned-topic"; final int numPartitions = 5; - persistentTopics.createPartitionedTopic(testTenant, testNamespace, partitionedTopicName, numPartitions); - + AsyncResponse response = mock(AsyncResponse.class); + persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, partitionedTopicName, numPartitions); + verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); String role = "role"; Set expectActions = new HashSet<>(); expectActions.add(AuthAction.produce); From d49e6800e09bbb7ed7414350ade86d4ba50619e3 Mon Sep 17 00:00:00 2001 From: penghui Date: Thu, 5 Mar 2020 19:04:48 +0800 Subject: [PATCH 2/7] Fix tests --- .../broker/admin/v2/PersistentTopics.java | 6 +++ .../pulsar/broker/admin/AdminApiTest.java | 4 +- .../apache/pulsar/broker/admin/AdminTest.java | 4 +- .../broker/admin/PersistentTopicsTest.java | 39 ++++++++++++++----- .../broker/admin/v1/V1_AdminApiTest.java | 4 +- 5 files changed, 40 insertions(+), 17 deletions(-) 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 f3e5a2cd35fd7..2794fce04d5c9 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 @@ -55,6 +55,9 @@ import io.swagger.annotations.ApiResponse; import io.swagger.annotations.ApiResponses; import io.swagger.annotations.ApiParam; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import static org.apache.pulsar.common.util.Codec.decode; /** @@ -207,6 +210,7 @@ public void createPartitionedTopic( validateAdminAccessForTenant(topicName.getTenant()); internalCreatePartitionedTopic(asyncResponse, numPartitions); } catch (Exception e) { + log.error("Unexpected exception on create partitioned topic.", e); asyncResponse.resume(e); } } @@ -1078,4 +1082,6 @@ public MessageId getLastMessageId( validateTopicName(tenant, namespace, encodedTopic); return internalGetLastMessageId(authoritative); } + + private static final Logger log = LoggerFactory.getLogger(PersistentTopics.class); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java index 5a86d9b714f06..d8326285076fb 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminApiTest.java @@ -919,9 +919,7 @@ public void partitionedTopics(String topicName) throws Exception { try { admin.topics().createPartitionedTopic(partitionedTopicName, 32); fail("Should have failed as the partitioned topic already exists"); - } catch (PreconditionFailedException e) { - // Expecting PreconditionFailedException instead of ConflictException as it'll - // fail validation before actually try to create metadata in ZK. + } catch (ConflictException ignore) { } producer = client.newProducer(Schema.BYTES) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java index b839e6827477e..3cae752f436b8 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AdminTest.java @@ -678,8 +678,10 @@ void persistentTopics() throws Exception { // create topic assertEquals(persistentTopics.getPartitionedTopicList(property, cluster, namespace), Lists.newArrayList()); response = mock(AsyncResponse.class); + ArgumentCaptor responseCaptor = ArgumentCaptor.forClass(Response.class); persistentTopics.createPartitionedTopic(response, property, cluster, namespace, topic, 5); - verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); + verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); + assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); assertEquals(persistentTopics.getPartitionedTopicList(property, cluster, namespace), Lists .newArrayList(String.format("persistent://%s/%s/%s/%s", property, cluster, namespace, topic))); 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 164f3e893e558..cac35a6e1eb18 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 @@ -56,6 +56,7 @@ import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.CompletableFuture; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyBoolean; @@ -158,14 +159,16 @@ public void testGetSubscriptions() { // 3) Create the partitioned topic response = mock(AsyncResponse.class); + ArgumentCaptor responseCaptor = ArgumentCaptor.forClass(Response.class); persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, testLocalTopicName, 3); - verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); + verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); + Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); // 4) Create a subscription response = mock(AsyncResponse.class); persistentTopics.createSubscription(response, testTenant, testNamespace, testLocalTopicName, "test", true, (MessageIdImpl) MessageId.earliest, false); - ArgumentCaptor responseCaptor = ArgumentCaptor.forClass(Response.class); + responseCaptor = ArgumentCaptor.forClass(Response.class); verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); @@ -241,7 +244,7 @@ public void testCreateNonPartitionedTopicWithInvalidName() { persistentTopics.createNonPartitionedTopic(testTenant, testNamespace, topicName, true); } - @Test(expectedExceptions = RestException.class) + @Test public void testCreatePartitionedTopicHavingNonPartitionTopicWithPartitionSuffix() throws KeeperException, InterruptedException { // Test the case in which user already has topic like topic-name-partition-123 created before we enforce the validation. final String nonPartitionTopicName1 = "standard-topic"; @@ -252,9 +255,12 @@ public void testCreatePartitionedTopicHavingNonPartitionTopicWithPartitionSuffix doReturn(mockLocalZooKeeperCacheService).when(pulsar).getLocalZkCacheService(); doReturn(mockZooKeeperChildrenCache).when(mockLocalZooKeeperCacheService).managedLedgerListCache(); doReturn(ImmutableSet.of(nonPartitionTopicName1, nonPartitionTopicName2)).when(mockZooKeeperChildrenCache).get(anyString()); + doReturn(CompletableFuture.completedFuture(ImmutableSet.of(nonPartitionTopicName1, nonPartitionTopicName2))).when(mockZooKeeperChildrenCache).getAsync(anyString()); AsyncResponse response = mock(AsyncResponse.class); + ArgumentCaptor errCaptor = ArgumentCaptor.forClass(RestException.class); persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, partitionedTopicName, 5); - verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); + verify(response, timeout(5000).times(1)).resume(errCaptor.capture()); + Assert.assertEquals(errCaptor.getValue().getResponse().getStatus(), Response.Status.CONFLICT.getStatusCode()); } @Test(expectedExceptions = RestException.class) @@ -267,6 +273,7 @@ public void testUpdatePartitionedTopicHavingNonPartitionTopicWithPartitionSuffix doReturn(mockLocalZooKeeperCacheService).when(pulsar).getLocalZkCacheService(); doReturn(mockZooKeeperChildrenCache).when(mockLocalZooKeeperCacheService).managedLedgerListCache(); doReturn(ImmutableSet.of(nonPartitionTopicName2)).when(mockZooKeeperChildrenCache).get(anyString()); + doReturn(CompletableFuture.completedFuture(ImmutableSet.of(nonPartitionTopicName2))).when(mockZooKeeperChildrenCache).getAsync(anyString()); doAnswer(invocation -> { persistentTopics.namespaceName = NamespaceName.get("tenant", "namespace"); persistentTopics.topicName = TopicName.get("persistent", "tenant", "cluster", "namespace", "topicname"); @@ -274,8 +281,10 @@ public void testUpdatePartitionedTopicHavingNonPartitionTopicWithPartitionSuffix }).when(persistentTopics).validatePartitionedTopicName(any(), any(), any()); doNothing().when(persistentTopics).validateAdminAccessForTenant(anyString()); AsyncResponse response = mock(AsyncResponse.class); + ArgumentCaptor responseCaptor = ArgumentCaptor.forClass(Response.class); persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, partitionedTopicName, 5); - verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); + verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); + Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); persistentTopics.updatePartitionedTopic(testTenant, testNamespace, partitionedTopicName, true, 10); } @@ -301,8 +310,10 @@ public void testUnloadTopic() { // 3) create partitioned topic and unload response = mock(AsyncResponse.class); + responseCaptor = ArgumentCaptor.forClass(Response.class); persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, partitionTopicName, 6); - verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); + verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); + Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); response = mock(AsyncResponse.class); persistentTopics.unloadTopic(response, testTenant, testNamespace, partitionTopicName, true); responseCaptor = ArgumentCaptor.forClass(Response.class); @@ -329,12 +340,16 @@ public void testUnloadTopicShallThrowNotFoundWhenTopicNotExist() { @Test public void testGetPartitionedTopicsList() throws KeeperException, InterruptedException, PulsarAdminException { AsyncResponse response = mock(AsyncResponse.class); + ArgumentCaptor responseCaptor = ArgumentCaptor.forClass(Response.class); persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, "test-topic1", 3); - verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); + verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); + Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); response = mock(AsyncResponse.class); + responseCaptor = ArgumentCaptor.forClass(Response.class); nonPersistentTopic.createPartitionedTopic(response, testTenant, testNamespace, "test-topic2", 3); - verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); + verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); + Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); List persistentPartitionedTopics = persistentTopics.getPartitionedTopicList(testTenant, testNamespace); @@ -363,8 +378,10 @@ public void testGrantPartitionedTopic() { final String partitionedTopicName = "partitioned-topic"; final int numPartitions = 5; AsyncResponse response = mock(AsyncResponse.class); + ArgumentCaptor responseCaptor = ArgumentCaptor.forClass(Response.class); persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, partitionedTopicName, numPartitions); - verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); + verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); + Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); String role = "role"; Set expectActions = new HashSet<>(); @@ -401,8 +418,10 @@ public void testRevokePartitionedTopic() { final String partitionedTopicName = "partitioned-topic"; final int numPartitions = 5; AsyncResponse response = mock(AsyncResponse.class); + ArgumentCaptor responseCaptor = ArgumentCaptor.forClass(Response.class); persistentTopics.createPartitionedTopic(response, testTenant, testNamespace, partitionedTopicName, numPartitions); - verify(response, timeout(5000).times(1)).resume(Response.noContent().build()); + verify(response, timeout(5000).times(1)).resume(responseCaptor.capture()); + Assert.assertEquals(responseCaptor.getValue().getStatus(), Response.Status.NO_CONTENT.getStatusCode()); String role = "role"; Set expectActions = new HashSet<>(); expectActions.add(AuthAction.produce); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/v1/V1_AdminApiTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/v1/V1_AdminApiTest.java index 373a6b7c48766..78e3dc16592ba 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/v1/V1_AdminApiTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/v1/V1_AdminApiTest.java @@ -879,9 +879,7 @@ public void partitionedTopics(String topicName) throws Exception { try { admin.topics().createPartitionedTopic(partitionedTopicName, 32); fail("Should have failed as the partitioned topic exists with its partition created"); - } catch (PreconditionFailedException e) { - // Expecting PreconditionFailedException instead of ConflictException as it'll - // fail validation before actually try to create metadata in ZK. + } catch (ConflictException ignore) { } producer = client.newProducer(Schema.BYTES) From 2495e5cea3721feb9358c8e2ccb32fda51d0dea2 Mon Sep 17 00:00:00 2001 From: penghui Date: Thu, 5 Mar 2020 23:02:25 +0800 Subject: [PATCH 3/7] handling resume async response exceptionally --- .../pulsar/broker/admin/AdminResource.java | 24 +++++++++++++++---- 1 file changed, 19 insertions(+), 5 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 092c645ce6906..6d120ded3c666 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 @@ -711,7 +711,13 @@ protected List getPartitionedTopicList(TopicDomain topicDomain) { } protected void internalCreatePartitionedTopic(AsyncResponse asyncResponse, int numPartitions) { - validateAdminAccessForTenant(topicName.getTenant()); + try { + validateAdminAccessForTenant(topicName.getTenant()); + } catch (Exception e) { + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); + resumeAsyncResponseExceptionally(asyncResponse, e); + return; + } if (numPartitions <= 0) { asyncResponse.resume(new RestException(Status.NOT_ACCEPTABLE, "Number of partitions should be more than 0")); return; @@ -734,8 +740,8 @@ protected void internalCreatePartitionedTopic(AsyncResponse asyncResponse, int n log.info("[{}] Successfully created partitioned topic {}", clientAppId(), topicName); asyncResponse.resume(Response.noContent().build()); } else { - log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, KeeperException.create(KeeperException.Code.get(rc))); - asyncResponse.resume(new RestException(KeeperException.create(KeeperException.Code.get(rc)))); + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, KeeperException.create(KeeperException.Code.get(rc2))); + asyncResponse.resume(new RestException(KeeperException.create(KeeperException.Code.get(rc2)))); } }, null); } else if (KeeperException.Code.NODEEXISTS.intValue() == rc) { @@ -752,12 +758,12 @@ protected void internalCreatePartitionedTopic(AsyncResponse asyncResponse, int n }); } catch (Exception e) { log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); - asyncResponse.resume(new RestException(e)); + resumeAsyncResponseExceptionally(asyncResponse, e); } } }).exceptionally(ex -> { log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, ex); - asyncResponse.resume(new RestException(ex)); + resumeAsyncResponseExceptionally(asyncResponse, ex); return null; }); } @@ -784,4 +790,12 @@ protected CompletableFuture checkTopicExistsAsync(TopicName topicName) return CompletableFuture.completedFuture(exists); }); } + + protected void resumeAsyncResponseExceptionally(AsyncResponse asyncResponse, Throwable throwable) { + if (throwable instanceof WebApplicationException) { + asyncResponse.resume(throwable); + } else { + asyncResponse.resume(new RestException(throwable)); + } + } } From 12825ba9f93e53b15e690d6e4790e6f2f9638cda Mon Sep 17 00:00:00 2001 From: penghui Date: Thu, 5 Mar 2020 23:07:23 +0800 Subject: [PATCH 4/7] Add more logs --- .../apache/pulsar/broker/admin/v1/NonPersistentTopics.java | 3 ++- .../org/apache/pulsar/broker/admin/v1/PersistentTopics.java | 3 ++- .../apache/pulsar/broker/admin/v2/NonPersistentTopics.java | 3 ++- .../org/apache/pulsar/broker/admin/v2/PersistentTopics.java | 4 ++-- 4 files changed, 8 insertions(+), 5 deletions(-) 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 bcbde9606d4bb..2338b0f055af5 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 @@ -131,7 +131,8 @@ public void createPartitionedTopic(@Suspended final AsyncResponse asyncResponse, validateTopicName(property, cluster, namespace, encodedTopic); internalCreatePartitionedTopic(asyncResponse, numPartitions); } catch (Exception e) { - asyncResponse.resume(e); + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); + resumeAsyncResponseExceptionally(asyncResponse, e); } } 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 7c814085f85d2..417510bfcb153 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 @@ -154,7 +154,8 @@ public void createPartitionedTopic(@Suspended final AsyncResponse asyncResponse, validateTopicName(property, cluster, namespace, encodedTopic); internalCreatePartitionedTopic(asyncResponse, numPartitions); } catch (Exception e) { - asyncResponse.resume(e); + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); + resumeAsyncResponseExceptionally(asyncResponse, e); } } 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 6493c35942845..3756f8285e7d4 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 @@ -180,7 +180,8 @@ public void createPartitionedTopic( validateAdminAccessForTenant(topicName.getTenant()); internalCreatePartitionedTopic(asyncResponse, numPartitions); } catch (Exception e) { - asyncResponse.resume(e); + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); + resumeAsyncResponseExceptionally(asyncResponse, e); } } 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 2794fce04d5c9..e0da4da9195b7 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 @@ -210,8 +210,8 @@ public void createPartitionedTopic( validateAdminAccessForTenant(topicName.getTenant()); internalCreatePartitionedTopic(asyncResponse, numPartitions); } catch (Exception e) { - log.error("Unexpected exception on create partitioned topic.", e); - asyncResponse.resume(e); + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, e); + resumeAsyncResponseExceptionally(asyncResponse, e); } } From 99f5486f2577571ec6ae1ba281c6b52553e737c1 Mon Sep 17 00:00:00 2001 From: penghui Date: Thu, 5 Mar 2020 23:15:02 +0800 Subject: [PATCH 5/7] Fix log --- .../org/apache/pulsar/broker/admin/v1/PersistentTopics.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) 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 417510bfcb153..362adc81c4a52 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 @@ -57,7 +57,8 @@ import io.swagger.annotations.ApiOperation; import io.swagger.annotations.ApiResponse; import io.swagger.annotations.ApiResponses; -import javax.ws.rs.container.AsyncResponse; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** */ @@ -66,7 +67,7 @@ @Api(value = "/persistent", description = "Persistent topic admin apis", tags = "persistent topic", hidden = true) @SuppressWarnings("deprecation") public class PersistentTopics extends PersistentTopicsBase { - + private static final Logger log = LoggerFactory.getLogger(PersistentTopics.class); @GET @Path("/{property}/{cluster}/{namespace}") @ApiOperation(hidden = true, value = "Get the list of topics under a namespace.", response = String.class, responseContainer = "List") From dd0077c44bdb5544e33e1b2dbddb525290dfcf8a Mon Sep 17 00:00:00 2001 From: penghui Date: Fri, 6 Mar 2020 10:39:27 +0800 Subject: [PATCH 6/7] Fix comments --- .../pulsar/broker/admin/AdminResource.java | 46 ++++++++++++------- .../admin/impl/PersistentTopicsBase.java | 21 +++++++-- .../broker/admin/v2/PersistentTopics.java | 11 +++-- 3 files changed, 53 insertions(+), 25 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 6d120ded3c666..32421d9c5ed93 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 @@ -27,6 +27,7 @@ import java.net.MalformedURLException; import java.net.URI; +import java.util.ArrayList; import java.util.List; import java.util.Set; import java.util.concurrent.CompletableFuture; @@ -257,16 +258,19 @@ protected List getListOfNamespaces(String property) throws Exception { return namespaces; } - protected void tryCreatePartitionsAsync(int numPartitions) { + protected CompletableFuture tryCreatePartitionsAsync(int numPartitions) { if (!topicName.isPersistent()) { - return; + return CompletableFuture.completedFuture(null); } + List> futures = new ArrayList<>(numPartitions); for (int i = 0; i < numPartitions; i++) { - tryCreatePartitionAsync(i); + futures.add(tryCreatePartitionAsync(i, null)); } + return FutureUtil.waitForAll(futures); } - private void tryCreatePartitionAsync(final int partition) { + private CompletableFuture tryCreatePartitionAsync(final int partition, CompletableFuture reuseFuture) { + CompletableFuture result = reuseFuture == null ? new CompletableFuture<>() : reuseFuture; zkCreateOptimisticAsync(localZk(), ZkAdminPaths.managedLedgerPath(topicName.getPartition(partition)), new byte[0], (rc, s, o, s1) -> { if (KeeperException.Code.OK.intValue() == rc) { @@ -274,18 +278,22 @@ private void tryCreatePartitionAsync(final int partition) { log.debug("[{}] Topic partition {} created.", clientAppId(), topicName.getPartition(partition)); } + result.complete(null); } else if (KeeperException.Code.NODEEXISTS.intValue() == rc) { log.info("[{}] Topic partition {} is exists, doing nothing.", clientAppId(), topicName.getPartition(partition)); + result.completeExceptionally(KeeperException.create(KeeperException.Code.NODEEXISTS)); } else if (KeeperException.Code.BADVERSION.intValue() == rc) { log.warn("[{}] Fail to create topic partition {} with concurrent modification, retry now.", clientAppId(), topicName.getPartition(partition)); - tryCreatePartitionAsync(partition); + tryCreatePartitionAsync(partition, result); } else { log.error("[{}] Fail to create topic partition {}", clientAppId(), topicName.getPartition(partition), KeeperException.create(KeeperException.Code.get(rc))); + result.completeExceptionally(KeeperException.create(KeeperException.Code.get(rc))); } }); + return result; } protected NamespaceName namespaceName; @@ -732,18 +740,22 @@ protected void internalCreatePartitionedTopic(AsyncResponse asyncResponse, int n byte[] data = jsonMapper().writeValueAsBytes(new PartitionedTopicMetadata(numPartitions)); zkCreateOptimisticAsync(globalZk(), path, data, (rc, s, o, s1) -> { if (KeeperException.Code.OK.intValue() == rc) { - if (topicName.isPersistent()) { - tryCreatePartitionsAsync(numPartitions); - } - globalZk().sync(path, (rc2, s2, ctx) -> { - if (KeeperException.Code.OK.intValue() == rc2) { - log.info("[{}] Successfully created partitioned topic {}", clientAppId(), topicName); - asyncResponse.resume(Response.noContent().build()); - } else { - log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, KeeperException.create(KeeperException.Code.get(rc2))); - asyncResponse.resume(new RestException(KeeperException.create(KeeperException.Code.get(rc2)))); - } - }, null); + tryCreatePartitionsAsync(numPartitions).thenAccept(v -> { + globalZk().sync(path, (rc2, s2, ctx) -> { + if (KeeperException.Code.OK.intValue() == rc2) { + log.info("[{}] Successfully created partitioned topic {}", clientAppId(), topicName); + asyncResponse.resume(Response.noContent().build()); + } else { + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, KeeperException.create(KeeperException.Code.get(rc2))); + asyncResponse.resume(new RestException(KeeperException.create(KeeperException.Code.get(rc2)))); + } + }, null); + }).exceptionally(e -> { + log.error("[{}] Failed to create partitions for topic {}", clientAppId(), topicName); + // The partitioned topic is created but there are some partitions create failed + asyncResponse.resume(new RestException(e)); + return null; + }); } else if (KeeperException.Code.NODEEXISTS.intValue() == rc) { log.warn("[{}] Failed to create already existing partitioned topic {}", clientAppId(), topicName); asyncResponse.resume(new RestException(Status.CONFLICT, "Partitioned topic already exists")); 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 d265ee8c20965..10b797b12d825 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 @@ -500,11 +500,22 @@ protected void internalUpdatePartitionedTopic(int numPartitions, boolean updateL } } - protected void internalCreateMissedPartitions() { - PartitionedTopicMetadata metadata = getPartitionedTopicMetadata(topicName, false, false); - if (metadata != null) { - tryCreatePartitionsAsync(metadata.partitions); - } + protected void internalCreateMissedPartitions(AsyncResponse asyncResponse) { + getPartitionedTopicMetadataAsync(topicName, false, false).thenAccept(metadata -> { + if (metadata != null) { + tryCreatePartitionsAsync(metadata.partitions).thenAccept(v -> { + asyncResponse.resume(Response.noContent().build()); + }).exceptionally(e -> { + log.error("[{}] Failed to create partitions for topic {}", clientAppId(), topicName); + resumeAsyncResponseExceptionally(asyncResponse, e); + return null; + }); + } + }).exceptionally(e -> { + log.error("[{}] Failed to create partitions for topic {}", clientAppId(), topicName); + resumeAsyncResponseExceptionally(asyncResponse, e); + return null; + }); } private CompletableFuture updatePartitionInOtherCluster(int numPartitions, Set clusters) { 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 e0da4da9195b7..b2fc28b010a69 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 @@ -286,7 +286,7 @@ public void updatePartitionedTopic( @POST @Path("/{tenant}/{namespace}/{topic}/createMissedPartitions") - @ApiOperation(value = "Create missed partitions of an existing partitioned topic.", notes = "This is a best-effort operation for create missed partitions of existing non-global partitioned-topic and does't throw any exceptions when create failed") + @ApiOperation(value = "Create missed partitions of an existing partitioned topic.") @ApiResponses(value = { @ApiResponse(code = 307, message = "Current broker doesn't serve the namespace of this topic"), @ApiResponse(code = 401, message = "Don't have permission to adminisActions to be grantedtrate resources on this tenant"), @@ -297,6 +297,7 @@ public void updatePartitionedTopic( @ApiResponse(code = 500, message = "Internal server error") }) public void createMissedPartitions( + @Suspended final AsyncResponse asyncResponse, @ApiParam(value = "Specify the tenant", required = true) @PathParam("tenant") String tenant, @ApiParam(value = "Specify the namespace", required = true) @@ -304,8 +305,12 @@ public void createMissedPartitions( @ApiParam(value = "Specify topic name", required = true) @PathParam("topic") @Encoded String encodedTopic) { - validatePartitionedTopicName(tenant, namespace, encodedTopic); - internalCreateMissedPartitions(); + try { + validatePartitionedTopicName(tenant, namespace, encodedTopic); + internalCreateMissedPartitions(asyncResponse); + } catch (Exception e) { + resumeAsyncResponseExceptionally(asyncResponse, e); + } } @GET From a3f4470502ade725a03526662062147c0cd3a7d8 Mon Sep 17 00:00:00 2001 From: penghui Date: Fri, 6 Mar 2020 10:43:44 +0800 Subject: [PATCH 7/7] sync global zk first --- .../pulsar/broker/admin/AdminResource.java | 31 ++++++++++--------- 1 file changed, 16 insertions(+), 15 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 32421d9c5ed93..a21698298488b 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 @@ -740,22 +740,23 @@ protected void internalCreatePartitionedTopic(AsyncResponse asyncResponse, int n byte[] data = jsonMapper().writeValueAsBytes(new PartitionedTopicMetadata(numPartitions)); zkCreateOptimisticAsync(globalZk(), path, data, (rc, s, o, s1) -> { if (KeeperException.Code.OK.intValue() == rc) { - tryCreatePartitionsAsync(numPartitions).thenAccept(v -> { - globalZk().sync(path, (rc2, s2, ctx) -> { - if (KeeperException.Code.OK.intValue() == rc2) { - log.info("[{}] Successfully created partitioned topic {}", clientAppId(), topicName); + globalZk().sync(path, (rc2, s2, ctx) -> { + if (KeeperException.Code.OK.intValue() == rc2) { + log.info("[{}] Successfully created partitioned topic {}", clientAppId(), topicName); + tryCreatePartitionsAsync(numPartitions).thenAccept(v -> { + log.info("[{}] Successfully created partitions for topic {}", clientAppId(), topicName); asyncResponse.resume(Response.noContent().build()); - } else { - log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, KeeperException.create(KeeperException.Code.get(rc2))); - asyncResponse.resume(new RestException(KeeperException.create(KeeperException.Code.get(rc2)))); - } - }, null); - }).exceptionally(e -> { - log.error("[{}] Failed to create partitions for topic {}", clientAppId(), topicName); - // The partitioned topic is created but there are some partitions create failed - asyncResponse.resume(new RestException(e)); - return null; - }); + }).exceptionally(e -> { + log.error("[{}] Failed to create partitions for topic {}", clientAppId(), topicName); + // The partitioned topic is created but there are some partitions create failed + asyncResponse.resume(new RestException(e)); + return null; + }); + } else { + log.error("[{}] Failed to create partitioned topic {}", clientAppId(), topicName, KeeperException.create(KeeperException.Code.get(rc2))); + asyncResponse.resume(new RestException(KeeperException.create(KeeperException.Code.get(rc2)))); + } + }, null); } else if (KeeperException.Code.NODEEXISTS.intValue() == rc) { log.warn("[{}] Failed to create already existing partitioned topic {}", clientAppId(), topicName); asyncResponse.resume(new RestException(Status.CONFLICT, "Partitioned topic already exists"));