From 7dd9e0e76826e8cde6cc5f75394a2de1d75470d4 Mon Sep 17 00:00:00 2001 From: zhanghao1 Date: Fri, 14 Aug 2020 11:16:18 +0800 Subject: [PATCH 1/2] [ISSUE 7757] Support Persistence Policy on topic level. --- .../pulsar/broker/admin/AdminResource.java | 24 +++- .../broker/admin/impl/NamespacesBase.java | 18 --- .../admin/impl/PersistentTopicsBase.java | 43 +++++++ .../broker/admin/v2/PersistentTopics.java | 105 ++++++++++++++++-- .../admin/TopicPoliciesDisableTest.java | 25 ++++- .../broker/admin/TopicPoliciesTest.java | 78 +++++++++++++ .../apache/pulsar/client/admin/Topics.java | 51 ++++++++- .../client/admin/internal/TopicsImpl.java | 81 +++++++++++++- .../apache/pulsar/admin/cli/CmdTopics.java | 57 ++++++++++ 9 files changed, 440 insertions(+), 42 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 55f79571b1a31..ebb2f30babbc0 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 @@ -23,7 +23,6 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.collect.Lists; - import java.net.MalformedURLException; import java.net.URI; import java.util.ArrayList; @@ -35,14 +34,12 @@ import java.util.concurrent.ExecutionException; import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.Collectors; - 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; - import org.apache.bookkeeper.util.ZkUtils; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; @@ -65,6 +62,7 @@ import org.apache.pulsar.common.policies.data.DispatchRate; import org.apache.pulsar.common.policies.data.FailureDomain; import org.apache.pulsar.common.policies.data.LocalPolicies; +import org.apache.pulsar.common.policies.data.PersistencePolicies; import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.policies.data.SubscribeRate; @@ -76,8 +74,8 @@ import org.apache.pulsar.common.util.ObjectMapperFactory; import org.apache.pulsar.zookeeper.ZooKeeperCache; import org.apache.pulsar.zookeeper.ZooKeeperChildrenCache; -import org.apache.pulsar.zookeeper.ZooKeeperManagedLedgerCache; import org.apache.pulsar.zookeeper.ZooKeeperDataCache; +import org.apache.pulsar.zookeeper.ZooKeeperManagedLedgerCache; import org.apache.zookeeper.AsyncCallback; import org.apache.zookeeper.CreateMode; import org.apache.zookeeper.KeeperException; @@ -894,4 +892,22 @@ protected void checkArgument(boolean b, String errorMessage) { throw new RestException(Status.BAD_REQUEST, errorMessage); } } + + protected void validatePersistencePolicies(PersistencePolicies persistence) { + checkNotNull(persistence, "persistence policies should not be null"); + final ServiceConfiguration config = pulsar().getConfiguration(); + checkArgument(persistence.getBookkeeperEnsemble() <= config.getManagedLedgerMaxEnsembleSize(), + "Bookkeeper-Ensemble must be <= " + config.getManagedLedgerMaxEnsembleSize()); + checkArgument(persistence.getBookkeeperWriteQuorum() <= config.getManagedLedgerMaxWriteQuorum(), + "Bookkeeper-WriteQuorum must be <= " + config.getManagedLedgerMaxWriteQuorum()); + checkArgument(persistence.getBookkeeperAckQuorum() <= config.getManagedLedgerMaxAckQuorum(), + "Bookkeeper-AckQuorum must be <= " + config.getManagedLedgerMaxAckQuorum()); + checkArgument( + (persistence.getBookkeeperEnsemble() >= persistence.getBookkeeperWriteQuorum()) + && (persistence.getBookkeeperWriteQuorum() >= persistence.getBookkeeperAckQuorum()), + String.format("Bookkeeper Ensemble (%s) >= WriteQuorum (%s) >= AckQuoru (%s)", + persistence.getBookkeeperEnsemble(), persistence.getBookkeeperWriteQuorum(), + persistence.getBookkeeperAckQuorum())); + + } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java index dc5b0afa12723..a1d7e3fe4b625 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java @@ -2005,24 +2005,6 @@ protected List internalGetAntiAffinityNamespaces(String cluster, String } } - private void validatePersistencePolicies(PersistencePolicies persistence) { - checkNotNull(persistence, "persistence policies should not be null"); - final ServiceConfiguration config = pulsar().getConfiguration(); - checkArgument(persistence.getBookkeeperEnsemble() <= config.getManagedLedgerMaxEnsembleSize(), - "Bookkeeper-Ensemble must be <= " + config.getManagedLedgerMaxEnsembleSize()); - checkArgument(persistence.getBookkeeperWriteQuorum() <= config.getManagedLedgerMaxWriteQuorum(), - "Bookkeeper-WriteQuorum must be <= " + config.getManagedLedgerMaxWriteQuorum()); - checkArgument(persistence.getBookkeeperAckQuorum() <= config.getManagedLedgerMaxAckQuorum(), - "Bookkeeper-AckQuorum must be <= " + config.getManagedLedgerMaxAckQuorum()); - checkArgument( - (persistence.getBookkeeperEnsemble() >= persistence.getBookkeeperWriteQuorum()) - && (persistence.getBookkeeperWriteQuorum() >= persistence.getBookkeeperAckQuorum()), - String.format("Bookkeeper Ensemble (%s) >= WriteQuorum (%s) >= AckQuoru (%s)", - persistence.getBookkeeperEnsemble(), persistence.getBookkeeperWriteQuorum(), - persistence.getBookkeeperAckQuorum())); - - } - protected RetentionPolicies internalGetRetention() { validateNamespacePolicyOperation(namespaceName, PolicyName.RETENTION, PolicyOperation.READ); 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 520a18153b413..5ba65b58703c4 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 @@ -95,6 +95,7 @@ import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.common.allocator.PulsarByteBufAllocator; import org.apache.pulsar.common.naming.PartitionedManagedLedgerInfo; +import org.apache.pulsar.common.policies.data.PersistencePolicies; import org.apache.pulsar.common.policies.data.PolicyName; import org.apache.pulsar.common.policies.data.PolicyOperation; import org.apache.pulsar.common.protocol.Commands; @@ -2253,6 +2254,48 @@ protected CompletableFuture internalRemoveRetention() { return pulsar().getTopicPoliciesService().updateTopicPoliciesAsync(topicName, topicPolicies.get()); } + protected void internalGetPersistence(AsyncResponse asyncResponse){ + validateAdminAccessForTenant(namespaceName.getTenant()); + validatePoliciesReadOnlyAccess(); + if (topicName.isGlobal()) { + validateGlobalNamespaceOwnership(namespaceName); + } + Optional retention = getTopicPolicies(topicName) + .map(TopicPolicies::getPersistence); + if (!retention.isPresent()) { + asyncResponse.resume(Response.noContent().build()); + }else { + asyncResponse.resume(retention.get()); + } + } + + protected CompletableFuture internalSetPersistence(PersistencePolicies persistencePolicies) { + validateAdminAccessForTenant(namespaceName.getTenant()); + validatePoliciesReadOnlyAccess(); + if (topicName.isGlobal()) { + validateGlobalNamespaceOwnership(namespaceName); + } + validatePersistencePolicies(persistencePolicies); + + TopicPolicies topicPolicies = getTopicPolicies(topicName).orElseGet(TopicPolicies::new); + topicPolicies.setPersistence(persistencePolicies); + return pulsar().getTopicPoliciesService().updateTopicPoliciesAsync(topicName, topicPolicies); + } + + protected CompletableFuture internalRemovePersistence() { + validateAdminAccessForTenant(namespaceName.getTenant()); + validatePoliciesReadOnlyAccess(); + if (topicName.isGlobal()) { + validateGlobalNamespaceOwnership(namespaceName); + } + Optional topicPolicies = getTopicPolicies(topicName); + if (!topicPolicies.isPresent()) { + return CompletableFuture.completedFuture(null); + } + topicPolicies.get().setPersistence(null); + return pulsar().getTopicPoliciesService().updateTopicPoliciesAsync(topicName, topicPolicies.get()); + } + protected MessageId internalTerminate(boolean authoritative) { if (topicName.isGlobal()) { validateGlobalNamespaceOwnership(namespaceName); 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 3b5b07f718a54..d4d32cff6d59d 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 @@ -18,11 +18,19 @@ */ package org.apache.pulsar.broker.admin.v2; +import static org.apache.pulsar.common.util.Codec.decode; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.google.common.collect.Maps; +import io.swagger.annotations.Api; +import io.swagger.annotations.ApiOperation; +import io.swagger.annotations.ApiParam; +import io.swagger.annotations.ApiResponse; +import io.swagger.annotations.ApiResponses; import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Set; - import javax.ws.rs.DELETE; import javax.ws.rs.DefaultValue; import javax.ws.rs.Encoded; @@ -38,9 +46,6 @@ import javax.ws.rs.container.Suspended; import javax.ws.rs.core.MediaType; import javax.ws.rs.core.Response; - -import com.fasterxml.jackson.core.JsonProcessingException; -import com.google.common.collect.Maps; import org.apache.pulsar.broker.admin.impl.PersistentTopicsBase; import org.apache.pulsar.broker.web.RestException; import org.apache.pulsar.client.admin.LongRunningProcessStatus; @@ -50,22 +55,15 @@ import org.apache.pulsar.common.partition.PartitionedTopicMetadata; import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.BacklogQuota; +import org.apache.pulsar.common.policies.data.PersistencePolicies; import org.apache.pulsar.common.policies.data.PersistentOfflineTopicStats; import org.apache.pulsar.common.policies.data.PersistentTopicInternalStats; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.policies.data.TopicPolicies; import org.apache.pulsar.common.policies.data.TopicStats; - -import io.swagger.annotations.Api; -import io.swagger.annotations.ApiOperation; -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; - /** */ @Path("/persistent") @@ -1168,6 +1166,89 @@ public void removeRetention(@Suspended final AsyncResponse asyncResponse, }); } + @GET + @Path("/{tenant}/{namespace}/{topic}/persistence") + @ApiOperation(value = "Get configuration of persistence policies for specified topic.") + @ApiResponses(value = {@ApiResponse(code = 403, message = "Don't have admin permission"), + @ApiResponse(code = 404, message = "Topic does not exist"), + @ApiResponse(code = 405, message = "Topic level policy is disabled, to enable the topic level policy and retry"), + @ApiResponse(code = 409, message = "Concurrent modification")}) + public void getPersistence(@Suspended final AsyncResponse asyncResponse, + @PathParam("tenant") String tenant, + @PathParam("namespace") String namespace, + @PathParam("topic") @Encoded String encodedTopic) { + validateTopicName(tenant, namespace, encodedTopic); + try { + internalGetPersistence(asyncResponse); + } catch (RestException e) { + asyncResponse.resume(e); + } catch (Exception e) { + asyncResponse.resume(new RestException(e)); + } + } + + @POST + @Path("/{tenant}/{namespace}/{topic}/persistence") + @ApiOperation(value = "Set configuration of persistence policies for specified topic.") + @ApiResponses(value = {@ApiResponse(code = 403, message = "Don't have admin permission"), + @ApiResponse(code = 404, message = "Topic does not exist"), + @ApiResponse(code = 405, message = "Topic level policy is disabled, to enable the topic level policy and retry"), + @ApiResponse(code = 409, message = "Concurrent modification"), + @ApiResponse(code = 400, message = "Invalid persistence policies")}) + public void setPersistence(@Suspended final AsyncResponse asyncResponse, + @PathParam("tenant") String tenant, + @PathParam("namespace") String namespace, + @PathParam("topic") @Encoded String encodedTopic, + @ApiParam(value = "Bookkeeper persistence policies for specified topic") PersistencePolicies persistencePolicies) { + validateTopicName(tenant, namespace, encodedTopic); + internalSetPersistence(persistencePolicies).whenComplete((r, ex) -> { + if (ex instanceof RestException) { + log.error("Failed updated persistence policies", ex); + asyncResponse.resume(ex); + } else if (ex != null) { + log.error("Failed updated persistence policies", ex); + asyncResponse.resume(new RestException(ex)); + } else { + try { + log.info("[{}] Successfully updated persistence policies: namespace={}, topic={}, persistencePolicies={}", + clientAppId(), + namespaceName, + topicName.getLocalName(), + jsonMapper().writeValueAsString(persistencePolicies)); + } catch (JsonProcessingException ignore) { + } + asyncResponse.resume(Response.noContent().build()); + } + }); + } + + @DELETE + @Path("/{tenant}/{namespace}/{topic}/persistence") + @ApiOperation(value = "Remove configuration of persistence policies for specified topic.") + @ApiResponses(value = {@ApiResponse(code = 403, message = "Don't have admin permission"), + @ApiResponse(code = 404, message = "Topic does not exist"), + @ApiResponse(code = 405, message = "Topic level policy is disabled, to enable the topic level policy and retry"), + @ApiResponse(code = 409, message = "Concurrent modification")}) + public void removePersistence(@Suspended final AsyncResponse asyncResponse, + @PathParam("tenant") String tenant, + @PathParam("namespace") String namespace, + @PathParam("topic") @Encoded String encodedTopic) { + validateTopicName(tenant, namespace, encodedTopic); + internalRemovePersistence().whenComplete((r, ex) -> { + if (ex != null) { + log.error("Failed updated retention", ex); + asyncResponse.resume(new RestException(ex)); + } else { + log.info("[{}] Successfully remove persistence policies: namespace={}, topic={}", + clientAppId(), + namespaceName, + topicName.getLocalName()); + asyncResponse.resume(Response.noContent().build()); + } + }); + } + + @POST @Path("/{tenant}/{namespace}/{topic}/terminate") @ApiOperation(value = "Terminate a topic. A topic that is terminated will not accept any more " diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesDisableTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesDisableTest.java index e473382584d19..819cb5cf8f8e2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesDisableTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesDisableTest.java @@ -24,6 +24,7 @@ import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.ClusterData; +import org.apache.pulsar.common.policies.data.PersistencePolicies; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.Assert; @@ -63,7 +64,7 @@ public void cleanup() throws Exception { } @Test - public void testBacklogQuotaDisabled() throws Exception { + public void testBacklogQuotaDisabled() { BacklogQuota backlogQuota = new BacklogQuota(1024, BacklogQuota.RetentionPolicy.consumer_backlog_eviction); log.info("Backlog quota: {} will set to the topic: {}", backlogQuota, testTopic); @@ -90,7 +91,7 @@ public void testBacklogQuotaDisabled() throws Exception { } @Test - public void testRetentionDisabled() throws Exception { + public void testRetentionDisabled() { RetentionPolicies retention = new RetentionPolicies(); log.info("Retention: {} will set to the topic: {}", retention, testTopic); @@ -108,4 +109,24 @@ public void testRetentionDisabled() throws Exception { Assert.assertEquals(e.getStatusCode(), 405); } } + + @Test + public void testPersistenceDisabled() { + PersistencePolicies persistencePolicies = new PersistencePolicies(); + log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, testTopic); + + try { + admin.topics().setPersistence(testTopic, persistencePolicies); + Assert.fail(); + } catch (PulsarAdminException e) { + Assert.assertEquals(e.getStatusCode(), 405); + } + + try { + admin.topics().getPersistence(testTopic); + Assert.fail(); + } catch (PulsarAdminException e) { + Assert.assertEquals(e.getStatusCode(), 405); + } + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java index 8e5aa5118c26c..d306b010bc1bf 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java @@ -27,6 +27,7 @@ import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.ClusterData; +import org.apache.pulsar.common.policies.data.PersistencePolicies; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.policies.data.TenantInfo; import org.testng.Assert; @@ -237,4 +238,81 @@ public void testRemoveRetention() throws Exception { admin.topics().deletePartitionedTopic(testTopic, true); } + + @Test + public void testCheckPersistence() throws Exception { + PersistencePolicies persistencePolicies = new PersistencePolicies(6, 2, 2, 0.0); + log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, testTopic); + try { + admin.topics().setPersistence(testTopic, persistencePolicies); + Assert.fail(); + } catch (PulsarAdminException e) { + Assert.assertEquals(e.getStatusCode(), 400); + } + + persistencePolicies = new PersistencePolicies(2, 6, 2, 0.0); + log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, testTopic); + try { + admin.topics().setPersistence(testTopic, persistencePolicies); + Assert.fail(); + } catch (PulsarAdminException e) { + Assert.assertEquals(e.getStatusCode(), 400); + } + + persistencePolicies = new PersistencePolicies(2, 2, 6, 0.0); + log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, testTopic); + try { + admin.topics().setPersistence(testTopic, persistencePolicies); + Assert.fail(); + } catch (PulsarAdminException e) { + Assert.assertEquals(e.getStatusCode(), 400); + } + + persistencePolicies = new PersistencePolicies(1, 2, 2, 0.0); + log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, testTopic); + try { + admin.topics().setPersistence(testTopic, persistencePolicies); + Assert.fail(); + } catch (PulsarAdminException e) { + Assert.assertEquals(e.getStatusCode(), 400); + } + + admin.topics().deletePartitionedTopic(testTopic, true); + } + + @Test + public void testSetPersistence() throws Exception { + PersistencePolicies persistencePolicies = new PersistencePolicies(2, 2, 2, 0.0); + log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, testTopic); + + admin.topics().setPersistence(testTopic, persistencePolicies); + Thread.sleep(3000); + PersistencePolicies getPersistencePolicies = admin.topics().getPersistence(testTopic); + log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, testTopic); + Assert.assertEquals(getPersistencePolicies, persistencePolicies); + + admin.topics().deletePartitionedTopic(testTopic, true); + } + + @Test + public void testRemovePersistence() throws Exception { + + PersistencePolicies persistencePolicies = new PersistencePolicies(2, 2, 2, 0.0); + log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, testTopic); + + admin.topics().setPersistence(testTopic, persistencePolicies); + Thread.sleep(3000); + PersistencePolicies getPersistencePolicies = admin.topics().getPersistence(testTopic); + + log.info("PersistencePolicies {} get on topic: {}", getPersistencePolicies, testTopic); + Assert.assertEquals(getPersistencePolicies, persistencePolicies); + + admin.topics().removePersistence(testTopic); + Thread.sleep(3000); + log.info("PersistencePolicies {} get on topic: {} after remove", getPersistencePolicies, testTopic); + getPersistencePolicies = admin.topics().getPersistence(testTopic); + Assert.assertNull(getPersistencePolicies); + + admin.topics().deletePartitionedTopic(testTopic, true); + } } diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java index 72098ee64056a..c49c395a6ce94 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Topics.java @@ -19,12 +19,10 @@ package org.apache.pulsar.client.admin; import com.google.gson.JsonObject; - import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.CompletableFuture; - import org.apache.pulsar.client.admin.PulsarAdminException.ConflictException; import org.apache.pulsar.client.admin.PulsarAdminException.NotAllowedException; import org.apache.pulsar.client.admin.PulsarAdminException.NotAuthorizedException; @@ -37,6 +35,7 @@ import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.PartitionedTopicInternalStats; import org.apache.pulsar.common.policies.data.PartitionedTopicStats; +import org.apache.pulsar.common.policies.data.PersistencePolicies; import org.apache.pulsar.common.policies.data.PersistentTopicInternalStats; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.policies.data.TopicStats; @@ -1615,4 +1614,52 @@ void createSubscription(String topic, String subscriptionName, MessageId message * Topic name */ CompletableFuture removeRetentionAsync(String topic); + + /** + * Set the configuration of persistence policies for specified topic. + * + * @param topic Topic name + * @param persistencePolicies Configuration of bookkeeper persistence policies + * @throws PulsarAdminException Unexpected error + */ + void setPersistence(String topic, PersistencePolicies persistencePolicies) throws PulsarAdminException; + + /** + * Set the configuration of persistence policies for specified topic asynchronously. + * + * @param topic Topic name + * @param persistencePolicies Configuration of bookkeeper persistence policies + */ + CompletableFuture setPersistenceAsync(String topic, PersistencePolicies persistencePolicies); + + /** + * Get the configuration of persistence policies for specified topic. + * + * @param topic Topic name + * @return Configuration of bookkeeper persistence policies + * @throws PulsarAdminException Unexpected error + */ + PersistencePolicies getPersistence(String topic) throws PulsarAdminException; + + /** + * Get the configuration of persistence policies for specified topic asynchronously. + * + * @param topic Topic name + */ + CompletableFuture getPersistenceAsync(String topic); + + /** + * Remove the configuration of persistence policies for specified topic. + * + * @param topic Topic name + * @throws PulsarAdminException Unexpected error + */ + void removePersistence(String topic) throws PulsarAdminException; + + /** + * Remove the configuration of persistence policies for specified topic asynchronously. + * + * @param topic Topic name + */ + CompletableFuture removePersistenceAsync(String topic); } diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java index 4df515af9bc99..83b1ed7bf4e03 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java @@ -24,10 +24,8 @@ import com.google.common.collect.Maps; import com.google.gson.Gson; import com.google.gson.JsonObject; - import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; - import java.io.InputStream; import java.util.ArrayList; import java.util.Collections; @@ -41,7 +39,6 @@ import java.util.concurrent.TimeoutException; import java.util.stream.Collectors; import java.util.stream.Stream; - import javax.ws.rs.client.Entity; import javax.ws.rs.client.InvocationCallback; import javax.ws.rs.client.WebTarget; @@ -50,7 +47,6 @@ import javax.ws.rs.core.MultivaluedMap; import javax.ws.rs.core.Response; import javax.ws.rs.core.Response.Status; - import org.apache.pulsar.client.admin.LongRunningProcessStatus; import org.apache.pulsar.client.admin.OffloadProcessStatus; import org.apache.pulsar.client.admin.PulsarAdminException; @@ -75,6 +71,7 @@ import org.apache.pulsar.common.policies.data.ErrorData; import org.apache.pulsar.common.policies.data.PartitionedTopicInternalStats; import org.apache.pulsar.common.policies.data.PartitionedTopicStats; +import org.apache.pulsar.common.policies.data.PersistencePolicies; import org.apache.pulsar.common.policies.data.PersistentTopicInternalStats; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.policies.data.TopicStats; @@ -1540,5 +1537,81 @@ public CompletableFuture removeRetentionAsync(String topic) { return asyncDeleteRequest(path); } + @Override + public void setPersistence(String topic, PersistencePolicies persistencePolicies) throws PulsarAdminException { + try { + setPersistenceAsync(topic, persistencePolicies).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); + } catch (ExecutionException e) { + throw (PulsarAdminException) e.getCause(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); + } + } + + @Override + public CompletableFuture setPersistenceAsync(String topic, PersistencePolicies persistencePolicies) { + TopicName tn = validateTopic(topic); + WebTarget path = topicPath(tn, "persistence"); + return asyncPostRequest(path, Entity.entity(persistencePolicies, MediaType.APPLICATION_JSON)); + } + + @Override + public PersistencePolicies getPersistence(String topic) throws PulsarAdminException { + try { + return getPersistenceAsync(topic).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); + } catch (ExecutionException e) { + throw (PulsarAdminException) e.getCause(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); + } + } + + @Override + public CompletableFuture getPersistenceAsync(String topic) { + TopicName tn = validateTopic(topic); + WebTarget path = topicPath(tn, "persistence"); + final CompletableFuture future = new CompletableFuture<>(); + asyncGetRequest(path, + new InvocationCallback() { + @Override + public void completed(PersistencePolicies persistencePolicies) { + future.complete(persistencePolicies); + } + + @Override + public void failed(Throwable throwable) { + future.completeExceptionally(getApiException(throwable.getCause())); + } + }); + return future; + } + + @Override + public void removePersistence(String topic) throws PulsarAdminException { + try { + removePersistenceAsync(topic).get(this.readTimeoutMs, TimeUnit.MILLISECONDS); + } catch (ExecutionException e) { + throw (PulsarAdminException) e.getCause(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new PulsarAdminException(e); + } catch (TimeoutException e) { + throw new PulsarAdminException.TimeoutException(e); + } + } + + @Override + public CompletableFuture removePersistenceAsync(String topic) { + TopicName tn = validateTopic(topic); + WebTarget path = topicPath(tn, "persistence"); + return asyncDeleteRequest(path); + } + private static final Logger log = LoggerFactory.getLogger(TopicsImpl.class); } diff --git a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java index 658233f3448b3..7ef1ab984b123 100644 --- a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java +++ b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.java @@ -48,6 +48,7 @@ import org.apache.pulsar.client.impl.BatchMessageIdImpl; import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.common.policies.data.BacklogQuota; +import org.apache.pulsar.common.policies.data.PersistencePolicies; import org.apache.pulsar.common.policies.data.PersistentTopicInternalStats; import org.apache.pulsar.common.policies.data.RetentionPolicies; import org.apache.pulsar.common.util.RelativeTimeUtil; @@ -111,6 +112,9 @@ public CmdTopics(PulsarAdmin admin) { jcommander.addCommand("get-retention", new GetRetention()); jcommander.addCommand("set-retention", new SetRetention()); jcommander.addCommand("remove-retention", new RemoveRetention()); + jcommander.addCommand("get-persistence", new GetPersistence()); + jcommander.addCommand("set-persistence", new SetPersistence()); + jcommander.addCommand("remove-persistence", new RemovePersistence()); } @Parameters(commandDescription = "Get the list of topics under a namespace.") @@ -992,4 +996,57 @@ void run() throws PulsarAdminException { admin.topics().removeRetention(persistentTopic); } } + + @Parameters(commandDescription = "Get the persistence policies for a topic") + private class GetPersistence extends CliCommand { + @Parameter(description = "persistent://tenant/namespace/topic", required = true) + private java.util.List params; + + @Override + void run() throws PulsarAdminException { + String persistentTopic = validatePersistentTopic(params); + print(admin.topics().getPersistence(persistentTopic)); + } + } + + @Parameters(commandDescription = "Set the persistence policies for a topic") + private class SetPersistence extends CliCommand { + @Parameter(description = "persistent://tenant/namespace/topic", required = true) + private java.util.List params; + + @Parameter(names = { "-e", + "--bookkeeper-ensemble" }, description = "Number of bookies to use for a topic", required = true) + private int bookkeeperEnsemble; + + @Parameter(names = { "-w", + "--bookkeeper-write-quorum" }, description = "How many writes to make of each entry", required = true) + private int bookkeeperWriteQuorum; + + @Parameter(names = { "-a", + "--bookkeeper-ack-quorum" }, description = "Number of acks (garanteed copies) to wait for each entry", required = true) + private int bookkeeperAckQuorum; + + @Parameter(names = { "-r", + "--ml-mark-delete-max-rate" }, description = "Throttling rate of mark-delete operation (0 means no throttle)", required = true) + private double managedLedgerMaxMarkDeleteRate; + + @Override + void run() throws PulsarAdminException { + String persistentTopic = validatePersistentTopic(params); + admin.topics().setPersistence(persistentTopic, new PersistencePolicies(bookkeeperEnsemble, + bookkeeperWriteQuorum, bookkeeperAckQuorum, managedLedgerMaxMarkDeleteRate)); + } + } + + @Parameters(commandDescription = "Remove the persistence policy for a topic") + private class RemovePersistence extends CliCommand { + @Parameter(description = "persistent://tenant/namespace/topic", required = true) + private java.util.List params; + + @Override + void run() throws PulsarAdminException { + String persistentTopic = validatePersistentTopic(params); + admin.topics().removePersistence(persistentTopic); + } + } } From 4600c1fb57edf68c118428c8032318f3133d728b Mon Sep 17 00:00:00 2001 From: zhanghao1 Date: Sat, 15 Aug 2020 14:45:16 +0800 Subject: [PATCH 2/2] Bug fix. --- .../admin/impl/PersistentTopicsBase.java | 6 +-- .../broker/admin/v2/PersistentTopics.java | 1 - .../pulsar/broker/service/BrokerService.java | 48 +++++++++++++------ .../broker/admin/TopicPoliciesTest.java | 32 +++++++++++-- 4 files changed, 64 insertions(+), 23 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/PersistentTopicsBase.java index 5e9241e39cc1e..650209ac0677a 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 @@ -2301,12 +2301,12 @@ protected void internalGetPersistence(AsyncResponse asyncResponse){ if (topicName.isGlobal()) { validateGlobalNamespaceOwnership(namespaceName); } - Optional retention = getTopicPolicies(topicName) + Optional persistencePolicies = getTopicPolicies(topicName) .map(TopicPolicies::getPersistence); - if (!retention.isPresent()) { + if (!persistencePolicies.isPresent()) { asyncResponse.resume(Response.noContent().build()); }else { - asyncResponse.resume(retention.get()); + asyncResponse.resume(persistencePolicies.get()); } } 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 2553d6445146a..93b0cc252a09c 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 @@ -50,7 +50,6 @@ import org.apache.pulsar.broker.web.RestException; import org.apache.pulsar.client.admin.LongRunningProcessStatus; import org.apache.pulsar.client.admin.OffloadProcessStatus; -import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.common.partition.PartitionedTopicMetadata; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index e6b1c819980eb..fa2f7ecbb2bea 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -31,7 +31,6 @@ import com.google.common.collect.Lists; import com.google.common.collect.Maps; import com.google.common.collect.Queues; - import io.netty.bootstrap.ServerBootstrap; import io.netty.buffer.ByteBuf; import io.netty.channel.AdaptiveRecvByteBufAllocator; @@ -42,7 +41,6 @@ import io.netty.channel.socket.SocketChannel; import io.netty.handler.ssl.SslContext; import io.netty.util.concurrent.DefaultThreadFactory; - import java.io.Closeable; import java.io.IOException; import java.lang.reflect.Field; @@ -70,11 +68,9 @@ import java.util.concurrent.locks.ReentrantReadWriteLock; import java.util.function.Consumer; import java.util.function.Predicate; - import lombok.AccessLevel; import lombok.Getter; import lombok.Setter; - import org.apache.bookkeeper.common.util.OrderedExecutor; import org.apache.bookkeeper.common.util.OrderedScheduler; import org.apache.bookkeeper.mledger.AsyncCallbacks.DeleteLedgerCallback; @@ -118,7 +114,6 @@ import org.apache.pulsar.broker.zookeeper.aspectj.ClientCnxnAspect.EventListner; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminBuilder; - import org.apache.pulsar.client.api.ClientBuilder; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; @@ -144,6 +139,7 @@ import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.PublishRate; import org.apache.pulsar.common.policies.data.RetentionPolicies; +import org.apache.pulsar.common.policies.data.TopicPolicies; import org.apache.pulsar.common.policies.data.TopicStats; import org.apache.pulsar.common.policies.data.TopicType; import org.apache.pulsar.common.stats.Metrics; @@ -1070,6 +1066,26 @@ public CompletableFuture getManagedLedgerConfig(TopicName t // Get persistence policy for this topic Optional policies = Optional.empty(); Optional localPolicies = Optional.empty(); + + PersistencePolicies persistencePolicies = null; + RetentionPolicies retentionPolicies = null; + + if (pulsar.getConfig().isTopicLevelPoliciesEnabled()) { + TopicName cloneTopicName = topicName; + if (topicName.isPartitioned()) { + cloneTopicName = TopicName.get(topicName.getPartitionedTopicName()); + } + try { + TopicPolicies topicPolicies = pulsar.getTopicPoliciesService().getTopicPolicies(cloneTopicName); + if (topicPolicies != null) { + persistencePolicies = topicPolicies.getPersistence(); + retentionPolicies = topicPolicies.getRetentionPolicies(); + } + } catch (BrokerServiceException.TopicPoliciesCacheNotInitException e) { + log.warn("Topic {} policies cache have not init.", topicName); + } + } + try { policies = pulsar .getConfigurationCache().policiesCache().get(AdminResource.path(POLICIES, @@ -1083,16 +1099,20 @@ public CompletableFuture getManagedLedgerConfig(TopicName t return; } - PersistencePolicies persistencePolicies = policies.map(p -> p.persistence).orElseGet( - () -> new PersistencePolicies(serviceConfig.getManagedLedgerDefaultEnsembleSize(), - serviceConfig.getManagedLedgerDefaultWriteQuorum(), - serviceConfig.getManagedLedgerDefaultAckQuorum(), - serviceConfig.getManagedLedgerDefaultMarkDeleteRateLimit())); + if (persistencePolicies == null) { + persistencePolicies = policies.map(p -> p.persistence).orElseGet( + () -> new PersistencePolicies(serviceConfig.getManagedLedgerDefaultEnsembleSize(), + serviceConfig.getManagedLedgerDefaultWriteQuorum(), + serviceConfig.getManagedLedgerDefaultAckQuorum(), + serviceConfig.getManagedLedgerDefaultMarkDeleteRateLimit())); + } - RetentionPolicies retentionPolicies = policies.map(p -> p.retention_policies).orElseGet( - () -> new RetentionPolicies(serviceConfig.getDefaultRetentionTimeInMinutes(), - serviceConfig.getDefaultRetentionSizeInMB()) - ); + if (retentionPolicies == null) { + retentionPolicies = policies.map(p -> p.retention_policies).orElseGet( + () -> new RetentionPolicies(serviceConfig.getDefaultRetentionTimeInMinutes(), + serviceConfig.getDefaultRetentionSizeInMB()) + ); + } ManagedLedgerConfig managedLedgerConfig = new ManagedLedgerConfig(); managedLedgerConfig.setEnsembleSize(persistencePolicies.getBookkeeperEnsemble()); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java index d306b010bc1bf..903ddcee9ae67 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java @@ -18,10 +18,15 @@ */ package org.apache.pulsar.broker.admin; +import static org.testng.Assert.assertEquals; + import com.google.common.collect.Sets; import lombok.extern.slf4j.Slf4j; +import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest; import org.apache.pulsar.broker.service.BacklogQuotaManager; +import org.apache.pulsar.broker.service.Topic; +import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.common.naming.TopicName; @@ -46,6 +51,8 @@ public class TopicPoliciesTest extends MockedPulsarServiceBaseTest { private final String testTopic = "persistent://" + myNamespace + "/test-set-backlog-quota"; + private final String persistenceTopic = "persistent://" + myNamespace + "/test-set-persistence"; + @BeforeMethod @Override protected void setup() throws Exception { @@ -282,15 +289,30 @@ public void testCheckPersistence() throws Exception { @Test public void testSetPersistence() throws Exception { - PersistencePolicies persistencePolicies = new PersistencePolicies(2, 2, 2, 0.0); - log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, testTopic); + PersistencePolicies persistencePolicies = new PersistencePolicies(3, 3, 3, 0.1); + log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, persistenceTopic); - admin.topics().setPersistence(testTopic, persistencePolicies); + admin.topics().setPersistence(persistenceTopic, persistencePolicies); Thread.sleep(3000); - PersistencePolicies getPersistencePolicies = admin.topics().getPersistence(testTopic); - log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, testTopic); + + admin.topics().createPartitionedTopic(persistenceTopic, 2); + Producer producer = pulsarClient.newProducer().topic(persistenceTopic).create(); + producer.close(); + + admin.lookups().lookupTopic(persistenceTopic); + Topic t = pulsar.getBrokerService().getOrCreateTopic(persistenceTopic).get(); + PersistentTopic persistentTopic = (PersistentTopic) t; + ManagedLedgerConfig managedLedgerConfig = persistentTopic.getManagedLedger().getConfig(); + assertEquals(managedLedgerConfig.getEnsembleSize(), 3); + assertEquals(managedLedgerConfig.getWriteQuorumSize(), 3); + assertEquals(managedLedgerConfig.getAckQuorumSize(), 3); + assertEquals(managedLedgerConfig.getThrottleMarkDelete(), 0.1); + + PersistencePolicies getPersistencePolicies = admin.topics().getPersistence(persistenceTopic); + log.info("PersistencePolicies: {} will set to the topic: {}", persistencePolicies, persistenceTopic); Assert.assertEquals(getPersistencePolicies, persistencePolicies); + admin.topics().deletePartitionedTopic(persistenceTopic, true); admin.topics().deletePartitionedTopic(testTopic, true); }