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..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 @@ -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,5 @@ public int getBrokerDeleteInactiveTopicsMaxInactiveDurationSeconds() { return brokerDeleteInactiveTopicsMaxInactiveDurationSeconds; } } + } 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 89b72179bde76..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 @@ -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,25 @@ 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(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), "192.0.0.1"); + + config = new ServiceConfiguration(); + config.setBrokerServicePortTls(Optional.of(6651)); + config.setAdvertisedAddress("192.0.0.2"); + assertEquals(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), "192.0.0.2"); + + config.setAdvertisedAddress(null); + assertEquals(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), + ServiceConfigurationUtils.getDefaultOrConfiguredAddress(null)); + } + + @Test(expectedExceptions = IllegalArgumentException.class) public void testListenerDuplicate_1() { ServiceConfiguration config = new ServiceConfiguration(); 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..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 @@ -271,22 +271,11 @@ 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.advertisedListeners = MultipleListenerValidator.validateAndAnalysisAdvertisedListener(config); + this.advertisedAddress = ServiceConfigurationUtils.getAppliedAdvertisedAddress(config); 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); @@ -1326,18 +1315,10 @@ 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 ServiceConfigurationUtils.getDefaultOrConfiguredAddress(config.getAdvertisedAddress()); - } - - private String brokerUrl(ServiceConfiguration config) { + protected String brokerUrl(ServiceConfiguration config) { if (config.getBrokerServicePort().isPresent()) { - return brokerUrl(advertisedAddress(config), getBrokerListenPort().get()); + return brokerUrl(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), + getBrokerListenPort().get()); } else { return null; } @@ -1349,7 +1330,8 @@ 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(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), + getBrokerListenPortTls().get()); } else { return null; } @@ -1361,7 +1343,8 @@ 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(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), + getListenPortHTTP().get()); } else { return null; } @@ -1373,7 +1356,8 @@ 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(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), + getListenPortHTTPS().get()); } else { return null; } @@ -1514,7 +1498,7 @@ public static WorkerConfig initializeWorkerConfigFromBrokerConfig(ServiceConfigu // worker talks to local broker String hostname = ServiceConfigurationUtils.getDefaultOrConfiguredAddress( - brokerConfig.getAdvertisedAddress()); + 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 9df80e25cff2f..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.getAdvertisedAddress()); + 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 c688549b7fea4..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(PulsarService.advertisedAddress(brokerConfig), + .serviceUrl(PulsarService.brokerUrlTls(ServiceConfigurationUtils + .getAppliedAdvertisedAddress(brokerConfig), brokerConfig.getBrokerServicePortTls().get())) .allowTlsInsecureConnection(brokerConfig.isTlsAllowInsecureConnection()) .tlsTrustCertsFilePath(brokerConfig.getTlsCertificateFilePath()); } else { - clientBuilder.serviceUrl(PulsarService.brokerUrl(PulsarService.advertisedAddress(brokerConfig), + clientBuilder.serviceUrl(PulsarService.brokerUrl(ServiceConfigurationUtils + .getAppliedAdvertisedAddress(brokerConfig), brokerConfig.getBrokerServicePort().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"); + } + }