From 540a01d85000f6505c9542f5db55131521a1280a Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Sun, 20 Feb 2022 02:11:12 +0800 Subject: [PATCH 1/7] Make stadnalone use setMetadataStoreUrl and fix functions incorrect zkString --- conf/standalone.conf | 18 ++++++++++--- .../pulsar/PulsarStandaloneBuilder.java | 7 +++-- .../pulsar/PulsarStandaloneStarter.java | 7 +++-- .../conf/InternalConfigurationData.java | 26 +++++++++---------- .../functions/worker/PulsarWorkerService.java | 14 ++++++---- .../pulsar/functions/worker/WorkerUtils.java | 4 ++- 6 files changed, 49 insertions(+), 27 deletions(-) diff --git a/conf/standalone.conf b/conf/standalone.conf index d774b7b0bf65b..800a01f490345 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -19,11 +19,15 @@ ### --- General broker settings --- ### -# Zookeeper quorum connection string -zookeeperServers= +# The metadata store URL +# Examples: +# * zk:my-zk-1:2181,my-zk-2:2181,my-zk-3:2181 +# * my-zk-1:2181,my-zk-2:2181,my-zk-3:2181 (will default to ZooKeeper when the schema is not specified) +# * zk:my-zk-1:2181,my-zk-2:2181,my-zk-3:2181/my-chroot-path (to add a ZK chroot path) +metadataStoreUrl= -# Configuration Store connection string -configurationStoreServers= +# The metadata store URL for the configuration data. If empty, we fall back to use metadataStoreUrl +configurationMetadataStoreUrl= brokerServicePort=6650 @@ -1076,3 +1080,9 @@ zooKeeperOperationTimeoutSeconds=-1 # ZooKeeper cache expiry time in seconds # Deprecated: use metadataStoreCacheExpirySeconds zooKeeperCacheExpirySeconds=-1 + +# Zookeeper quorum connection string +zookeeperServers= + +# Configuration Store connection string +configurationStoreServers= diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneBuilder.java b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneBuilder.java index 581db63bb0edd..70d457dbb4652 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneBuilder.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneBuilder.java @@ -21,6 +21,7 @@ import static org.apache.commons.lang3.StringUtils.isBlank; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.ServiceConfigurationUtils; +import org.apache.pulsar.metadata.impl.ZKMetadataStore; public final class PulsarStandaloneBuilder { @@ -114,8 +115,10 @@ public PulsarStandalone build() { } // Set ZK server's host to localhost - pulsarStandalone.getConfig().setZookeeperServers(zkServers + ":" + pulsarStandalone.getZkPort()); - pulsarStandalone.getConfig().setConfigurationStoreServers(zkServers + ":" + pulsarStandalone.getZkPort()); + final String metadataStoreUrl = + ZKMetadataStore.ZK_SCHEME_IDENTIFIER + zkServers + ":" + pulsarStandalone.getZkPort(); + pulsarStandalone.getConfig().setMetadataStoreUrl(metadataStoreUrl); + pulsarStandalone.getConfig().setConfigurationMetadataStoreUrl(metadataStoreUrl); pulsarStandalone.getConfig().setRunningStandalone(true); return pulsarStandalone; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java index 92b3e3e64eebe..cb72c679bcb7a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java @@ -28,6 +28,7 @@ import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.common.configuration.PulsarConfigurationLoader; import org.apache.pulsar.common.util.CmdGenerateDocs; +import org.apache.pulsar.metadata.impl.ZKMetadataStore; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -100,8 +101,10 @@ public PulsarStandaloneStarter(String[] args) throws Exception { } } } - config.setZookeeperServers(zkServers + ":" + this.getZkPort()); - config.setConfigurationStoreServers(zkServers + ":" + this.getZkPort()); + final String metadataStoreUrl = + ZKMetadataStore.ZK_SCHEME_IDENTIFIER + zkServers + ":" + this.getZkPort(); + config.setMetadataStoreUrl(metadataStoreUrl); + config.setConfigurationMetadataStoreUrl(metadataStoreUrl); config.setRunningStandalone(true); diff --git a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java index 80da3caa4bab7..2927aad8afe41 100644 --- a/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java +++ b/pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/conf/InternalConfigurationData.java @@ -27,8 +27,8 @@ @ToString public class InternalConfigurationData { - private String zookeeperServers; - private String configurationStoreServers; + private String metadataStoreUrl; + private String configurationMetadataStoreUrl; @Deprecated private String ledgersRootPath; private String bookkeeperMetadataServiceUri; @@ -38,23 +38,23 @@ public InternalConfigurationData() { } public InternalConfigurationData(String zookeeperServers, - String configurationStoreServers, + String configurationMetadataStoreUrl, String ledgersRootPath, String bookkeeperMetadataServiceUri, String stateStorageServiceUrl) { - this.zookeeperServers = zookeeperServers; - this.configurationStoreServers = configurationStoreServers; + this.metadataStoreUrl = zookeeperServers; + this.configurationMetadataStoreUrl = configurationMetadataStoreUrl; this.ledgersRootPath = ledgersRootPath; this.bookkeeperMetadataServiceUri = bookkeeperMetadataServiceUri; this.stateStorageServiceUrl = stateStorageServiceUrl; } - public String getZookeeperServers() { - return zookeeperServers; + public String getMetadataStoreUrl() { + return metadataStoreUrl; } - public String getConfigurationStoreServers() { - return configurationStoreServers; + public String getConfigurationMetadataStoreUrl() { + return configurationMetadataStoreUrl; } /** @deprecated */ @@ -77,8 +77,8 @@ public boolean equals(Object obj) { return false; } InternalConfigurationData other = (InternalConfigurationData) obj; - return Objects.equals(zookeeperServers, other.zookeeperServers) - && Objects.equals(configurationStoreServers, other.configurationStoreServers) + return Objects.equals(metadataStoreUrl, other.metadataStoreUrl) + && Objects.equals(configurationMetadataStoreUrl, other.configurationMetadataStoreUrl) && Objects.equals(ledgersRootPath, other.ledgersRootPath) && Objects.equals(bookkeeperMetadataServiceUri, other.bookkeeperMetadataServiceUri) && Objects.equals(stateStorageServiceUrl, other.stateStorageServiceUrl); @@ -86,8 +86,8 @@ public boolean equals(Object obj) { @Override public int hashCode() { - return Objects.hash(zookeeperServers, - configurationStoreServers, + return Objects.hash(metadataStoreUrl, + configurationMetadataStoreUrl, ledgersRootPath, bookkeeperMetadataServiceUri, stateStorageServiceUrl); diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java index 436e27e00a82c..9b97d84bb0196 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java @@ -69,6 +69,7 @@ import org.apache.pulsar.functions.worker.service.api.Sources; import org.apache.pulsar.functions.worker.service.api.Workers; import org.apache.pulsar.metadata.api.MetadataStoreException.AlreadyExistsException; +import org.apache.pulsar.metadata.impl.ZKMetadataStore; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -272,13 +273,14 @@ private static URI initializeStandaloneWorkerService(PulsarClientCreator clientC URI dlogURI; try { if (workerConfig.isInitializedDlogMetadata()) { - dlogURI = WorkerUtils.newDlogNamespaceURI(internalConf.getZookeeperServers()); + dlogURI = WorkerUtils.newDlogNamespaceURI(internalConf.getMetadataStoreUrl() + .substring(ZKMetadataStore.ZK_SCHEME_IDENTIFIER.length())); } else { dlogURI = WorkerUtils.initializeDlogNamespace(internalConf); } } catch (IOException ioe) { log.error("Failed to initialize dlog namespace with zookeeper {} at metadata service uri {} for storing " - + "function packages", internalConf.getZookeeperServers(), + + "function packages", internalConf.getMetadataStoreUrl(), internalConf.getBookkeeperMetadataServiceUri(), ioe); throw ioe; } @@ -354,15 +356,17 @@ public void initInBroker(ServiceConfiguration brokerConfig, URI dlogURI; try { // initializing dlog namespace for function worker - if (workerConfig.isInitializedDlogMetadata()) { - dlogURI = WorkerUtils.newDlogNamespaceURI(internalConf.getZookeeperServers()); + if (workerConfig.isInitializedDlogMetadata()){ + dlogURI = + WorkerUtils.newDlogNamespaceURI(internalConf.getMetadataStoreUrl() + .substring(ZKMetadataStore.ZK_SCHEME_IDENTIFIER.length())); } else { dlogURI = WorkerUtils.initializeDlogNamespace(internalConf); } } catch (IOException ioe) { LOG.error("Failed to initialize dlog namespace with zookeeper {} at at metadata service uri {} for " + "storing function packages", - internalConf.getZookeeperServers(), internalConf.getBookkeeperMetadataServiceUri(), ioe); + internalConf.getMetadataStoreUrl(), internalConf.getBookkeeperMetadataServiceUri(), ioe); throw ioe; } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java index bf86867208560..1ff6f6108f61a 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java @@ -62,6 +62,7 @@ import org.apache.pulsar.functions.utils.FunctionCommon; import org.apache.pulsar.functions.worker.dlog.DLInputStream; import org.apache.pulsar.functions.worker.dlog.DLOutputStream; +import org.apache.pulsar.metadata.impl.ZKMetadataStore; import org.apache.zookeeper.KeeperException.Code; @Slf4j @@ -173,7 +174,8 @@ public static URI initializeDlogNamespace(InternalConfigurationData internalConf // for BC purposes if (internalConf.getBookkeeperMetadataServiceUri() == null) { ledgersRootPath = internalConf.getLedgersRootPath(); - ledgersStoreServers = internalConf.getZookeeperServers(); + ledgersStoreServers = internalConf.getMetadataStoreUrl() + .substring(ZKMetadataStore.ZK_SCHEME_IDENTIFIER.length()); chrootPath = ""; } else { URI metadataServiceUri = URI.create(internalConf.getBookkeeperMetadataServiceUri()); From e03fb7e8c6f0d9eb2b91ffaa47664115d338f0b0 Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Fri, 25 Feb 2022 14:10:39 +0800 Subject: [PATCH 2/7] fix unit test error --- .../functions/worker/PulsarWorkerService.java | 17 +++++++++++------ .../pulsar/functions/worker/WorkerUtils.java | 9 ++++++--- .../bookkeeper/BookKeeperPackagesStorage.java | 9 ++++++--- .../BookKeeperPackagesStorageConfiguration.java | 4 ++-- .../BookKeeperPackagesStorageTest.java | 11 +++++++---- 5 files changed, 32 insertions(+), 18 deletions(-) diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java index 9b97d84bb0196..bb465193ba104 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java @@ -19,6 +19,7 @@ package org.apache.pulsar.functions.worker; import static org.apache.pulsar.common.policies.data.PoliciesUtil.getBundles; +import static org.apache.pulsar.metadata.impl.ZKMetadataStore.ZK_SCHEME_IDENTIFIER; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.annotations.VisibleForTesting; @@ -69,7 +70,6 @@ import org.apache.pulsar.functions.worker.service.api.Sources; import org.apache.pulsar.functions.worker.service.api.Workers; import org.apache.pulsar.metadata.api.MetadataStoreException.AlreadyExistsException; -import org.apache.pulsar.metadata.impl.ZKMetadataStore; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -273,8 +273,11 @@ private static URI initializeStandaloneWorkerService(PulsarClientCreator clientC URI dlogURI; try { if (workerConfig.isInitializedDlogMetadata()) { - dlogURI = WorkerUtils.newDlogNamespaceURI(internalConf.getMetadataStoreUrl() - .substring(ZKMetadataStore.ZK_SCHEME_IDENTIFIER.length())); + String metadataStoreUrl = internalConf.getMetadataStoreUrl(); + if (metadataStoreUrl.startsWith(ZK_SCHEME_IDENTIFIER)) { + metadataStoreUrl = metadataStoreUrl.substring(ZK_SCHEME_IDENTIFIER.length()); + } + dlogURI = WorkerUtils.newDlogNamespaceURI(metadataStoreUrl); } else { dlogURI = WorkerUtils.initializeDlogNamespace(internalConf); } @@ -357,9 +360,11 @@ public void initInBroker(ServiceConfiguration brokerConfig, try { // initializing dlog namespace for function worker if (workerConfig.isInitializedDlogMetadata()){ - dlogURI = - WorkerUtils.newDlogNamespaceURI(internalConf.getMetadataStoreUrl() - .substring(ZKMetadataStore.ZK_SCHEME_IDENTIFIER.length())); + String metadataStoreUrl = internalConf.getMetadataStoreUrl(); + if (metadataStoreUrl.startsWith(ZK_SCHEME_IDENTIFIER)) { + metadataStoreUrl = metadataStoreUrl.substring(ZK_SCHEME_IDENTIFIER.length()); + } + dlogURI = WorkerUtils.newDlogNamespaceURI(metadataStoreUrl); } else { dlogURI = WorkerUtils.initializeDlogNamespace(internalConf); } diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java index 1ff6f6108f61a..2f2c667ba47cb 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java @@ -20,6 +20,7 @@ import static java.nio.file.StandardCopyOption.REPLACE_EXISTING; import static org.apache.commons.lang3.StringUtils.isNotBlank; +import static org.apache.pulsar.metadata.impl.ZKMetadataStore.ZK_SCHEME_IDENTIFIER; import java.io.File; import java.io.FileInputStream; import java.io.FileOutputStream; @@ -62,7 +63,6 @@ import org.apache.pulsar.functions.utils.FunctionCommon; import org.apache.pulsar.functions.worker.dlog.DLInputStream; import org.apache.pulsar.functions.worker.dlog.DLOutputStream; -import org.apache.pulsar.metadata.impl.ZKMetadataStore; import org.apache.zookeeper.KeeperException.Code; @Slf4j @@ -174,8 +174,11 @@ public static URI initializeDlogNamespace(InternalConfigurationData internalConf // for BC purposes if (internalConf.getBookkeeperMetadataServiceUri() == null) { ledgersRootPath = internalConf.getLedgersRootPath(); - ledgersStoreServers = internalConf.getMetadataStoreUrl() - .substring(ZKMetadataStore.ZK_SCHEME_IDENTIFIER.length()); + String metadataStoreUrl = internalConf.getMetadataStoreUrl(); + if (metadataStoreUrl.startsWith(ZK_SCHEME_IDENTIFIER)) { + metadataStoreUrl = metadataStoreUrl.substring(ZK_SCHEME_IDENTIFIER.length()); + } + ledgersStoreServers = metadataStoreUrl; chrootPath = ""; } else { URI metadataServiceUri = URI.create(internalConf.getBookkeeperMetadataServiceUri()); diff --git a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java index 3da138475ce8c..1a78ec5ef71e6 100644 --- a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java +++ b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java @@ -49,6 +49,7 @@ public class BookKeeperPackagesStorage implements PackagesStorage { private static final String NS_CLIENT_ID = "packages-management"; + public static final String ZK_SCHEME_IDENTIFIER = "zk:"; final BookKeeperPackagesStorageConfiguration configuration; private Namespace namespace; @@ -92,12 +93,14 @@ private URI initializeDlogNamespace() throws IOException { ledgersRootPath = metadataServiceUri.getPath(); } else { ledgersRootPath = configuration.getPackagesManagementLedgerRootPath(); - ledgersStoreServers = configuration.getZookeeperServers(); + ledgersStoreServers = configuration.getMetadataStoreUrl(); + if (ledgersStoreServers.startsWith(ZK_SCHEME_IDENTIFIER)) { + ledgersStoreServers = ledgersStoreServers.substring(ZK_SCHEME_IDENTIFIER.length()); + } } BKDLConfig bkdlConfig = new BKDLConfig(ledgersStoreServers, ledgersRootPath); DLMetadata dlMetadata = DLMetadata.create(bkdlConfig); - URI dlogURI = URI.create(String.format("distributedlog://%s/pulsar/packages", - configuration.getZookeeperServers())); + URI dlogURI = URI.create(String.format("distributedlog://%s/pulsar/packages", ledgersStoreServers)); try { dlMetadata.create(dlogURI); } catch (ZKException e) { diff --git a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java index 226b80abeaa30..00ce9b53fef95 100644 --- a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java +++ b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java @@ -38,8 +38,8 @@ int getPackagesReplicas() { return Integer.parseInt(getProperty("packagesReplicas")); } - String getZookeeperServers() { - return getProperty("zookeeperServers"); + String getMetadataStoreUrl() { + return getProperty("metadataStoreUrl"); } String getPackagesManagementLedgerRootPath() { diff --git a/pulsar-package-management/bookkeeper-storage/src/test/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageTest.java b/pulsar-package-management/bookkeeper-storage/src/test/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageTest.java index 69312076410ca..8715b35bc2809 100644 --- a/pulsar-package-management/bookkeeper-storage/src/test/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageTest.java +++ b/pulsar-package-management/bookkeeper-storage/src/test/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageTest.java @@ -18,6 +18,7 @@ */ package org.apache.pulsar.packages.management.storage.bookkeeper; +import static org.apache.pulsar.metadata.impl.ZKMetadataStore.ZK_SCHEME_IDENTIFIER; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; @@ -50,7 +51,7 @@ public void setup() throws Exception { PackagesStorageProvider provider = PackagesStorageProvider .newProvider(BookKeeperPackagesStorageProvider.class.getName()); DefaultPackagesStorageConfiguration configuration = new DefaultPackagesStorageConfiguration(); - configuration.setProperty("zookeeperServers", zkUtil.getZooKeeperConnectString()); + configuration.setProperty("metadataStoreUrl", ZK_SCHEME_IDENTIFIER + zkUtil.getZooKeeperConnectString()); configuration.setProperty("packagesReplicas", "1"); configuration.setProperty("packagesManagementLedgerRootPath", "/ledgers"); storage = provider.getStorage(configuration); @@ -68,7 +69,8 @@ public void teardown() throws Exception { public void testConfiguration() { assertTrue(storage instanceof BookKeeperPackagesStorage); BookKeeperPackagesStorage bkStorage = (BookKeeperPackagesStorage) storage; - assertEquals(bkStorage.configuration.getZookeeperServers(), zkUtil.getZooKeeperConnectString()); + assertEquals(bkStorage.configuration.getMetadataStoreUrl() + .substring(ZK_SCHEME_IDENTIFIER.length()), zkUtil.getZooKeeperConnectString()); assertEquals(bkStorage.configuration.getPackagesReplicas(), 1); assertEquals(bkStorage.configuration.getPackagesManagementLedgerRootPath(), "/ledgers"); } @@ -198,7 +200,8 @@ public void testReadWriteOperationsWithSeparatedBkCluster() throws Exception { .newProvider(BookKeeperPackagesStorageProvider.class.getName()); DefaultPackagesStorageConfiguration configuration = new DefaultPackagesStorageConfiguration(); // set the unavailable bk cluster with mock zookeeper path - configuration.setProperty("zookeeperServers", zkUtil.getZooKeeperConnectString() + "/mock"); + configuration.setProperty("metadataStoreUrl", ZK_SCHEME_IDENTIFIER + zkUtil.getZooKeeperConnectString() + + "/mock"); configuration.setProperty("packagesReplicas", "1"); configuration.setProperty("packagesManagementLedgerRootPath", "/ledgers"); PackagesStorage storage1 = provider.getStorage(configuration); @@ -221,7 +224,7 @@ public void testReadWriteOperationsWithSeparatedBkCluster() throws Exception { // set the available bk cluster with bookkeeperMetadataServiceUri using actual zookeeper path String bookkeeperMetadataServiceUri = String.format("zk+null://%s/ledgers", zkUtil.getZooKeeperConnectString()); DefaultPackagesStorageConfiguration configuration2 = new DefaultPackagesStorageConfiguration(); - configuration2.setProperty("zookeeperServers", zkUtil.getZooKeeperConnectString()); + configuration2.setProperty("metadataStoreUrl", ZK_SCHEME_IDENTIFIER + zkUtil.getZooKeeperConnectString()); configuration2.setProperty("bookkeeperMetadataServiceUri", bookkeeperMetadataServiceUri); configuration2.setProperty("packagesReplicas", "1"); PackagesStorage storage2 = provider.getStorage(configuration2); From 7074f4c077c06223df401c5cbed7191bd03b185a Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Fri, 25 Feb 2022 15:28:26 +0800 Subject: [PATCH 3/7] fix code style --- .../org/apache/pulsar/functions/worker/PulsarWorkerService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java index bb465193ba104..caa2a037609a5 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java @@ -359,7 +359,7 @@ public void initInBroker(ServiceConfiguration brokerConfig, URI dlogURI; try { // initializing dlog namespace for function worker - if (workerConfig.isInitializedDlogMetadata()){ + if (workerConfig.isInitializedDlogMetadata()) { String metadataStoreUrl = internalConf.getMetadataStoreUrl(); if (metadataStoreUrl.startsWith(ZK_SCHEME_IDENTIFIER)) { metadataStoreUrl = metadataStoreUrl.substring(ZK_SCHEME_IDENTIFIER.length()); From 67a5f6f08068a39fe5e6a9290023f62fa74000c1 Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Fri, 25 Feb 2022 19:01:16 +0800 Subject: [PATCH 4/7] use removeIdentifierFromMetadataURL --- .../pulsar/functions/worker/PulsarWorkerService.java | 12 +++--------- .../apache/pulsar/functions/worker/WorkerUtils.java | 8 ++------ 2 files changed, 5 insertions(+), 15 deletions(-) diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java index caa2a037609a5..730e5af4751c7 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/PulsarWorkerService.java @@ -19,7 +19,7 @@ package org.apache.pulsar.functions.worker; import static org.apache.pulsar.common.policies.data.PoliciesUtil.getBundles; -import static org.apache.pulsar.metadata.impl.ZKMetadataStore.ZK_SCHEME_IDENTIFIER; +import static org.apache.pulsar.metadata.impl.MetadataStoreFactoryImpl.removeIdentifierFromMetadataURL; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.annotations.VisibleForTesting; @@ -273,10 +273,7 @@ private static URI initializeStandaloneWorkerService(PulsarClientCreator clientC URI dlogURI; try { if (workerConfig.isInitializedDlogMetadata()) { - String metadataStoreUrl = internalConf.getMetadataStoreUrl(); - if (metadataStoreUrl.startsWith(ZK_SCHEME_IDENTIFIER)) { - metadataStoreUrl = metadataStoreUrl.substring(ZK_SCHEME_IDENTIFIER.length()); - } + String metadataStoreUrl = removeIdentifierFromMetadataURL(internalConf.getMetadataStoreUrl()); dlogURI = WorkerUtils.newDlogNamespaceURI(metadataStoreUrl); } else { dlogURI = WorkerUtils.initializeDlogNamespace(internalConf); @@ -360,10 +357,7 @@ public void initInBroker(ServiceConfiguration brokerConfig, try { // initializing dlog namespace for function worker if (workerConfig.isInitializedDlogMetadata()) { - String metadataStoreUrl = internalConf.getMetadataStoreUrl(); - if (metadataStoreUrl.startsWith(ZK_SCHEME_IDENTIFIER)) { - metadataStoreUrl = metadataStoreUrl.substring(ZK_SCHEME_IDENTIFIER.length()); - } + String metadataStoreUrl = removeIdentifierFromMetadataURL(internalConf.getMetadataStoreUrl()); dlogURI = WorkerUtils.newDlogNamespaceURI(metadataStoreUrl); } else { dlogURI = WorkerUtils.initializeDlogNamespace(internalConf); diff --git a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java index 2f2c667ba47cb..8519f0ca9e4f6 100644 --- a/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java +++ b/pulsar-functions/worker/src/main/java/org/apache/pulsar/functions/worker/WorkerUtils.java @@ -20,7 +20,7 @@ import static java.nio.file.StandardCopyOption.REPLACE_EXISTING; import static org.apache.commons.lang3.StringUtils.isNotBlank; -import static org.apache.pulsar.metadata.impl.ZKMetadataStore.ZK_SCHEME_IDENTIFIER; +import static org.apache.pulsar.metadata.impl.MetadataStoreFactoryImpl.removeIdentifierFromMetadataURL; import java.io.File; import java.io.FileInputStream; import java.io.FileOutputStream; @@ -174,11 +174,7 @@ public static URI initializeDlogNamespace(InternalConfigurationData internalConf // for BC purposes if (internalConf.getBookkeeperMetadataServiceUri() == null) { ledgersRootPath = internalConf.getLedgersRootPath(); - String metadataStoreUrl = internalConf.getMetadataStoreUrl(); - if (metadataStoreUrl.startsWith(ZK_SCHEME_IDENTIFIER)) { - metadataStoreUrl = metadataStoreUrl.substring(ZK_SCHEME_IDENTIFIER.length()); - } - ledgersStoreServers = metadataStoreUrl; + ledgersStoreServers = removeIdentifierFromMetadataURL(internalConf.getMetadataStoreUrl()); chrootPath = ""; } else { URI metadataServiceUri = URI.create(internalConf.getBookkeeperMetadataServiceUri()); From a26f647978c6a4653bb9d7a16a98fcc260f2dc8d Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Mon, 28 Feb 2022 14:31:36 +0800 Subject: [PATCH 5/7] fix npe --- .../org/apache/pulsar/PulsarStandaloneStarter.java | 3 +++ .../storage/bookkeeper/BookKeeperPackagesStorage.java | 11 ++++++++--- .../BookKeeperPackagesStorageConfiguration.java | 4 ++++ 3 files changed, 15 insertions(+), 3 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java index cb72c679bcb7a..eb9b5978147a1 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java @@ -107,6 +107,9 @@ public PulsarStandaloneStarter(String[] args) throws Exception { config.setConfigurationMetadataStoreUrl(metadataStoreUrl); config.setRunningStandalone(true); + config.getProperties().setProperty("metadataStoreUrl", metadataStoreUrl); + config.getProperties().setProperty("configurationMetadataStoreUrl", metadataStoreUrl); + Runtime.getRuntime().addShutdownHook(new Thread() { public void run() { diff --git a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java index 1a78ec5ef71e6..68eb126815cbf 100644 --- a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java +++ b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java @@ -93,14 +93,19 @@ private URI initializeDlogNamespace() throws IOException { ledgersRootPath = metadataServiceUri.getPath(); } else { ledgersRootPath = configuration.getPackagesManagementLedgerRootPath(); - ledgersStoreServers = configuration.getMetadataStoreUrl(); - if (ledgersStoreServers.startsWith(ZK_SCHEME_IDENTIFIER)) { - ledgersStoreServers = ledgersStoreServers.substring(ZK_SCHEME_IDENTIFIER.length()); + if (configuration.getMetadataStoreUrl() != null) { + ledgersStoreServers = configuration.getMetadataStoreUrl(); + if (ledgersStoreServers.startsWith(ZK_SCHEME_IDENTIFIER)) { + ledgersStoreServers = ledgersStoreServers.substring(ZK_SCHEME_IDENTIFIER.length()); + } + } else { + ledgersStoreServers = configuration.getZookeeperServers(); } } BKDLConfig bkdlConfig = new BKDLConfig(ledgersStoreServers, ledgersRootPath); DLMetadata dlMetadata = DLMetadata.create(bkdlConfig); URI dlogURI = URI.create(String.format("distributedlog://%s/pulsar/packages", ledgersStoreServers)); + log.info("Tried to initial:{}", String.format("distributedlog://%s/pulsar/packages", ledgersStoreServers)); try { dlMetadata.create(dlogURI); } catch (ZKException e) { diff --git a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java index 00ce9b53fef95..5153c76ea0cee 100644 --- a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java +++ b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java @@ -42,6 +42,10 @@ String getMetadataStoreUrl() { return getProperty("metadataStoreUrl"); } + String getZookeeperServers() { + return getProperty("zookeeperServers"); + } + String getPackagesManagementLedgerRootPath() { return getProperty("packagesManagementLedgerRootPath"); } From f09866da5d8deddad743c92b4abf9d51d042072f Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Mon, 28 Feb 2022 17:23:18 +0800 Subject: [PATCH 6/7] fix unit test error --- .../org/apache/pulsar/PulsarStandaloneStarter.java | 3 --- .../bookkeeper/BookKeeperPackagesStorage.java | 14 +++----------- .../BookKeeperPackagesStorageConfiguration.java | 4 ---- .../bookkeeper/BookKeeperPackagesStorageTest.java | 11 ++++------- 4 files changed, 7 insertions(+), 25 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java index eb9b5978147a1..cb72c679bcb7a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneStarter.java @@ -107,9 +107,6 @@ public PulsarStandaloneStarter(String[] args) throws Exception { config.setConfigurationMetadataStoreUrl(metadataStoreUrl); config.setRunningStandalone(true); - config.getProperties().setProperty("metadataStoreUrl", metadataStoreUrl); - config.getProperties().setProperty("configurationMetadataStoreUrl", metadataStoreUrl); - Runtime.getRuntime().addShutdownHook(new Thread() { public void run() { diff --git a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java index 68eb126815cbf..3da138475ce8c 100644 --- a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java +++ b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java @@ -49,7 +49,6 @@ public class BookKeeperPackagesStorage implements PackagesStorage { private static final String NS_CLIENT_ID = "packages-management"; - public static final String ZK_SCHEME_IDENTIFIER = "zk:"; final BookKeeperPackagesStorageConfiguration configuration; private Namespace namespace; @@ -93,19 +92,12 @@ private URI initializeDlogNamespace() throws IOException { ledgersRootPath = metadataServiceUri.getPath(); } else { ledgersRootPath = configuration.getPackagesManagementLedgerRootPath(); - if (configuration.getMetadataStoreUrl() != null) { - ledgersStoreServers = configuration.getMetadataStoreUrl(); - if (ledgersStoreServers.startsWith(ZK_SCHEME_IDENTIFIER)) { - ledgersStoreServers = ledgersStoreServers.substring(ZK_SCHEME_IDENTIFIER.length()); - } - } else { - ledgersStoreServers = configuration.getZookeeperServers(); - } + ledgersStoreServers = configuration.getZookeeperServers(); } BKDLConfig bkdlConfig = new BKDLConfig(ledgersStoreServers, ledgersRootPath); DLMetadata dlMetadata = DLMetadata.create(bkdlConfig); - URI dlogURI = URI.create(String.format("distributedlog://%s/pulsar/packages", ledgersStoreServers)); - log.info("Tried to initial:{}", String.format("distributedlog://%s/pulsar/packages", ledgersStoreServers)); + URI dlogURI = URI.create(String.format("distributedlog://%s/pulsar/packages", + configuration.getZookeeperServers())); try { dlMetadata.create(dlogURI); } catch (ZKException e) { diff --git a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java index 5153c76ea0cee..226b80abeaa30 100644 --- a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java +++ b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java @@ -38,10 +38,6 @@ int getPackagesReplicas() { return Integer.parseInt(getProperty("packagesReplicas")); } - String getMetadataStoreUrl() { - return getProperty("metadataStoreUrl"); - } - String getZookeeperServers() { return getProperty("zookeeperServers"); } diff --git a/pulsar-package-management/bookkeeper-storage/src/test/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageTest.java b/pulsar-package-management/bookkeeper-storage/src/test/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageTest.java index 8715b35bc2809..69312076410ca 100644 --- a/pulsar-package-management/bookkeeper-storage/src/test/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageTest.java +++ b/pulsar-package-management/bookkeeper-storage/src/test/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageTest.java @@ -18,7 +18,6 @@ */ package org.apache.pulsar.packages.management.storage.bookkeeper; -import static org.apache.pulsar.metadata.impl.ZKMetadataStore.ZK_SCHEME_IDENTIFIER; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; @@ -51,7 +50,7 @@ public void setup() throws Exception { PackagesStorageProvider provider = PackagesStorageProvider .newProvider(BookKeeperPackagesStorageProvider.class.getName()); DefaultPackagesStorageConfiguration configuration = new DefaultPackagesStorageConfiguration(); - configuration.setProperty("metadataStoreUrl", ZK_SCHEME_IDENTIFIER + zkUtil.getZooKeeperConnectString()); + configuration.setProperty("zookeeperServers", zkUtil.getZooKeeperConnectString()); configuration.setProperty("packagesReplicas", "1"); configuration.setProperty("packagesManagementLedgerRootPath", "/ledgers"); storage = provider.getStorage(configuration); @@ -69,8 +68,7 @@ public void teardown() throws Exception { public void testConfiguration() { assertTrue(storage instanceof BookKeeperPackagesStorage); BookKeeperPackagesStorage bkStorage = (BookKeeperPackagesStorage) storage; - assertEquals(bkStorage.configuration.getMetadataStoreUrl() - .substring(ZK_SCHEME_IDENTIFIER.length()), zkUtil.getZooKeeperConnectString()); + assertEquals(bkStorage.configuration.getZookeeperServers(), zkUtil.getZooKeeperConnectString()); assertEquals(bkStorage.configuration.getPackagesReplicas(), 1); assertEquals(bkStorage.configuration.getPackagesManagementLedgerRootPath(), "/ledgers"); } @@ -200,8 +198,7 @@ public void testReadWriteOperationsWithSeparatedBkCluster() throws Exception { .newProvider(BookKeeperPackagesStorageProvider.class.getName()); DefaultPackagesStorageConfiguration configuration = new DefaultPackagesStorageConfiguration(); // set the unavailable bk cluster with mock zookeeper path - configuration.setProperty("metadataStoreUrl", ZK_SCHEME_IDENTIFIER + zkUtil.getZooKeeperConnectString() + - "/mock"); + configuration.setProperty("zookeeperServers", zkUtil.getZooKeeperConnectString() + "/mock"); configuration.setProperty("packagesReplicas", "1"); configuration.setProperty("packagesManagementLedgerRootPath", "/ledgers"); PackagesStorage storage1 = provider.getStorage(configuration); @@ -224,7 +221,7 @@ public void testReadWriteOperationsWithSeparatedBkCluster() throws Exception { // set the available bk cluster with bookkeeperMetadataServiceUri using actual zookeeper path String bookkeeperMetadataServiceUri = String.format("zk+null://%s/ledgers", zkUtil.getZooKeeperConnectString()); DefaultPackagesStorageConfiguration configuration2 = new DefaultPackagesStorageConfiguration(); - configuration2.setProperty("metadataStoreUrl", ZK_SCHEME_IDENTIFIER + zkUtil.getZooKeeperConnectString()); + configuration2.setProperty("zookeeperServers", zkUtil.getZooKeeperConnectString()); configuration2.setProperty("bookkeeperMetadataServiceUri", bookkeeperMetadataServiceUri); configuration2.setProperty("packagesReplicas", "1"); PackagesStorage storage2 = provider.getStorage(configuration2); From 41c7d07b2500cb42aacb949c1a6c6f17f4467a7b Mon Sep 17 00:00:00 2001 From: gavingaozhangmin Date: Mon, 28 Feb 2022 20:40:47 +0800 Subject: [PATCH 7/7] fix npe --- .../storage/bookkeeper/BookKeeperPackagesStorage.java | 10 +++++++++- .../BookKeeperPackagesStorageConfiguration.java | 4 ++++ 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java index 3da138475ce8c..45fce3a62bce9 100644 --- a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java +++ b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorage.java @@ -49,6 +49,7 @@ public class BookKeeperPackagesStorage implements PackagesStorage { private static final String NS_CLIENT_ID = "packages-management"; + public static final String ZK_SCHEME_IDENTIFIER = "zk:"; final BookKeeperPackagesStorageConfiguration configuration; private Namespace namespace; @@ -92,7 +93,14 @@ private URI initializeDlogNamespace() throws IOException { ledgersRootPath = metadataServiceUri.getPath(); } else { ledgersRootPath = configuration.getPackagesManagementLedgerRootPath(); - ledgersStoreServers = configuration.getZookeeperServers(); + if (StringUtils.isNotBlank(configuration.getMetadataStoreUrl())) { + ledgersStoreServers = configuration.getMetadataStoreUrl(); + if (ledgersStoreServers.startsWith(ZK_SCHEME_IDENTIFIER)) { + ledgersStoreServers = ledgersStoreServers.substring(ZK_SCHEME_IDENTIFIER.length()); + } + } else { + ledgersStoreServers = configuration.getZookeeperServers(); + } } BKDLConfig bkdlConfig = new BKDLConfig(ledgersStoreServers, ledgersRootPath); DLMetadata dlMetadata = DLMetadata.create(bkdlConfig); diff --git a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java index 226b80abeaa30..d0c00b800f45a 100644 --- a/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java +++ b/pulsar-package-management/bookkeeper-storage/src/main/java/org/apache/pulsar/packages/management/storage/bookkeeper/BookKeeperPackagesStorageConfiguration.java @@ -42,6 +42,10 @@ String getZookeeperServers() { return getProperty("zookeeperServers"); } + String getMetadataStoreUrl() { + return getProperty("metadataStoreUrl"); + } + String getPackagesManagementLedgerRootPath() { return getProperty("packagesManagementLedgerRootPath"); }