diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index c92b3879b9df8..091221159f739 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -756,12 +756,6 @@ public class ServiceConfiguration implements PulsarConfiguration { doc = "Enable bookie secondary-isolation group if bookkeeperClientIsolationGroups doesn't have enough bookie available." ) private String bookkeeperClientSecondaryIsolationGroups; - @FieldContext( - category = CATEGORY_STORAGE_BK, - required = false, - doc = "Minimum bookies that should be available as part of bookkeeperClientIsolationGroups \n\n" - + "else broker will include bookkeeperClientSecondaryIsolationGroups bookies in isolated list.") - private int bookkeeperClientMinAvailableBookiesInIsolationGroups = 0; @FieldContext(category = CATEGORY_STORAGE_BK, doc = "Enable/disable having read operations for a ledger to be sticky to " + "a single bookie.\n" + "If this flag is enabled, the client will use one single bookie (by " + diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java index 2105af5fe3bb5..b51c572a07baf 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/BookKeeperClientFactoryImpl.java @@ -124,8 +124,6 @@ private void setDefaultEnsemblePlacementPolicy(ClientConfiguration bkConf, Servi conf.getBookkeeperClientIsolationGroups()); bkConf.setProperty(ZkIsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, conf.getBookkeeperClientSecondaryIsolationGroups()); - bkConf.setProperty(ZkIsolatedBookieEnsemblePlacementPolicy.MIN_AVAILABLE_PRIMARY_ISOLATED_BOOKIE, - conf.getBookkeeperClientMinAvailableBookiesInIsolationGroups()); if (bkConf.getProperty(ZooKeeperCache.ZK_CACHE_INSTANCE) == null) { ZooKeeperCache zkc = new ZooKeeperCache(zkClient, conf.getZooKeeperOperationTimeoutSeconds()) { }; 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 48a3c5d9d3499..12fb2431d2fdd 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 @@ -72,6 +72,7 @@ import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.BacklogQuota.BacklogQuotaType; +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; import org.apache.pulsar.common.policies.data.BundlesData; import org.apache.pulsar.common.policies.data.ClusterData; import org.apache.pulsar.common.policies.data.DispatchRate; @@ -583,7 +584,7 @@ protected void internalUnloadNamespace() { } - protected void internalSetBookieAffinityGroup(String bookieAffinityGroup) { + protected void internalSetBookieAffinityGroup(BookieAffinityGroupData bookieAffinityGroup) { log.info("[{}] Setting bookie-affinity-group {} for namespace {}", clientAppId(), bookieAffinityGroup, this.namespaceName); @@ -614,7 +615,7 @@ protected void internalSetBookieAffinityGroup(String bookieAffinityGroup) { .get(pulsar().getConfiguration().getZooKeeperOperationTimeoutSeconds(), SECONDS); localPolicies = new LocalPolicies(); } - localPolicies.bookkeeperAffinityGroup = bookieAffinityGroup; + localPolicies.bookieAffinityGroup = bookieAffinityGroup; byte[] data = ObjectMapperFactory.getThreadLocal().writeValueAsBytes(localPolicies); pulsar().getLocalZkCache().getZooKeeper().setData(path, data, Math.toIntExact(version)); // invalidate namespace's local-policies @@ -636,7 +637,7 @@ protected void internalSetBookieAffinityGroup(String bookieAffinityGroup) { } } - protected String internalGetBookieAffinityGroup() { + protected BookieAffinityGroupData internalGetBookieAffinityGroup() { validateSuperUserAccess(); if (namespaceName.isGlobal()) { @@ -650,9 +651,9 @@ protected String internalGetBookieAffinityGroup() { String path = joinPath(LOCAL_POLICIES_ROOT, this.namespaceName.toString()); try { Optional policies = pulsar().getLocalZkCacheService().policiesCache().get(path); - final String bookkeeperAffinityGroup = policies.orElseThrow(() -> new RestException(Status.NOT_FOUND, - "Namespace local-policies does not exist")).bookkeeperAffinityGroup; - if (StringUtils.isBlank(bookkeeperAffinityGroup)) { + final BookieAffinityGroupData bookkeeperAffinityGroup = policies.orElseThrow(() -> new RestException(Status.NOT_FOUND, + "Namespace local-policies does not exist")).bookieAffinityGroup; + if (bookkeeperAffinityGroup == null) { throw new RestException(Status.NOT_FOUND, "bookie-affinity group does not exist"); } return bookkeeperAffinityGroup; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/Namespaces.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/Namespaces.java index 4dfa6ce38681d..d74feaec280f5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/Namespaces.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v1/Namespaces.java @@ -28,6 +28,7 @@ import org.apache.pulsar.common.api.proto.PulsarApi.CommandGetTopicsOfNamespace.Mode; import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.BacklogQuota; +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; import org.apache.pulsar.common.policies.data.BacklogQuota.BacklogQuotaType; import org.apache.pulsar.common.policies.data.BundlesData; import org.apache.pulsar.common.policies.data.DispatchRate; @@ -552,13 +553,13 @@ public void setPersistence(@PathParam("property") String property, @PathParam("c } @POST - @Path("/{property}/{cluster}/{namespace}/persistence/bookieAffinity/{bookieAffinityGroup}") + @Path("/{property}/{cluster}/{namespace}/persistence/bookieAffinity") @ApiOperation(hidden = true, value = "Set the bookie-affinity-group to namespace-local policy.") @ApiResponses(value = { @ApiResponse(code = 403, message = "Don't have admin permission"), @ApiResponse(code = 404, message = "Namespace does not exist"), @ApiResponse(code = 409, message = "Concurrent modification") }) public void setBookieAffinityGroup(@PathParam("property") String property, @PathParam("cluster") String cluster, - @PathParam("namespace") String namespace, @PathParam("bookieAffinityGroup") String bookieAffinityGroup) { + @PathParam("namespace") String namespace, BookieAffinityGroupData bookieAffinityGroup) { validateNamespaceName(property, cluster, namespace); internalSetBookieAffinityGroup(bookieAffinityGroup); } @@ -569,7 +570,7 @@ public void setBookieAffinityGroup(@PathParam("property") String property, @Path @ApiResponses(value = { @ApiResponse(code = 403, message = "Don't have admin permission"), @ApiResponse(code = 404, message = "Namespace does not exist"), @ApiResponse(code = 409, message = "Concurrent modification") }) - public String getBookieAffinityGroup(@PathParam("property") String property, + public BookieAffinityGroupData getBookieAffinityGroup(@PathParam("property") String property, @PathParam("cluster") String cluster, @PathParam("namespace") String namespace) { validateNamespaceName(property, cluster, namespace); return internalGetBookieAffinityGroup(); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Namespaces.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Namespaces.java index d6b9042a0ed1a..09b99281b4d97 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Namespaces.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/Namespaces.java @@ -40,6 +40,7 @@ import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.BacklogQuota.BacklogQuotaType; +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; import org.apache.pulsar.common.policies.data.BundlesData; import org.apache.pulsar.common.policies.data.DispatchRate; import org.apache.pulsar.common.policies.data.PersistencePolicies; @@ -492,13 +493,13 @@ public void setPersistence(@PathParam("tenant") String tenant, @PathParam("names } @POST - @Path("/{tenant}/{namespace}/persistence/bookieAffinity/{bookieAffinityGroup}") + @Path("/{tenant}/{namespace}/persistence/bookieAffinity") @ApiOperation(value = "Set the bookie-affinity-group to namespace-persistent policy.") @ApiResponses(value = { @ApiResponse(code = 403, message = "Don't have admin permission"), @ApiResponse(code = 404, message = "Namespace does not exist"), @ApiResponse(code = 409, message = "Concurrent modification") }) public void setBookieAffinityGroup(@PathParam("tenant") String tenant, @PathParam("namespace") String namespace, - @PathParam("bookieAffinityGroup") String bookieAffinityGroup) { + BookieAffinityGroupData bookieAffinityGroup) { validateNamespaceName(tenant, namespace); internalSetBookieAffinityGroup(bookieAffinityGroup); } @@ -509,7 +510,7 @@ public void setBookieAffinityGroup(@PathParam("tenant") String tenant, @PathPara @ApiResponses(value = { @ApiResponse(code = 403, message = "Don't have admin permission"), @ApiResponse(code = 404, message = "Namespace does not exist"), @ApiResponse(code = 409, message = "Concurrent modification") }) - public String getBookieAffinityGroup(@PathParam("property") String property, + public BookieAffinityGroupData getBookieAffinityGroup(@PathParam("property") String property, @PathParam("namespace") String namespace) { validateNamespaceName(property, namespace); return internalGetBookieAffinityGroup(); 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 1254088da95e7..077e737617f01 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 @@ -746,12 +746,14 @@ public CompletableFuture getManagedLedgerConfig(TopicName t managedLedgerConfig.setEnsembleSize(persistencePolicies.getBookkeeperEnsemble()); managedLedgerConfig.setWriteQuorumSize(persistencePolicies.getBookkeeperWriteQuorum()); managedLedgerConfig.setAckQuorumSize(persistencePolicies.getBookkeeperAckQuorum()); - if (localPolicies.isPresent() && StringUtils.isNotBlank(localPolicies.get().bookkeeperAffinityGroup)) { + if (localPolicies.isPresent() && localPolicies.get().bookieAffinityGroup != null) { managedLedgerConfig .setBookKeeperEnsemblePlacementPolicyClassName(ZkIsolatedBookieEnsemblePlacementPolicy.class); Map properties = Maps.newHashMap(); properties.put(ZkIsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, - localPolicies.get().bookkeeperAffinityGroup); + localPolicies.get().bookieAffinityGroup.bookkeeperAffinityGroupPrimary); + properties.put(ZkIsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, + localPolicies.get().bookieAffinityGroup.bookkeeperAffinityGroupSecondary); managedLedgerConfig.setBookKeeperEnsemblePlacementPolicyProperties(properties); } managedLedgerConfig.setThrottleMarkDelete(persistencePolicies.getManagedLedgerMaxMarkDeleteRate()); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBookieIsolationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBookieIsolationTest.java index 5fe3f1dadbece..2ada2ecc0946a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBookieIsolationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/BrokerBookieIsolationTest.java @@ -19,6 +19,7 @@ package org.apache.pulsar.broker.service; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.fail; import java.lang.reflect.Field; import java.net.URL; @@ -52,6 +53,8 @@ import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.ProducerBuilder; import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.PulsarClientException.BrokerPersistenceException; +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; import org.apache.pulsar.common.policies.data.BookieInfo; import org.apache.pulsar.common.policies.data.BookiesRackConfiguration; import org.apache.pulsar.common.policies.data.ClusterData; @@ -173,13 +176,19 @@ public void testBookieIsolation() throws Exception { admin.namespaces().createNamespace(ns2); admin.namespaces().createNamespace(ns3); admin.namespaces().createNamespace(ns4); - admin.namespaces().setBookieAffinityGroup(ns2, tenantNamespaceIsolationGroups); - admin.namespaces().setBookieAffinityGroup(ns3, tenantNamespaceIsolationGroups); - admin.namespaces().setBookieAffinityGroup(ns4, tenantNamespaceIsolationGroups); - - assertEquals(admin.namespaces().getBookieAffinityGroup(ns2), tenantNamespaceIsolationGroups); - assertEquals(admin.namespaces().getBookieAffinityGroup(ns3), tenantNamespaceIsolationGroups); - assertEquals(admin.namespaces().getBookieAffinityGroup(ns4), tenantNamespaceIsolationGroups); + admin.namespaces().setBookieAffinityGroup(ns2, + new BookieAffinityGroupData(tenantNamespaceIsolationGroups, null)); + admin.namespaces().setBookieAffinityGroup(ns3, + new BookieAffinityGroupData(tenantNamespaceIsolationGroups, null)); + admin.namespaces().setBookieAffinityGroup(ns4, + new BookieAffinityGroupData(tenantNamespaceIsolationGroups, null)); + + assertEquals(admin.namespaces().getBookieAffinityGroup(ns2), + new BookieAffinityGroupData(tenantNamespaceIsolationGroups, null)); + assertEquals(admin.namespaces().getBookieAffinityGroup(ns3), + new BookieAffinityGroupData(tenantNamespaceIsolationGroups, null)); + assertEquals(admin.namespaces().getBookieAffinityGroup(ns4), + new BookieAffinityGroupData(tenantNamespaceIsolationGroups, null)); try { admin.namespaces().getBookieAffinityGroup(ns1); @@ -230,10 +239,139 @@ public void testBookieIsolation() throws Exception { // broker should create only 1 bk-client and factory per isolation-group assertEquals(bkPlacementPolicyToBkClientMap.size(), 1); - Class clazz = bkPlacementPolicyToBkClientMap.keySet().iterator().next() - .getPolicyClass(); - System.out.println(clazz); + } + + /** + * It verifies that "ZkIsolatedBookieEnsemblePlacementPolicy" considers secondary affinity-group if primary group + * doesn't have enough non-faulty bookies. + * + * @throws Exception + */ + @Test + public void testBookieIsilationWithSecondaryGroup() throws Exception { + final String tenant1 = "tenant1"; + final String cluster = "use"; + final String ns1 = String.format("%s/%s/%s", tenant1, cluster, "ns1"); + final String ns2 = String.format("%s/%s/%s", tenant1, cluster, "ns2"); + final String ns3 = String.format("%s/%s/%s", tenant1, cluster, "ns3"); + final String ns4 = String.format("%s/%s/%s", tenant1, cluster, "ns4"); + final int totalPublish = 100; + + final String brokerBookkeeperClientIsolationGroups = "default-group"; + final String tenantNamespaceIsolationGroupsPrimary = "tenant1-isolation-primary"; + final String tenantNamespaceIsolationGroupsSecondary = "tenant1-isolation=secondary"; + + BookieServer[] bookies = bkEnsemble.getBookies(); + ZooKeeper zkClient = bkEnsemble.getZkClient(); + + Set defaultBookies = Sets.newHashSet(bookies[0].getLocalAddress(), + bookies[1].getLocalAddress()); + Set isolatedBookies = Sets.newHashSet(bookies[2].getLocalAddress(), + bookies[3].getLocalAddress()); + Set downedBookies = Sets.newHashSet(new BookieSocketAddress("1.1.1.1:1111"), + new BookieSocketAddress("1.1.1.1:1112")); + + setDefaultIsolationGroup(brokerBookkeeperClientIsolationGroups, zkClient, defaultBookies); + // primary group empty + setDefaultIsolationGroup(tenantNamespaceIsolationGroupsPrimary, zkClient, downedBookies); + setDefaultIsolationGroup(tenantNamespaceIsolationGroupsSecondary, zkClient, isolatedBookies); + + ServiceConfiguration config = new ServiceConfiguration(); + config.setLoadManagerClassName(ModularLoadManagerImpl.class.getName()); + config.setClusterName(cluster); + config.setWebServicePort(Optional.of(PRIMARY_BROKER_WEBSERVICE_PORT)); + config.setZookeeperServers("127.0.0.1" + ":" + ZOOKEEPER_PORT); + config.setBrokerServicePort(Optional.of(PRIMARY_BROKER_PORT)); + config.setAdvertisedAddress("localhost"); + config.setBookkeeperClientIsolationGroups(brokerBookkeeperClientIsolationGroups); + + config.setManagedLedgerDefaultEnsembleSize(2); + config.setManagedLedgerDefaultWriteQuorum(2); + config.setManagedLedgerDefaultAckQuorum(2); + + int totalEntriesPerLedger = 20; + int totalLedgers = totalPublish / totalEntriesPerLedger; + config.setManagedLedgerMaxEntriesPerLedger(totalEntriesPerLedger); + config.setManagedLedgerMinLedgerRolloverTimeMinutes(0); + pulsarService = new PulsarService(config); + pulsarService.start(); + + URL brokerUrl = new URL("http://127.0.0.1" + ":" + PRIMARY_BROKER_WEBSERVICE_PORT); + PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl(brokerUrl.toString()).build(); + + ClusterData clusterData = new ClusterData(pulsarService.getWebServiceAddress()); + admin.clusters().createCluster(cluster, clusterData); + TenantInfo tenantInfo = new TenantInfo(null, Sets.newHashSet(cluster)); + admin.tenants().createTenant(tenant1, tenantInfo); + admin.namespaces().createNamespace(ns1); + admin.namespaces().createNamespace(ns2); + admin.namespaces().createNamespace(ns3); + admin.namespaces().createNamespace(ns4); + admin.namespaces().setBookieAffinityGroup(ns2, new BookieAffinityGroupData( + tenantNamespaceIsolationGroupsPrimary, tenantNamespaceIsolationGroupsSecondary)); + admin.namespaces().setBookieAffinityGroup(ns3, new BookieAffinityGroupData( + tenantNamespaceIsolationGroupsPrimary, tenantNamespaceIsolationGroupsSecondary)); + admin.namespaces().setBookieAffinityGroup(ns4, + new BookieAffinityGroupData(tenantNamespaceIsolationGroupsPrimary, null)); + + assertEquals(admin.namespaces().getBookieAffinityGroup(ns2), new BookieAffinityGroupData( + tenantNamespaceIsolationGroupsPrimary, tenantNamespaceIsolationGroupsSecondary)); + assertEquals(admin.namespaces().getBookieAffinityGroup(ns3), new BookieAffinityGroupData( + tenantNamespaceIsolationGroupsPrimary, tenantNamespaceIsolationGroupsSecondary)); + assertEquals(admin.namespaces().getBookieAffinityGroup(ns4), + new BookieAffinityGroupData(tenantNamespaceIsolationGroupsPrimary, null)); + + try { + admin.namespaces().getBookieAffinityGroup(ns1); + } catch (PulsarAdminException.NotFoundException e) { + // Ok + } + + PulsarClient pulsarClient = PulsarClient.builder().serviceUrl(brokerUrl.toString()) + .statsInterval(-1, TimeUnit.SECONDS).build(); + + PersistentTopic topic1 = (PersistentTopic) createTopicAndPublish(pulsarClient, ns1, "topic1", totalPublish); + PersistentTopic topic2 = (PersistentTopic) createTopicAndPublish(pulsarClient, ns2, "topic1", totalPublish); + PersistentTopic topic3 = (PersistentTopic) createTopicAndPublish(pulsarClient, ns3, "topic1", totalPublish); + + Bookie bookie1 = bookies[0].getBookie(); + Field ledgerManagerField = Bookie.class.getDeclaredField("ledgerManager"); + ledgerManagerField.setAccessible(true); + LedgerManager ledgerManager = (LedgerManager) ledgerManagerField.get(bookie1); + + // namespace: ns1 + ManagedLedgerImpl ml = (ManagedLedgerImpl) topic1.getManagedLedger(); + assertEquals(ml.getLedgersInfoAsList().size(), totalLedgers); + // validate ledgers' ensemble with affinity bookies + assertAffinityBookies(ledgerManager, ml.getLedgersInfoAsList(), defaultBookies); + + // namespace: ns2 + ml = (ManagedLedgerImpl) topic2.getManagedLedger(); + assertEquals(ml.getLedgersInfoAsList().size(), totalLedgers); + // validate ledgers' ensemble with affinity bookies + assertAffinityBookies(ledgerManager, ml.getLedgersInfoAsList(), isolatedBookies); + + // namespace: ns3 + ml = (ManagedLedgerImpl) topic3.getManagedLedger(); + assertEquals(ml.getLedgersInfoAsList().size(), totalLedgers); + // validate ledgers' ensemble with affinity bookies + assertAffinityBookies(ledgerManager, ml.getLedgersInfoAsList(), isolatedBookies); + + ManagedLedgerClientFactory mlFactory = pulsarService.getManagedLedgerClientFactory(); + Map bkPlacementPolicyToBkClientMap = mlFactory + .getBkEnsemblePolicyToBookKeeperMap(); + + // broker should create only 1 bk-client and factory per isolation-group + assertEquals(bkPlacementPolicyToBkClientMap.size(), 1); + + // ns4 doesn't have secondary group so, publish should fail + try { + PersistentTopic topic4 = (PersistentTopic) createTopicAndPublish(pulsarClient, ns4, "topic1", 1); + fail("should have failed due to not enough non-faulty bookie"); + } catch (BrokerPersistenceException e) { + // Ok.. + } } private void assertAffinityBookies(LedgerManager ledgerManager, List ledgers1, @@ -257,7 +395,7 @@ private Topic createTopicAndPublish(PulsarClient pulsarClient, String ns, String .subscribe(); consumer.close(); - ProducerBuilder producerBuilder = pulsarClient.newProducer().topic(topicName); + ProducerBuilder producerBuilder = pulsarClient.newProducer().topic(topicName).sendTimeout(5, TimeUnit.SECONDS); Producer producer = producerBuilder.create(); for (int i = 0; i < totalPublish; i++) { diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Namespaces.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Namespaces.java index 098a8f817cd47..3866c4990526a 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Namespaces.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/Namespaces.java @@ -29,6 +29,7 @@ import org.apache.pulsar.client.admin.PulsarAdminException.PreconditionFailedException; import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.BacklogQuota; +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; import org.apache.pulsar.common.policies.data.BundlesData; import org.apache.pulsar.common.policies.data.DispatchRate; import org.apache.pulsar.common.policies.data.PersistencePolicies; @@ -761,9 +762,28 @@ List getAntiAffinityNamespaces(String tenant, String cluster, String nam */ PersistencePolicies getPersistence(String namespace) throws PulsarAdminException; - void setBookieAffinityGroup(String namespace, String bookieAffinityGroup) throws PulsarAdminException; + /** + * Set bookie affinity group for a namespace to isolate namespace write to bookies that are part of given affinity + * group. + * + * @param namespace + * namespace name + * @param bookieAffinityGroup + * bookie affinity group + * @throws PulsarAdminException + */ + void setBookieAffinityGroup(String namespace, BookieAffinityGroupData bookieAffinityGroup) + throws PulsarAdminException; + - String getBookieAffinityGroup(String namespace) throws PulsarAdminException; + /** + * Get bookie affinity group configured for a namespace. + * + * @param namespace + * @return + * @throws PulsarAdminException + */ + BookieAffinityGroupData getBookieAffinityGroup(String namespace) throws PulsarAdminException; /** * Set the retention configuration for all the topics on a namespace. diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/NamespacesImpl.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/NamespacesImpl.java index a577e53a8eddc..19368d9c13b35 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/NamespacesImpl.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/NamespacesImpl.java @@ -36,6 +36,7 @@ import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.BacklogQuota; +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; import org.apache.pulsar.common.policies.data.BacklogQuota.BacklogQuotaType; import org.apache.pulsar.common.policies.data.BundlesData; import org.apache.pulsar.common.policies.data.DispatchRate; @@ -404,22 +405,22 @@ public void setPersistence(String namespace, PersistencePolicies persistence) th } @Override - public void setBookieAffinityGroup(String namespace, String bookieAffinityGroup) throws PulsarAdminException { + public void setBookieAffinityGroup(String namespace, BookieAffinityGroupData bookieAffinityGroup) throws PulsarAdminException { try { NamespaceName ns = NamespaceName.get(namespace); - WebTarget path = namespacePath(ns, "persistence", "bookieAffinity", bookieAffinityGroup); - request(path).post(Entity.entity("", MediaType.APPLICATION_JSON), ErrorData.class); + WebTarget path = namespacePath(ns, "persistence", "bookieAffinity"); + request(path).post(Entity.entity(bookieAffinityGroup, MediaType.APPLICATION_JSON), ErrorData.class); } catch (Exception e) { throw getApiException(e); } } @Override - public String getBookieAffinityGroup(String namespace) throws PulsarAdminException { + public BookieAffinityGroupData getBookieAffinityGroup(String namespace) throws PulsarAdminException { try { NamespaceName ns = NamespaceName.get(namespace); WebTarget path = namespacePath(ns, "persistence", "bookieAffinity"); - return request(path).get(String.class); + return request(path).get(BookieAffinityGroupData.class); } catch (Exception e) { throw getApiException(e); } diff --git a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java index 47219290fd970..96d47419b96b0 100644 --- a/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java +++ b/pulsar-client-tools-test/src/test/java/org/apache/pulsar/admin/cli/PulsarAdminToolTest.java @@ -55,6 +55,7 @@ import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.BacklogQuota.RetentionPolicy; +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; import org.apache.pulsar.common.policies.data.BookieInfo; import org.apache.pulsar.common.policies.data.BundlesData; import org.apache.pulsar.common.policies.data.ClusterData; @@ -280,8 +281,10 @@ void namespaces() throws Exception { namespaces.run(split("get-clusters myprop/clust/ns1")); verify(mockNamespaces).getNamespaceReplicationClusters("myprop/clust/ns1"); - namespaces.run(split("set-bookie-affinity-group myprop/clust/ns1 --group test")); - verify(mockNamespaces).setBookieAffinityGroup("myprop/clust/ns1", "test"); + namespaces + .run(split("set-bookie-affinity-group myprop/clust/ns1 --primary-group test1 --secondary-group test2")); + verify(mockNamespaces).setBookieAffinityGroup("myprop/clust/ns1", + new BookieAffinityGroupData("test1", "test2")); namespaces.run(split("get-bookie-affinity-group myprop/clust/ns1")); verify(mockNamespaces).getBookieAffinityGroup("myprop/clust/ns1"); diff --git a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaces.java b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaces.java index 359f25bb8660a..a837f6c1eeb2a 100644 --- a/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaces.java +++ b/pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdNamespaces.java @@ -37,6 +37,7 @@ import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.policies.data.BacklogQuota; +import org.apache.pulsar.common.policies.data.BookieAffinityGroupData; import org.apache.pulsar.common.policies.data.BundlesData; import org.apache.pulsar.common.policies.data.DispatchRate; import org.apache.pulsar.common.policies.data.PersistencePolicies; @@ -426,15 +427,19 @@ private class SetBookieAffinityGroup extends CliCommand { @Parameter(description = "tenant/namespace", required = true) private java.util.List params; - @Parameter(names = { "--group", - "-g" }, description = "Bookie-affinity group name where namespace messages should be written", required = true) - private String bookieAffinityGroupName; + @Parameter(names = { "--primary-group", + "-pg" }, description = "Bookie-affinity primary-groups (comma separated) name where namespace messages should be written", required = true) + private String bookieAffinityGroupNamePrimary; + @Parameter(names = { "--secondary-group", + "-sg" }, description = "Bookie-affinity secondary-group (comma separated) name where namespace messages should be written", required = false) + private String bookieAffinityGroupNameSecondary; @Override void run() throws PulsarAdminException { String namespace = validateNamespace(params); - admin.namespaces().setBookieAffinityGroup(namespace, bookieAffinityGroupName); + admin.namespaces().setBookieAffinityGroup(namespace, + new BookieAffinityGroupData(bookieAffinityGroupNamePrimary, bookieAffinityGroupNameSecondary)); } } diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/BookieAffinityGroupData.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/BookieAffinityGroupData.java new file mode 100644 index 0000000000000..098d4d33ccd67 --- /dev/null +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/BookieAffinityGroupData.java @@ -0,0 +1,56 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.common.policies.data; + +import com.google.common.base.Objects; + +public class BookieAffinityGroupData { + + public String bookkeeperAffinityGroupPrimary; + public String bookkeeperAffinityGroupSecondary; + + public BookieAffinityGroupData() { + } + + public BookieAffinityGroupData(String bookkeeperAffinityGroupPrimary, String bookkeeperAffinityGroupSecondary) { + this.bookkeeperAffinityGroupPrimary = bookkeeperAffinityGroupPrimary; + this.bookkeeperAffinityGroupSecondary = bookkeeperAffinityGroupSecondary; + } + + @Override + public int hashCode() { + return Objects.hashCode(bookkeeperAffinityGroupPrimary, bookkeeperAffinityGroupSecondary); + } + + @Override + public boolean equals(Object obj) { + if (obj instanceof BookieAffinityGroupData) { + BookieAffinityGroupData other = (BookieAffinityGroupData) obj; + return Objects.equal(bookkeeperAffinityGroupPrimary, other.bookkeeperAffinityGroupPrimary) + && Objects.equal(bookkeeperAffinityGroupSecondary, other.bookkeeperAffinityGroupSecondary); + } + return false; + } + + @Override + public String toString() { + return String.format("bookkeeperAffinityGroupPrimary=%s bookkeeperAffinityGroupSecondary=%s", + bookkeeperAffinityGroupPrimary, bookkeeperAffinityGroupSecondary); + } +} diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/LocalPolicies.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/LocalPolicies.java index dc8b6b18f7f7c..c8ae18a372777 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/LocalPolicies.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/LocalPolicies.java @@ -25,7 +25,7 @@ public class LocalPolicies { public BundlesData bundles; - public String bookkeeperAffinityGroup; + public BookieAffinityGroupData bookieAffinityGroup; public LocalPolicies() { bundles = defaultBundle(); @@ -33,7 +33,7 @@ public LocalPolicies() { @Override public int hashCode() { - return Objects.hashCode(bundles, bookkeeperAffinityGroup); + return Objects.hashCode(bundles, bookieAffinityGroup); } @Override @@ -41,7 +41,7 @@ public boolean equals(Object obj) { if (obj instanceof LocalPolicies) { LocalPolicies other = (LocalPolicies) obj; return Objects.equal(bundles, other.bundles) - && Objects.equal(bookkeeperAffinityGroup, other.bookkeeperAffinityGroup); + && Objects.equal(bookieAffinityGroup, other.bookieAffinityGroup); } return false; } diff --git a/pulsar-zookeeper-utils/src/main/java/org/apache/pulsar/zookeeper/ZkIsolatedBookieEnsemblePlacementPolicy.java b/pulsar-zookeeper-utils/src/main/java/org/apache/pulsar/zookeeper/ZkIsolatedBookieEnsemblePlacementPolicy.java index a46dba0d82c3d..89b46427a369e 100644 --- a/pulsar-zookeeper-utils/src/main/java/org/apache/pulsar/zookeeper/ZkIsolatedBookieEnsemblePlacementPolicy.java +++ b/pulsar-zookeeper-utils/src/main/java/org/apache/pulsar/zookeeper/ZkIsolatedBookieEnsemblePlacementPolicy.java @@ -54,16 +54,12 @@ public class ZkIsolatedBookieEnsemblePlacementPolicy extends RackawareEnsemblePl private static final Logger LOG = LoggerFactory.getLogger(ZkIsolatedBookieEnsemblePlacementPolicy.class); public static final String ISOLATION_BOOKIE_GROUPS = "isolationBookieGroups"; - // if policy doesn't find min-available bookies in primary-isolationBookieGroups then it uses bookies from - // secondaryIsolationBookieGroups - public static final String MIN_AVAILABLE_PRIMARY_ISOLATED_BOOKIE = "minAvailablePrimaryIsolatedBookies"; public static final String SECONDARY_ISOLATION_BOOKIE_GROUPS = "secondaryIsolationBookieGroups"; private ZooKeeperCache bookieMappingCache = null; private final List primaryIsolationGroups = new ArrayList(); private final List secondaryIsolationGroups = new ArrayList(); - private int minAvailablePrimaryIsolatedBookies = 0; private final ObjectMapper jsonMapper = ObjectMapperFactory.create(); public ZkIsolatedBookieEnsemblePlacementPolicy() { @@ -91,9 +87,6 @@ public RackawareEnsemblePlacementPolicyImpl initialize(ClientConfiguration conf, } } } - minAvailablePrimaryIsolatedBookies = conf.getProperty(MIN_AVAILABLE_PRIMARY_ISOLATED_BOOKIE) != null - ? (int) conf.getProperty(MIN_AVAILABLE_PRIMARY_ISOLATED_BOOKIE) - : 0; return super.initialize(conf, optionalDnsResolver, timer, featureProvider, statsLogger); } @@ -127,7 +120,7 @@ private ZooKeeperCache getAndSetZkCache(Configuration conf) { public PlacementResult> newEnsemble(int ensembleSize, int writeQuorumSize, int ackQuorumSize, Map customMetadata, Set excludeBookies) throws BKNotEnoughBookiesException { - Set blacklistedBookies = getBlacklistedBookies(); + Set blacklistedBookies = getBlacklistedBookies(ensembleSize); if (excludeBookies == null) { excludeBookies = new HashSet(); } @@ -140,7 +133,7 @@ public PlacementResult replaceBookie(int ensembleSize, int Map customMetadata, List currentEnsemble, BookieSocketAddress bookieToReplace, Set excludeBookies) throws BKNotEnoughBookiesException { - Set blacklistedBookies = getBlacklistedBookies(); + Set blacklistedBookies = getBlacklistedBookies(ensembleSize); if (excludeBookies == null) { excludeBookies = new HashSet(); } @@ -149,7 +142,7 @@ public PlacementResult replaceBookie(int ensembleSize, int bookieToReplace, excludeBookies); } - private Set getBlacklistedBookies() { + private Set getBlacklistedBookies(int ensembleSize) { Set blacklistedBookies = new HashSet(); try { if (bookieMappingCache != null) { @@ -158,11 +151,18 @@ private Set getBlacklistedBookies() { .orElseThrow(() -> new KeeperException.NoNodeException( ZkBookieRackAffinityMapping.BOOKIE_INFO_ROOT_PATH)); Set allBookies = allGroupsBookieMapping.keySet(); + int totalAvailableBookiesInPrimaryGroup = 0; for (String group : allBookies) { + Set bookiesInGroup = allGroupsBookieMapping.get(group).keySet(); if (!primaryIsolationGroups.contains(group)) { - for (String bookieAddress : allGroupsBookieMapping.get(group).keySet()) { + for (String bookieAddress : bookiesInGroup) { blacklistedBookies.add(new BookieSocketAddress(bookieAddress)); } + } else { + for (String groupBookie : bookiesInGroup) { + totalAvailableBookiesInPrimaryGroup += knownBookies + .containsKey(new BookieSocketAddress(groupBookie)) ? 1 : 0; + } } } // sometime while doing isolation, user might not want to remove isolated bookies from other default @@ -178,7 +178,7 @@ private Set getBlacklistedBookies() { } } // if primary-isolated-bookies are not enough then add consider secondary isolated bookie group as well. - if ((allBookies.size() - blacklistedBookies.size()) < minAvailablePrimaryIsolatedBookies) { + if (totalAvailableBookiesInPrimaryGroup < ensembleSize) { for (String group : secondaryIsolationGroups) { Map bookieGroup = allGroupsBookieMapping.get(group); if (bookieGroup != null && !bookieGroup.isEmpty()) { diff --git a/pulsar-zookeeper-utils/src/test/java/org/apache/pulsar/zookeeper/ZkIsolatedBookieEnsemblePlacementPolicyTest.java b/pulsar-zookeeper-utils/src/test/java/org/apache/pulsar/zookeeper/ZkIsolatedBookieEnsemblePlacementPolicyTest.java index 8c2576b1a0462..6187f1f50ebe4 100644 --- a/pulsar-zookeeper-utils/src/test/java/org/apache/pulsar/zookeeper/ZkIsolatedBookieEnsemblePlacementPolicyTest.java +++ b/pulsar-zookeeper-utils/src/test/java/org/apache/pulsar/zookeeper/ZkIsolatedBookieEnsemblePlacementPolicyTest.java @@ -390,7 +390,6 @@ public void testSecondaryIsolationGroupsBookies() throws Exception { }); bkClientConf.setProperty(ZkIsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, isolatedGroup); bkClientConf.setProperty(ZkIsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, secondaryIsolatedGroup); - bkClientConf.setProperty(ZkIsolatedBookieEnsemblePlacementPolicy.MIN_AVAILABLE_PRIMARY_ISOLATED_BOOKIE, 2); isolationPolicy.initialize(bkClientConf, Optional.empty(), timer, SettableFeatureProvider.DISABLE_ALL, NullStatsLogger.INSTANCE); isolationPolicy.onClusterChanged(writableBookies, readOnlyBookies); @@ -436,7 +435,6 @@ public void testSecondaryIsolationGroupsBookiesNegative() throws Exception { bkClientConf.setProperty(ZkIsolatedBookieEnsemblePlacementPolicy.ISOLATION_BOOKIE_GROUPS, isolatedGroup); bkClientConf.setProperty(ZkIsolatedBookieEnsemblePlacementPolicy.SECONDARY_ISOLATION_BOOKIE_GROUPS, secondaryIsolatedGroup); - bkClientConf.setProperty(ZkIsolatedBookieEnsemblePlacementPolicy.MIN_AVAILABLE_PRIMARY_ISOLATED_BOOKIE, 2); isolationPolicy.initialize(bkClientConf, Optional.empty(), timer, SettableFeatureProvider.DISABLE_ALL, NullStatsLogger.INSTANCE); isolationPolicy.onClusterChanged(writableBookies, readOnlyBookies);