From 95e635d0ee0560c941e447138e7ea5cc8fee2b12 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Thu, 22 Apr 2021 00:37:28 +0800 Subject: [PATCH 1/9] Fix the inconsistency of AdvertisedAddress --- .../pulsar/broker/ServiceConfiguration.java | 20 +++++++++++++++++++ .../apache/pulsar/broker/PulsarService.java | 16 ++------------- 2 files changed, 22 insertions(+), 14 deletions(-) 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 358345c4cb824..8e2257849eaad 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 @@ -25,6 +25,7 @@ import java.util.ArrayList; import java.util.HashSet; import java.util.List; +import java.util.Map; import java.util.Optional; import java.util.Properties; import java.util.Set; @@ -33,6 +34,7 @@ import lombok.Setter; import org.apache.bookkeeper.client.api.DigestType; import org.apache.pulsar.broker.authorization.PulsarAuthorizationProvider; +import org.apache.pulsar.broker.validator.MultipleListenerValidator; import org.apache.pulsar.common.configuration.Category; import org.apache.pulsar.common.configuration.FieldContext; import org.apache.pulsar.common.configuration.PulsarConfiguration; @@ -43,6 +45,7 @@ import org.apache.pulsar.common.policies.data.TopicType; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.sasl.SaslConstants; +import org.apache.pulsar.policies.data.loadbalancer.AdvertisedListener; /** * Pulsar service configuration object. @@ -2216,4 +2219,21 @@ public int getBrokerDeleteInactiveTopicsMaxInactiveDurationSeconds() { return brokerDeleteInactiveTopicsMaxInactiveDurationSeconds; } } + + public String getAdvertisedAddress() { + Map result = MultipleListenerValidator + .validateAndAnalysisAdvertisedListener(this); + + if (advertisedAddress != null) { + return advertisedAddress; + } + + String address = result.get(internalListenerName).getBrokerServiceUrl().getHost(); + if (address != null) { + return address; + } + + return ServiceConfigurationUtils.getDefaultOrConfiguredAddress(advertisedAddress); + + } } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index 86630bb6bcb81..b4cf71c4a2805 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -271,22 +271,10 @@ public PulsarService(ServiceConfiguration config, // Validate correctness of configuration PulsarConfigurationLoader.isComplete(config); // validate `advertisedAddress`, `advertisedListeners`, `internalListenerName` - Map result = - MultipleListenerValidator.validateAndAnalysisAdvertisedListener(config); - if (result != null) { - this.advertisedListeners = Collections.unmodifiableMap(result); - } else { - this.advertisedListeners = Collections.unmodifiableMap(Collections.emptyMap()); - } + this.advertisedAddress = config.getAdvertisedAddress(); state = State.Init; // use `internalListenerName` listener as `advertisedAddress` this.bindAddress = ServiceConfigurationUtils.getDefaultOrConfiguredAddress(config.getBindAddress()); - if (!this.advertisedListeners.isEmpty()) { - this.advertisedAddress = this.advertisedListeners.get( - config.getInternalListenerName()).getBrokerServiceUrl().getHost(); - } else { - this.advertisedAddress = advertisedAddress(config); - } this.brokerVersion = PulsarVersion.getVersion(); this.config = config; this.shutdownService = new MessagingServiceShutdownHook(this, processTerminator); @@ -1332,7 +1320,7 @@ public ShutdownService getShutdownService() { * @return Hostname or IP address the service advertises to the outside world. */ public static String advertisedAddress(ServiceConfiguration config) { - return ServiceConfigurationUtils.getDefaultOrConfiguredAddress(config.getAdvertisedAddress()); + return config.getAdvertisedAddress(); } private String brokerUrl(ServiceConfiguration config) { From 724b8fd3da2675cbe26180674460f29cbfb99b69 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Thu, 22 Apr 2021 00:57:01 +0800 Subject: [PATCH 2/9] Remove method --- .../apache/pulsar/broker/PulsarService.java | 18 ++++-------------- .../impl/ModularLoadManagerImpl.java | 2 +- .../pulsar/compaction/CompactorTool.java | 4 ++-- 3 files changed, 7 insertions(+), 17 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index b4cf71c4a2805..e079ec1cc8ec5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -105,7 +105,6 @@ import org.apache.pulsar.broker.storage.ManagedLedgerStorage; import org.apache.pulsar.broker.transaction.buffer.TransactionBufferProvider; import org.apache.pulsar.broker.transaction.buffer.impl.TransactionBufferClientImpl; -import org.apache.pulsar.broker.validator.MultipleListenerValidator; import org.apache.pulsar.broker.web.WebService; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminBuilder; @@ -1314,18 +1313,9 @@ public ShutdownService getShutdownService() { return shutdownService; } - /** - * Advertised service address. - * - * @return Hostname or IP address the service advertises to the outside world. - */ - public static String advertisedAddress(ServiceConfiguration config) { - return config.getAdvertisedAddress(); - } - private String brokerUrl(ServiceConfiguration config) { if (config.getBrokerServicePort().isPresent()) { - return brokerUrl(advertisedAddress(config), getBrokerListenPort().get()); + return brokerUrl(config.getAdvertisedAddress(), getBrokerListenPort().get()); } else { return null; } @@ -1337,7 +1327,7 @@ public static String brokerUrl(String host, int port) { public String brokerUrlTls(ServiceConfiguration config) { if (config.getBrokerServicePortTls().isPresent()) { - return brokerUrlTls(advertisedAddress(config), getBrokerListenPortTls().get()); + return brokerUrlTls(config.getAdvertisedAddress(), getBrokerListenPortTls().get()); } else { return null; } @@ -1349,7 +1339,7 @@ public static String brokerUrlTls(String host, int port) { public String webAddress(ServiceConfiguration config) { if (config.getWebServicePort().isPresent()) { - return webAddress(advertisedAddress(config), getListenPortHTTP().get()); + return webAddress(config.getAdvertisedAddress(), getListenPortHTTP().get()); } else { return null; } @@ -1361,7 +1351,7 @@ public static String webAddress(String host, int port) { public String webAddressTls(ServiceConfiguration config) { if (config.getWebServicePortTls().isPresent()) { - return webAddressTls(advertisedAddress(config), getListenPortHTTPS().get()); + return webAddressTls(config.getAdvertisedAddress(), getListenPortHTTPS().get()); } else { return null; } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java index 9df80e25cff2f..a61cbc0b35558 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java @@ -974,7 +974,7 @@ private void updateLoadBalancingMetrics(final SystemResourceUsage systemResource List metrics = Lists.newArrayList(); Map dimensions = new HashMap<>(); - dimensions.put("broker", conf.getAdvertisedAddress()); + dimensions.put("broker", pulsar.getAdvertisedAddress()); dimensions.put("metric", "loadBalancing"); Metrics m = Metrics.create(dimensions); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java index c688549b7fea4..d263f7b2df16d 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java @@ -100,13 +100,13 @@ public static void main(String[] args) throws Exception { log.info("Found `brokerServicePortTls` in configuration file. \n" + "Will connect pulsar use TLS."); clientBuilder - .serviceUrl(PulsarService.brokerUrlTls(PulsarService.advertisedAddress(brokerConfig), + .serviceUrl(PulsarService.brokerUrlTls(brokerConfig.getAdvertisedAddress(), brokerConfig.getBrokerServicePortTls().get())) .allowTlsInsecureConnection(brokerConfig.isTlsAllowInsecureConnection()) .tlsTrustCertsFilePath(brokerConfig.getTlsCertificateFilePath()); } else { - clientBuilder.serviceUrl(PulsarService.brokerUrl(PulsarService.advertisedAddress(brokerConfig), + clientBuilder.serviceUrl(PulsarService.brokerUrl(brokerConfig.getAdvertisedAddress(), brokerConfig.getBrokerServicePort().get())); } From 4e3f23c5aca2593e8297e1440043e6e0a69d52c1 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Thu, 22 Apr 2021 13:10:36 +0800 Subject: [PATCH 3/9] fix unit test --- .../pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java index a61cbc0b35558..9df80e25cff2f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java @@ -974,7 +974,7 @@ private void updateLoadBalancingMetrics(final SystemResourceUsage systemResource List metrics = Lists.newArrayList(); Map dimensions = new HashMap<>(); - dimensions.put("broker", pulsar.getAdvertisedAddress()); + dimensions.put("broker", conf.getAdvertisedAddress()); dimensions.put("metric", "loadBalancing"); Metrics m = Metrics.create(dimensions); From 8fce08b003f180cfb9c6f100da991c5ae0494577 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Thu, 22 Apr 2021 13:19:15 +0800 Subject: [PATCH 4/9] restore field --- .../src/main/java/org/apache/pulsar/broker/PulsarService.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index e079ec1cc8ec5..b875423fa922c 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -105,6 +105,7 @@ import org.apache.pulsar.broker.storage.ManagedLedgerStorage; import org.apache.pulsar.broker.transaction.buffer.TransactionBufferProvider; import org.apache.pulsar.broker.transaction.buffer.impl.TransactionBufferClientImpl; +import org.apache.pulsar.broker.validator.MultipleListenerValidator; import org.apache.pulsar.broker.web.WebService; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminBuilder; @@ -270,6 +271,7 @@ public PulsarService(ServiceConfiguration config, // Validate correctness of configuration PulsarConfigurationLoader.isComplete(config); // validate `advertisedAddress`, `advertisedListeners`, `internalListenerName` + this.advertisedListeners = MultipleListenerValidator.validateAndAnalysisAdvertisedListener(config); this.advertisedAddress = config.getAdvertisedAddress(); state = State.Init; // use `internalListenerName` listener as `advertisedAddress` From 375df2f23d94b6cfa915f1ad6b881f058b8799e4 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Thu, 22 Apr 2021 15:33:56 +0800 Subject: [PATCH 5/9] fix test --- .../apache/pulsar/broker/ServiceConfiguration.java | 2 +- .../java/org/apache/pulsar/broker/PulsarService.java | 12 ++++++------ .../loadbalance/impl/ModularLoadManagerImpl.java | 2 +- .../org/apache/pulsar/compaction/CompactorTool.java | 4 ++-- 4 files changed, 10 insertions(+), 10 deletions(-) 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 8e2257849eaad..ca527bc899218 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 @@ -2220,7 +2220,7 @@ public int getBrokerDeleteInactiveTopicsMaxInactiveDurationSeconds() { } } - public String getAdvertisedAddress() { + public String getAppliedAdvertisedAddress() { Map result = MultipleListenerValidator .validateAndAnalysisAdvertisedListener(this); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index b875423fa922c..146c6f208d4bc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -272,7 +272,7 @@ public PulsarService(ServiceConfiguration config, PulsarConfigurationLoader.isComplete(config); // validate `advertisedAddress`, `advertisedListeners`, `internalListenerName` this.advertisedListeners = MultipleListenerValidator.validateAndAnalysisAdvertisedListener(config); - this.advertisedAddress = config.getAdvertisedAddress(); + this.advertisedAddress = config.getAppliedAdvertisedAddress(); state = State.Init; // use `internalListenerName` listener as `advertisedAddress` this.bindAddress = ServiceConfigurationUtils.getDefaultOrConfiguredAddress(config.getBindAddress()); @@ -1317,7 +1317,7 @@ public ShutdownService getShutdownService() { private String brokerUrl(ServiceConfiguration config) { if (config.getBrokerServicePort().isPresent()) { - return brokerUrl(config.getAdvertisedAddress(), getBrokerListenPort().get()); + return brokerUrl(config.getAppliedAdvertisedAddress(), getBrokerListenPort().get()); } else { return null; } @@ -1329,7 +1329,7 @@ public static String brokerUrl(String host, int port) { public String brokerUrlTls(ServiceConfiguration config) { if (config.getBrokerServicePortTls().isPresent()) { - return brokerUrlTls(config.getAdvertisedAddress(), getBrokerListenPortTls().get()); + return brokerUrlTls(config.getAppliedAdvertisedAddress(), getBrokerListenPortTls().get()); } else { return null; } @@ -1341,7 +1341,7 @@ public static String brokerUrlTls(String host, int port) { public String webAddress(ServiceConfiguration config) { if (config.getWebServicePort().isPresent()) { - return webAddress(config.getAdvertisedAddress(), getListenPortHTTP().get()); + return webAddress(config.getAppliedAdvertisedAddress(), getListenPortHTTP().get()); } else { return null; } @@ -1353,7 +1353,7 @@ public static String webAddress(String host, int port) { public String webAddressTls(ServiceConfiguration config) { if (config.getWebServicePortTls().isPresent()) { - return webAddressTls(config.getAdvertisedAddress(), getListenPortHTTPS().get()); + return webAddressTls(config.getAppliedAdvertisedAddress(), getListenPortHTTPS().get()); } else { return null; } @@ -1494,7 +1494,7 @@ public static WorkerConfig initializeWorkerConfigFromBrokerConfig(ServiceConfigu // worker talks to local broker String hostname = ServiceConfigurationUtils.getDefaultOrConfiguredAddress( - brokerConfig.getAdvertisedAddress()); + brokerConfig.getAppliedAdvertisedAddress()); workerConfig.setWorkerHostname(hostname); // inherit broker authorization setting workerConfig.setAuthenticationEnabled(brokerConfig.isAuthenticationEnabled()); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java index 9df80e25cff2f..c5126e1ef81c1 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java @@ -974,7 +974,7 @@ private void updateLoadBalancingMetrics(final SystemResourceUsage systemResource List metrics = Lists.newArrayList(); Map dimensions = new HashMap<>(); - dimensions.put("broker", conf.getAdvertisedAddress()); + dimensions.put("broker", conf.getAppliedAdvertisedAddress()); dimensions.put("metric", "loadBalancing"); Metrics m = Metrics.create(dimensions); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java index d263f7b2df16d..b45683bcd5934 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java @@ -100,13 +100,13 @@ public static void main(String[] args) throws Exception { log.info("Found `brokerServicePortTls` in configuration file. \n" + "Will connect pulsar use TLS."); clientBuilder - .serviceUrl(PulsarService.brokerUrlTls(brokerConfig.getAdvertisedAddress(), + .serviceUrl(PulsarService.brokerUrlTls(brokerConfig.getAppliedAdvertisedAddress(), brokerConfig.getBrokerServicePortTls().get())) .allowTlsInsecureConnection(brokerConfig.isTlsAllowInsecureConnection()) .tlsTrustCertsFilePath(brokerConfig.getTlsCertificateFilePath()); } else { - clientBuilder.serviceUrl(PulsarService.brokerUrl(brokerConfig.getAdvertisedAddress(), + clientBuilder.serviceUrl(PulsarService.brokerUrl(brokerConfig.getAppliedAdvertisedAddress(), brokerConfig.getBrokerServicePort().get())); } From f41a9449359ac178b803180de7839ab7d026d418 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Thu, 22 Apr 2021 15:47:05 +0800 Subject: [PATCH 6/9] add unit test --- .../pulsar/broker/ServiceConfiguration.java | 9 +++++--- .../MultipleListenerValidatorTest.java | 21 +++++++++++++++++++ 2 files changed, 27 insertions(+), 3 deletions(-) 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 ca527bc899218..9a6a76019122e 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 @@ -2228,9 +2228,12 @@ public String getAppliedAdvertisedAddress() { return advertisedAddress; } - String address = result.get(internalListenerName).getBrokerServiceUrl().getHost(); - if (address != null) { - return address; + AdvertisedListener advertisedListener = result.get(internalListenerName); + if (advertisedListener != null) { + String address = advertisedListener.getBrokerServiceUrl().getHost(); + if (address != null) { + return address; + } } return ServiceConfigurationUtils.getDefaultOrConfiguredAddress(advertisedAddress); diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/validator/MultipleListenerValidatorTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/validator/MultipleListenerValidatorTest.java index 89b72179bde76..064d86946f532 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/validator/MultipleListenerValidatorTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/validator/MultipleListenerValidatorTest.java @@ -19,10 +19,13 @@ package org.apache.pulsar.broker.validator; import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.ServiceConfigurationUtils; import org.testng.annotations.Test; import java.util.Optional; +import static org.testng.AssertJUnit.assertEquals; + /** * testcase for MultipleListenerValidator. */ @@ -37,6 +40,24 @@ public void testAppearTogether() { MultipleListenerValidator.validateAndAnalysisAdvertisedListener(config); } + @Test + public void testGetAppliedAdvertised() throws Exception { + ServiceConfiguration config = new ServiceConfiguration(); + config.setBrokerServicePortTls(Optional.of(6651)); + config.setAdvertisedListeners("internal:pulsar://192.0.0.1:6660, internal:pulsar+ssl://192.0.0.1:6651"); + config.setInternalListenerName("internal"); + assertEquals(config.getAppliedAdvertisedAddress(), "192.0.0.1"); + + config = new ServiceConfiguration(); + config.setBrokerServicePortTls(Optional.of(6651)); + config.setAdvertisedAddress("192.0.0.2"); + assertEquals(config.getAppliedAdvertisedAddress(), "192.0.0.2"); + + config.setAdvertisedAddress(null); + assertEquals(config.getAppliedAdvertisedAddress(), ServiceConfigurationUtils.getDefaultOrConfiguredAddress(null)); + } + + @Test(expectedExceptions = IllegalArgumentException.class) public void testListenerDuplicate_1() { ServiceConfiguration config = new ServiceConfiguration(); From d1199ba670f7806cb8af565eae03d3a8d6f7ba2f Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 7 May 2021 18:31:44 +0800 Subject: [PATCH 7/9] add comment line --- .../java/org/apache/pulsar/broker/ServiceConfiguration.java | 6 ++++++ 1 file changed, 6 insertions(+) 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 9a6a76019122e..45fb5ad30bb09 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 @@ -2220,6 +2220,12 @@ public int getBrokerDeleteInactiveTopicsMaxInactiveDurationSeconds() { } } + /** + * Get the address of Broker, first try to get it from AdvertisedAddress. + * If it is not set, try to get the address set by advertisedListener. + * If it is still not set, get it through InetAddress.getLocalHost(). + * @return + */ public String getAppliedAdvertisedAddress() { Map result = MultipleListenerValidator .validateAndAnalysisAdvertisedListener(this); From 71545adcd167d89147e262af722cae8ed9790cb2 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 7 May 2021 21:12:15 +0800 Subject: [PATCH 8/9] move to util --- .../pulsar/broker/ServiceConfiguration.java | 25 ---------------- .../broker/ServiceConfigurationUtils.java | 29 +++++++++++++++++++ .../MultipleListenerValidatorTest.java | 7 +++-- .../apache/pulsar/broker/PulsarService.java | 16 ++++++---- .../impl/ModularLoadManagerImpl.java | 3 +- .../pulsar/compaction/CompactorTool.java | 7 +++-- 6 files changed, 50 insertions(+), 37 deletions(-) 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 45fb5ad30bb09..ca95b5a42e49c 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 @@ -2220,29 +2220,4 @@ public int getBrokerDeleteInactiveTopicsMaxInactiveDurationSeconds() { } } - /** - * Get the address of Broker, first try to get it from AdvertisedAddress. - * If it is not set, try to get the address set by advertisedListener. - * If it is still not set, get it through InetAddress.getLocalHost(). - * @return - */ - public String getAppliedAdvertisedAddress() { - Map result = MultipleListenerValidator - .validateAndAnalysisAdvertisedListener(this); - - if (advertisedAddress != null) { - return advertisedAddress; - } - - AdvertisedListener advertisedListener = result.get(internalListenerName); - if (advertisedListener != null) { - String address = advertisedListener.getBrokerServiceUrl().getHost(); - if (address != null) { - return address; - } - } - - return ServiceConfigurationUtils.getDefaultOrConfiguredAddress(advertisedAddress); - - } } diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfigurationUtils.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfigurationUtils.java index 447b86207f02a..fe23e6b363a41 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfigurationUtils.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfigurationUtils.java @@ -18,11 +18,14 @@ */ package org.apache.pulsar.broker; +import org.apache.pulsar.broker.validator.MultipleListenerValidator; +import org.apache.pulsar.policies.data.loadbalancer.AdvertisedListener; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.net.InetAddress; import java.net.UnknownHostException; +import java.util.Map; import static org.apache.commons.lang3.StringUtils.isBlank; @@ -47,4 +50,30 @@ public static String unsafeLocalhostResolve() { } } + /** + * Get the address of Broker, first try to get it from AdvertisedAddress. + * If it is not set, try to get the address set by advertisedListener. + * If it is still not set, get it through InetAddress.getLocalHost(). + * @return + */ + public static String getAppliedAdvertisedAddress(ServiceConfiguration configuration) { + Map result = MultipleListenerValidator + .validateAndAnalysisAdvertisedListener(configuration); + + String advertisedAddress = configuration.getAdvertisedAddress(); + if (advertisedAddress != null) { + return advertisedAddress; + } + + AdvertisedListener advertisedListener = result.get(configuration.getInternalListenerName()); + if (advertisedListener != null) { + String address = advertisedListener.getBrokerServiceUrl().getHost(); + if (address != null) { + return address; + } + } + + return getDefaultOrConfiguredAddress(advertisedAddress); + } + } diff --git a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/validator/MultipleListenerValidatorTest.java b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/validator/MultipleListenerValidatorTest.java index 064d86946f532..8928e8223eeb3 100644 --- a/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/validator/MultipleListenerValidatorTest.java +++ b/pulsar-broker-common/src/test/java/org/apache/pulsar/broker/validator/MultipleListenerValidatorTest.java @@ -46,15 +46,16 @@ public void testGetAppliedAdvertised() throws Exception { config.setBrokerServicePortTls(Optional.of(6651)); config.setAdvertisedListeners("internal:pulsar://192.0.0.1:6660, internal:pulsar+ssl://192.0.0.1:6651"); config.setInternalListenerName("internal"); - assertEquals(config.getAppliedAdvertisedAddress(), "192.0.0.1"); + assertEquals(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), "192.0.0.1"); config = new ServiceConfiguration(); config.setBrokerServicePortTls(Optional.of(6651)); config.setAdvertisedAddress("192.0.0.2"); - assertEquals(config.getAppliedAdvertisedAddress(), "192.0.0.2"); + assertEquals(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), "192.0.0.2"); config.setAdvertisedAddress(null); - assertEquals(config.getAppliedAdvertisedAddress(), ServiceConfigurationUtils.getDefaultOrConfiguredAddress(null)); + assertEquals(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), + ServiceConfigurationUtils.getDefaultOrConfiguredAddress(null)); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index 146c6f208d4bc..078a5f9533610 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -272,7 +272,7 @@ public PulsarService(ServiceConfiguration config, PulsarConfigurationLoader.isComplete(config); // validate `advertisedAddress`, `advertisedListeners`, `internalListenerName` this.advertisedListeners = MultipleListenerValidator.validateAndAnalysisAdvertisedListener(config); - this.advertisedAddress = config.getAppliedAdvertisedAddress(); + this.advertisedAddress = ServiceConfigurationUtils.getAppliedAdvertisedAddress(config); state = State.Init; // use `internalListenerName` listener as `advertisedAddress` this.bindAddress = ServiceConfigurationUtils.getDefaultOrConfiguredAddress(config.getBindAddress()); @@ -1317,7 +1317,8 @@ public ShutdownService getShutdownService() { private String brokerUrl(ServiceConfiguration config) { if (config.getBrokerServicePort().isPresent()) { - return brokerUrl(config.getAppliedAdvertisedAddress(), getBrokerListenPort().get()); + return brokerUrl(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), + getBrokerListenPort().get()); } else { return null; } @@ -1329,7 +1330,8 @@ public static String brokerUrl(String host, int port) { public String brokerUrlTls(ServiceConfiguration config) { if (config.getBrokerServicePortTls().isPresent()) { - return brokerUrlTls(config.getAppliedAdvertisedAddress(), getBrokerListenPortTls().get()); + return brokerUrlTls(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), + getBrokerListenPortTls().get()); } else { return null; } @@ -1341,7 +1343,8 @@ public static String brokerUrlTls(String host, int port) { public String webAddress(ServiceConfiguration config) { if (config.getWebServicePort().isPresent()) { - return webAddress(config.getAppliedAdvertisedAddress(), getListenPortHTTP().get()); + return webAddress(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), + getListenPortHTTP().get()); } else { return null; } @@ -1353,7 +1356,8 @@ public static String webAddress(String host, int port) { public String webAddressTls(ServiceConfiguration config) { if (config.getWebServicePortTls().isPresent()) { - return webAddressTls(config.getAppliedAdvertisedAddress(), getListenPortHTTPS().get()); + return webAddressTls(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), + getListenPortHTTPS().get()); } else { return null; } @@ -1494,7 +1498,7 @@ public static WorkerConfig initializeWorkerConfigFromBrokerConfig(ServiceConfigu // worker talks to local broker String hostname = ServiceConfigurationUtils.getDefaultOrConfiguredAddress( - brokerConfig.getAppliedAdvertisedAddress()); + ServiceConfigurationUtils.getAppliedAdvertisedAddress(brokerConfig)); workerConfig.setWorkerHostname(hostname); // inherit broker authorization setting workerConfig.setAuthenticationEnabled(brokerConfig.isAuthenticationEnabled()); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java index c5126e1ef81c1..b7caab323c83a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/loadbalance/impl/ModularLoadManagerImpl.java @@ -47,6 +47,7 @@ import org.apache.pulsar.broker.PulsarServerException; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.ServiceConfigurationUtils; import org.apache.pulsar.broker.TimeAverageBrokerData; import org.apache.pulsar.broker.TimeAverageMessageData; import org.apache.pulsar.broker.loadbalance.BrokerFilter; @@ -974,7 +975,7 @@ private void updateLoadBalancingMetrics(final SystemResourceUsage systemResource List metrics = Lists.newArrayList(); Map dimensions = new HashMap<>(); - dimensions.put("broker", conf.getAppliedAdvertisedAddress()); + dimensions.put("broker", ServiceConfigurationUtils.getAppliedAdvertisedAddress(conf)); dimensions.put("metric", "loadBalancing"); Metrics m = Metrics.create(dimensions); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java index b45683bcd5934..26f32d9955a0e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/CompactorTool.java @@ -33,6 +33,7 @@ import org.apache.pulsar.broker.BookKeeperClientFactoryImpl; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; +import org.apache.pulsar.broker.ServiceConfigurationUtils; import org.apache.pulsar.client.api.ClientBuilder; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.common.configuration.PulsarConfigurationLoader; @@ -100,13 +101,15 @@ public static void main(String[] args) throws Exception { log.info("Found `brokerServicePortTls` in configuration file. \n" + "Will connect pulsar use TLS."); clientBuilder - .serviceUrl(PulsarService.brokerUrlTls(brokerConfig.getAppliedAdvertisedAddress(), + .serviceUrl(PulsarService.brokerUrlTls(ServiceConfigurationUtils + .getAppliedAdvertisedAddress(brokerConfig), brokerConfig.getBrokerServicePortTls().get())) .allowTlsInsecureConnection(brokerConfig.isTlsAllowInsecureConnection()) .tlsTrustCertsFilePath(brokerConfig.getTlsCertificateFilePath()); } else { - clientBuilder.serviceUrl(PulsarService.brokerUrl(brokerConfig.getAppliedAdvertisedAddress(), + clientBuilder.serviceUrl(PulsarService.brokerUrl(ServiceConfigurationUtils + .getAppliedAdvertisedAddress(brokerConfig), brokerConfig.getBrokerServicePort().get())); } From 3654795f3e66d1d3fd9e2594fea6ddab1f45e404 Mon Sep 17 00:00:00 2001 From: feynmanlin Date: Fri, 7 May 2021 22:56:25 +0800 Subject: [PATCH 9/9] add more unit test --- .../apache/pulsar/broker/PulsarService.java | 2 +- .../pulsar/broker/PulsarServiceTest.java | 50 ++++++++++++++++++- 2 files changed, 50 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java index 078a5f9533610..8f9bdd1ade469 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/PulsarService.java @@ -1315,7 +1315,7 @@ public ShutdownService getShutdownService() { return shutdownService; } - private String brokerUrl(ServiceConfiguration config) { + protected String brokerUrl(ServiceConfiguration config) { if (config.getBrokerServicePort().isPresent()) { return brokerUrl(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), getBrokerListenPort().get()); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceTest.java index a326d572447cf..9da4aa7492a5d 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/PulsarServiceTest.java @@ -21,17 +21,48 @@ import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertNull; import static org.testng.AssertJUnit.assertSame; import java.util.Optional; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; +import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest; import org.apache.pulsar.functions.worker.WorkerConfig; import org.apache.pulsar.functions.worker.WorkerService; +import org.testng.AssertJUnit; +import org.testng.annotations.AfterMethod; import org.testng.annotations.Test; @Slf4j -public class PulsarServiceTest { +public class PulsarServiceTest extends MockedPulsarServiceBaseTest { + + private boolean useListenerName = false; + + @Override + protected void setup() throws Exception { + super.internalSetup(); + } + + @AfterMethod(alwaysRun = true) + @Override + protected void cleanup() throws Exception { + super.internalCleanup(); + useListenerName = false; + resetConfig(); + } + + @Override + protected void doInitConf() throws Exception { + super.doInitConf(); + if (useListenerName) { + conf.setAdvertisedAddress(null); + conf.setBrokerServicePortTls(Optional.of(6651)); + conf.setBrokerServicePort(Optional.of(6660)); + conf.setWebServicePort(Optional.of(8081)); + conf.setWebServicePortTls(Optional.of(8082)); + } + } @Test public void testGetWorkerService() throws Exception { @@ -71,4 +102,21 @@ public void testGetWorkerServiceException() throws Exception { assertEquals(e.getMessage(), errorMessage); } } + + @Test + public void testAppliedAdvertised() throws Exception { + useListenerName = true; + conf.setAdvertisedListeners("internal:pulsar://127.0.0.1, internal:pulsar+ssl://127.0.0.1"); + conf.setInternalListenerName("internal"); + setup(); + + AssertJUnit.assertEquals(pulsar.getAdvertisedAddress(), "127.0.0.1"); + assertNull(pulsar.getConfiguration().getAdvertisedAddress()); + assertEquals(conf, pulsar.getConfiguration()); + assertEquals(pulsar.brokerUrlTls(conf), "pulsar+ssl://127.0.0.1:6651"); + assertEquals(pulsar.brokerUrl(conf), "pulsar://127.0.0.1:6660"); + assertEquals(pulsar.webAddress(conf), "http://127.0.0.1:8081"); + assertEquals(pulsar.webAddressTls(conf), "https://127.0.0.1:8082"); + } + }