diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandalone.java b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandalone.java index cf4c039564850..334a8830a6c10 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandalone.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandalone.java @@ -25,6 +25,7 @@ import com.google.common.collect.Sets; import io.netty.util.internal.PlatformDependent; import java.io.File; +import java.io.IOException; import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.Paths; @@ -32,12 +33,15 @@ import java.util.Optional; import lombok.CustomLog; import org.apache.bookkeeper.conf.ServerConfiguration; +import org.apache.commons.io.FileUtils; import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.resources.NamespaceResources; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.partition.PartitionedTopicMetadata; @@ -60,6 +64,7 @@ @CustomLog @Command(name = "standalone", showDefaultValues = true, scope = ScopeType.INHERIT) public class PulsarStandalone implements AutoCloseable { + private Path tempBaseDir; private static final String PULSAR_STANDALONE_USE_ZOOKEEPER = "PULSAR_STANDALONE_USE_ZOOKEEPER"; @@ -216,6 +221,25 @@ public boolean isHelp() { return help; } + public void setTempBaseDir(Path tempBaseDir) { + this.tempBaseDir = tempBaseDir; + } + + public Path getTempBaseDir() { + return this.tempBaseDir; + } + + public PulsarClient buildClient() throws PulsarClientException { + return PulsarClient.builder() + .serviceUrl(this.getBrokerServiceUrl()) + .build(); + } + + public PulsarAdmin buildAdmin() throws PulsarClientException { + return PulsarAdmin.builder() + .serviceHttpUrl(this.getWebServiceUrl()) + .build(); + } @Option(names = { "-c", "--config" }, description = "Configuration file path") private String configFile; @@ -451,6 +475,21 @@ public void close() { } } catch (Exception e) { log.error().exception(e).log("Shutdown failed"); + } finally { + deleteTempBaseDir(); + } + } + + private void deleteTempBaseDir() { + if (tempBaseDir == null) { + return; + } + try { + FileUtils.deleteDirectory(tempBaseDir.toFile()); + } catch (IOException e) { + log.error().exception(e).log("Failed to delete temp directory " + tempBaseDir); + } finally { + tempBaseDir = null; } } @@ -469,7 +508,11 @@ void startBookieWithMetadataStore() throws Exception { } ServerConfiguration bkServerConf = new ServerConfiguration(); - bkServerConf.loadConf(new File(configFile).toURI().toURL()); + if (StringUtils.isNotBlank(configFile)) { + bkServerConf.loadConf(new File(configFile).toURI().toURL()); + } else { + bkServerConf.setAllowLoopback(true); + } calculateCacheSize(bkServerConf); bkCluster = BKCluster.builder() .baseServerConfiguration(bkServerConf) 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 2cf91c564d5cf..78eb80af5ab68 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneBuilder.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/PulsarStandaloneBuilder.java @@ -19,10 +19,15 @@ package org.apache.pulsar; import static org.apache.commons.lang3.StringUtils.isBlank; +import java.io.IOException; +import java.io.UncheckedIOException; +import java.nio.file.Files; +import java.nio.file.Path; import org.apache.pulsar.broker.ServiceConfiguration; import org.apache.pulsar.broker.ServiceConfigurationUtils; public final class PulsarStandaloneBuilder { + private Path tempBaseDir; private PulsarStandalone pulsarStandalone; @@ -32,6 +37,18 @@ private PulsarStandaloneBuilder() { pulsarStandalone.setNoFunctionsWorker(true); } + public PulsarStandaloneBuilder withTempDirectory() { + try { + tempBaseDir = Files.createTempDirectory("pulsar-standalone"); + } catch (IOException e) { + throw new UncheckedIOException(e); + } + pulsarStandalone.setZkDir(tempBaseDir.resolve("zookeeper").toString()); + pulsarStandalone.setBkDir(tempBaseDir.resolve("bookkeeper").toString()); + pulsarStandalone.setTempBaseDir(tempBaseDir); + return this; + } + public static PulsarStandaloneBuilder instance() { return new PulsarStandaloneBuilder(); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/PulsarStandaloneBuilderTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/PulsarStandaloneBuilderTest.java new file mode 100644 index 0000000000000..bb4534ddb73a5 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/PulsarStandaloneBuilderTest.java @@ -0,0 +1,64 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.pulsar; + +import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertFalse; +import static org.testng.AssertJUnit.assertNotNull; +import static org.testng.AssertJUnit.assertTrue; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.concurrent.TimeUnit; +import org.apache.pulsar.client.api.Consumer; +import org.apache.pulsar.client.api.Message; +import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.Schema; +import org.testng.annotations.Test; + +public class PulsarStandaloneBuilderTest { + @Test + public void testStartFromJava() throws Exception { + PulsarStandalone standalone = PulsarStandaloneBuilder.instance() + .withTempDirectory() + .build(); + Path tempBaseDir = standalone.getTempBaseDir(); + try { + standalone.setNumOfBk(2); + standalone.start(); + assertTrue(Files.exists(tempBaseDir)); + try (PulsarClient client = standalone.buildClient()) { + Producer producer = client.newProducer(Schema.STRING) + .topic("test-topic").create(); + Consumer consumer = client.newConsumer(Schema.STRING) + .topic("test-topic").subscriptionName("sub").subscribe(); + + producer.send("hello"); + Message msg = consumer.receive(10, TimeUnit.SECONDS); + assertNotNull(msg); + assertEquals(msg.getValue(), "hello"); + } + + } finally { + standalone.close(); + } + assertFalse(Files.exists(tempBaseDir), "Temp directory was not cleaned up: " + tempBaseDir); + } +}