From 6ca5a042ee7ce61f4c8eaa62e05b314330ee4cec Mon Sep 17 00:00:00 2001 From: Sijie Guo Date: Tue, 9 Apr 2019 23:07:53 +0800 Subject: [PATCH 1/7] Adding offloader support for sql --- conf/presto/catalog/pulsar.properties | 6 + conf/presto/config.properties | 2 +- conf/standalone.conf | 6 + .../mledger/offload/OffloaderUtils.java | 9 +- .../pulsar/common/nar/NarClassLoader.java | 16 ++- pulsar-sql/presto-pulsar-plugin/pom.xml | 1 - pulsar-sql/presto-pulsar/pom.xml | 104 +++++++++++++----- .../pulsar/sql/presto/AvroSchemaHandler.java | 8 +- .../pulsar/sql/presto/JSONSchemaHandler.java | 4 +- .../sql/presto/PulsarConnectorCache.java | 99 ++++++++++++++++- .../sql/presto/PulsarConnectorConfig.java | 52 ++++++++- .../presto/PulsarConnectorMetricsTracker.java | 10 +- .../sql/presto/PulsarConnectorUtils.java | 10 ++ .../pulsar/sql/presto/PulsarRecordCursor.java | 24 ++-- .../pulsar/sql/presto/PulsarSplitManager.java | 17 +-- .../pulsar/sql/presto/SchemaHandler.java | 2 +- .../sql/presto/TestPulsarConnector.java | 29 ++--- .../pulsar/sql/presto/TestPulsarMetadata.java | 4 +- 18 files changed, 320 insertions(+), 83 deletions(-) diff --git a/conf/presto/catalog/pulsar.properties b/conf/presto/catalog/pulsar.properties index 77b22dc8bc15c..35d83500039f5 100644 --- a/conf/presto/catalog/pulsar.properties +++ b/conf/presto/catalog/pulsar.properties @@ -31,3 +31,9 @@ pulsar.target-num-splits=2 pulsar.max-split-message-queue-size=10000 # max entry queue size pulsar.max-split-entry-queue-size = 1000 + + +####### TIERED STORAGE OFFLOADER CONFIGS ####### +pulsar.managed-ledger-offload-driver = aws-s3 +pulsar.offloaders-directory = /Users/jerrypeng/workspace/incubator-pulsar/offloaders +pulsar.offloader-properties = {"s3ManagedLedgerOffloadBucket": "jerry-pulsar-test", "s3ManagedLedgerOffloadRegion": "us-west-2", "s3ManagedLedgerOffloadServiceEndpoint": "http://s3.amazonaws.com"} \ No newline at end of file diff --git a/conf/presto/config.properties b/conf/presto/config.properties index 9f17135523dfa..0d54d241dead6 100644 --- a/conf/presto/config.properties +++ b/conf/presto/config.properties @@ -38,5 +38,5 @@ query.client.timeout=5m query.min-expire-age=30m presto.version=testversion -distributed-joins-enabled=true +#distributed-joins-enabled=true node-scheduler.include-coordinator=true diff --git a/conf/standalone.conf b/conf/standalone.conf index 4cbf931c75748..5542eb1704ac7 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -17,6 +17,12 @@ # under the License. # +managedLedgerOffloadDriver=aws-s3 +s3ManagedLedgerOffloadBucket=jerry-pulsar-test +s3ManagedLedgerOffloadRegion=us-west-2 +//s3ManagedLedgerOffloadServiceEndpoint=https://apigateway.us-west-2.amazonaws.com +s3ManagedLedgerOffloadServiceEndpoint=http://s3.amazonaws.com + ### --- General broker settings --- ### # Zookeeper quorum connection string diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java index 8726704147d86..42ecb6966693f 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java @@ -20,6 +20,7 @@ import java.io.File; import java.io.IOException; +import java.net.URLClassLoader; import java.nio.file.DirectoryStream; import java.nio.file.Files; import java.nio.file.Path; @@ -49,7 +50,10 @@ public class OffloaderUtils { * @throws IOException when fail to retrieve the pulsar offloader class */ static Pair getOffloaderFactory(String narPath) throws IOException { - NarClassLoader ncl = NarClassLoader.getFromArchive(new File(narPath), Collections.emptySet()); + // need to load offloader NAR to the classloader that also loaded LedgerOffloaderFactory in case + // LedgerOffloaderFactory is loaded by a classloader that is not the default classloader + // as is the case for the pulsar presto plugin + NarClassLoader ncl = NarClassLoader.getFromArchive(new File(narPath), Collections.emptySet(), LedgerOffloaderFactory.class.getClassLoader()); String configStr = ncl.getServiceDefinition(PULSAR_OFFLOADER_SERVICE_NAME); OffloaderDefinition conf = ObjectMapperFactory.getThreadLocalYaml() @@ -66,10 +70,9 @@ static Pair getOffloaderFactory(String n CompletableFuture loadFuture = new CompletableFuture<>(); Thread loadingThread = new Thread(() -> { Thread.currentThread().setContextClassLoader(ncl); - - log.info("Loading offloader factory {} using class loader {}", factoryClass, ncl); try { Object offloader = factoryClass.newInstance(); + if (!(offloader instanceof LedgerOffloaderFactory)) { throw new IOException("Class " + conf.getOffloaderFactoryClass() + " does not implement interface " + LedgerOffloaderFactory.class.getName()); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/nar/NarClassLoader.java b/pulsar-common/src/main/java/org/apache/pulsar/common/nar/NarClassLoader.java index 1ba7ae775144c..0b78344502b3f 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/nar/NarClassLoader.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/nar/NarClassLoader.java @@ -144,7 +144,16 @@ public boolean accept(File pathname) { public static NarClassLoader getFromArchive(File narPath, Set additionalJars) throws IOException { File unpacked = NarUnpacker.unpackNar(narPath, NAR_CACHE_DIR); try { - return new NarClassLoader(unpacked, additionalJars); + return new NarClassLoader(unpacked, additionalJars, NarClassLoader.class.getClassLoader() ); + } catch (ClassNotFoundException e) { + throw new IOException(e); + } + } + + public static NarClassLoader getFromArchive(File narPath, Set additionalJars, ClassLoader parent) throws IOException { + File unpacked = NarUnpacker.unpackNar(narPath, NAR_CACHE_DIR); + try { + return new NarClassLoader(unpacked, additionalJars, parent); } catch (ClassNotFoundException e) { throw new IOException(e); } @@ -155,6 +164,7 @@ public static NarClassLoader getFromArchive(File narPath, Set additional * * @param narWorkingDirectory * directory to explode nar contents to + * @param parent * @throws IllegalArgumentException * if the NAR is missing the Java Services API file for FlowFileProcessor implementations. * @throws ClassNotFoundException @@ -163,9 +173,9 @@ public static NarClassLoader getFromArchive(File narPath, Set additional * @throws IOException * if an error occurs while loading the NAR. */ - private NarClassLoader(final File narWorkingDirectory, Set additionalJars) + private NarClassLoader(final File narWorkingDirectory, Set additionalJars, ClassLoader parent) throws ClassNotFoundException, IOException { - super(new URL[0]); + super(new URL[0], parent); this.narWorkingDirectory = narWorkingDirectory; // process the classpath diff --git a/pulsar-sql/presto-pulsar-plugin/pom.xml b/pulsar-sql/presto-pulsar-plugin/pom.xml index 8f65436a901cf..ebd214f8212f1 100644 --- a/pulsar-sql/presto-pulsar-plugin/pom.xml +++ b/pulsar-sql/presto-pulsar-plugin/pom.xml @@ -70,5 +70,4 @@ - \ No newline at end of file diff --git a/pulsar-sql/presto-pulsar/pom.xml b/pulsar-sql/presto-pulsar/pom.xml index 795d689cf8d60..30bb88ba9ca74 100644 --- a/pulsar-sql/presto-pulsar/pom.xml +++ b/pulsar-sql/presto-pulsar/pom.xml @@ -34,11 +34,6 @@ 0.170 - 0.35 - 4.2.0 - 1.1.0.Final - 1 - 24.1-jre 2.1.2 1.8.4 @@ -56,24 +51,6 @@ ${dep.airlift.version} - - com.google.inject - guice - ${dep.guice.version} - - - - javax.validation - validation-api - ${dep.javax-validation.version} - - - - javax.inject - javax.inject - ${dep.javax-inject.version} - - org.apache.avro avro @@ -82,13 +59,13 @@ org.apache.pulsar - pulsar-client-admin + pulsar-client-admin-original ${project.version} org.apache.pulsar - managed-ledger + managed-ledger-original ${project.version} @@ -127,4 +104,81 @@ + + + + org.apache.maven.plugins + maven-shade-plugin + + + package + + shade + + + true + true + + + + org.apache.pulsar:pulsar-client-original + org.apache.pulsar:pulsar-client-admin-original + org.apache.pulsar:managed-ledger-original + + org.glassfish.jersey*:* + javax.ws.rs:* + javax.annotation:* + org.glassfish.hk2*:* + + org.apache.httpcomponents:* + org.eclipse.jetty:* + + + + + + org.apache.pulsar:pulsar-client-original + + ** + + + + + + org.glassfish + org.apache.pulsar.shade.org.glassfish + + + javax.ws + org.apache.pulsar.shade.javax.ws + + + javax.annotation + org.apache.pulsar.shade.javax.annotation + + + jersey + org.apache.pulsar.shade.jersey + + + org.eclipse.jetty + org.apache.pulsar.shade.org.eclipse.jetty + + + org.apache.http + org.apache.pulsar.shade.org.apache.http + + + + + + + + + + + + + + \ No newline at end of file diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/AvroSchemaHandler.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/AvroSchemaHandler.java index 41c2f6f1fa4a2..5d55682acf105 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/AvroSchemaHandler.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/AvroSchemaHandler.java @@ -20,10 +20,10 @@ import io.airlift.log.Logger; -import org.apache.pulsar.shade.io.netty.buffer.ByteBuf; -import org.apache.pulsar.shade.io.netty.buffer.ByteBufAllocator; -import org.apache.pulsar.shade.io.netty.util.ReferenceCountUtil; -import org.apache.pulsar.shade.io.netty.util.concurrent.FastThreadLocal; +import io.netty.buffer.ByteBuf; +import io.netty.buffer.ByteBufAllocator; +import io.netty.util.ReferenceCountUtil; +import io.netty.util.concurrent.FastThreadLocal; import org.apache.avro.Schema; import org.apache.avro.generic.GenericDatumReader; import org.apache.avro.generic.GenericRecord; diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/JSONSchemaHandler.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/JSONSchemaHandler.java index ae1a7c4115f4d..5a12d3011f47f 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/JSONSchemaHandler.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/JSONSchemaHandler.java @@ -28,8 +28,8 @@ import java.util.List; import java.util.Map; -import org.apache.pulsar.shade.io.netty.buffer.ByteBuf; -import org.apache.pulsar.shade.io.netty.util.concurrent.FastThreadLocal; +import io.netty.buffer.ByteBuf; +import io.netty.util.concurrent.FastThreadLocal; public class JSONSchemaHandler implements SchemaHandler { diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java index ee775bb2fc2a1..d6d648fa0f99f 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java @@ -18,19 +18,45 @@ */ package org.apache.pulsar.sql.presto; +import com.google.common.collect.ImmutableMap; +import io.airlift.log.Logger; +import org.apache.bookkeeper.mledger.LedgerOffloader; +import org.apache.bookkeeper.mledger.LedgerOffloaderFactory; +import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.ManagedLedgerFactory; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; -import org.apache.pulsar.shade.org.apache.bookkeeper.conf.ClientConfiguration; -import org.apache.pulsar.shade.org.apache.bookkeeper.stats.StatsProvider; +import org.apache.bookkeeper.mledger.impl.NullLedgerOffloader; +import org.apache.bookkeeper.mledger.offload.OffloaderUtils; +import org.apache.bookkeeper.mledger.offload.Offloaders; +import org.apache.commons.lang3.StringUtils; +import org.apache.pulsar.PulsarVersion; +import org.apache.bookkeeper.conf.ClientConfiguration; +import org.apache.bookkeeper.stats.StatsProvider; +import org.apache.bookkeeper.common.util.OrderedScheduler; + +import java.io.IOException; +import java.util.Arrays; +import java.util.HashMap; +import java.util.Map; + +import static com.google.common.base.Preconditions.checkNotNull; public class PulsarConnectorCache { + private static final Logger log = Logger.get(PulsarConnectorCache.class); + private static PulsarConnectorCache instance; private final ManagedLedgerFactory managedLedgerFactory; private final StatsProvider statsProvider; + private OrderedScheduler offloaderScheduler; + private Offloaders offloaderManager; + private LedgerOffloader offloader; + + static final String METADATA_SOFTWARE_VERSION_KEY = "S3ManagedLedgerOffloaderSoftwareVersion"; + static final String METADATA_SOFTWARE_GITSHA_KEY = "S3ManagedLedgerOffloaderSoftwareGitSha"; private PulsarConnectorCache(PulsarConnectorConfig pulsarConnectorConfig) throws Exception { this.managedLedgerFactory = initManagedLedgerFactory(pulsarConnectorConfig); @@ -43,6 +69,8 @@ private PulsarConnectorCache(PulsarConnectorConfig pulsarConnectorConfig) throws pulsarConnectorConfig.getStatsProviderConfigs().forEach((key, value) -> clientConfiguration.setProperty(key, value)); this.statsProvider.start(clientConfiguration); + + this.offloader = initManagedLedgerOffloader(pulsarConnectorConfig); } public static PulsarConnectorCache getConnectorCache(PulsarConnectorConfig pulsarConnectorConfig) throws Exception { @@ -58,7 +86,7 @@ private static ManagedLedgerFactory initManagedLedgerFactory(PulsarConnectorConf ClientConfiguration bkClientConfiguration = new ClientConfiguration() .setZkServers(pulsarConnectorConfig.getZookeeperUri()) .setAllowShadedLedgerManagerFactoryClass(true) - .setShadedLedgerManagerFactoryClassPrefix("org.apache.pulsar.shade.") + .setShadedLedgerManagerFactoryClassPrefix("") .setClientTcpNoDelay(false) .setUseV2WireProtocol(true) .setStickyReadsEnabled(true) @@ -66,6 +94,65 @@ private static ManagedLedgerFactory initManagedLedgerFactory(PulsarConnectorConf return new ManagedLedgerFactoryImpl(bkClientConfiguration); } + public ManagedLedgerConfig getManagedLedgerConfig() { + + return new ManagedLedgerConfig() + .setLedgerOffloader(this.offloader); + } + + private synchronized OrderedScheduler getOffloaderScheduler(PulsarConnectorConfig pulsarConnectorConfig) { + if (this.offloaderScheduler == null) { + this.offloaderScheduler = OrderedScheduler.newSchedulerBuilder() + .numThreads(pulsarConnectorConfig.getManagedLedgerOffloadMaxThreads()) + .name("offloader").build(); + } + return this.offloaderScheduler; + } + + private LedgerOffloader initManagedLedgerOffloader(PulsarConnectorConfig conf) { + + log.info("driver: %s - %s", conf.getManagedLedgerOffloadDriver(), StringUtils.isNotBlank(conf.getManagedLedgerOffloadDriver())); + try { + if (StringUtils.isNotBlank(conf.getManagedLedgerOffloadDriver())) { + checkNotNull(conf.getOffloadersDirectory(), + "Offloader driver is configured to be '%s' but no offloaders directory is configured.", + conf.getManagedLedgerOffloadDriver()); + this.offloaderManager = OffloaderUtils.searchForOffloaders(conf.getOffloadersDirectory()); + LedgerOffloaderFactory offloaderFactory = this.offloaderManager.getOffloaderFactory( + conf.getManagedLedgerOffloadDriver()); + + log.info("offloaderFactory: %s", offloaderFactory.getClass().getName()); + + log.info("supported: %s", offloaderFactory.isDriverSupported("aws-s3")); + + log.info("methods: %s", Arrays.toString(offloaderFactory.getClass().getDeclaredMethods())); + + Map offloaderProperties = conf.getOffloaderProperties(); + offloaderProperties.put("offloadersDirectory", conf.getOffloadersDirectory()); + offloaderProperties.put("managedLedgerOffloadDriver", conf.getManagedLedgerOffloadDriver()); + offloaderProperties.put("managedLedgerOffloadMaxThreads", String.valueOf(conf.getManagedLedgerOffloadMaxThreads())); + + try { + return offloaderFactory.create( + PulsarConnectorUtils.getProperties(offloaderProperties), + ImmutableMap.of( + METADATA_SOFTWARE_VERSION_KEY.toLowerCase(), PulsarVersion.getVersion(), + METADATA_SOFTWARE_GITSHA_KEY.toLowerCase(), PulsarVersion.getGitSha() + ), + getOffloaderScheduler(conf)); + } catch (IOException ioe) { + log.error("Failed to create offloader: ", ioe); + throw new RuntimeException(ioe.getMessage(), ioe.getCause()); + } + } else { + log.info("No ledger offloader configured, using NULL instance"); + return NullLedgerOffloader.INSTANCE; + } + } catch (Throwable t) { + throw new RuntimeException(t); + } + } + public ManagedLedgerFactory getManagedLedgerFactory() { return managedLedgerFactory; } @@ -74,11 +161,13 @@ public StatsProvider getStatsProvider() { return statsProvider; } - public static void shutdown() throws ManagedLedgerException, InterruptedException { + public static void shutdown() throws Exception { synchronized (PulsarConnectorCache.class) { if (instance != null) { - instance.managedLedgerFactory.shutdown(); instance.statsProvider.stop(); + instance.managedLedgerFactory.shutdown(); + instance.offloaderScheduler.shutdown(); + instance.offloaderManager.close(); instance = null; } } diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorConfig.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorConfig.java index 34d332e54ddd4..23992334ca149 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorConfig.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorConfig.java @@ -22,7 +22,7 @@ import io.airlift.configuration.Config; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.api.PulsarClientException; -import org.apache.pulsar.shade.org.apache.bookkeeper.stats.NullStatsProvider; +import org.apache.bookkeeper.stats.NullStatsProvider; import javax.validation.constraints.NotNull; import java.io.IOException; @@ -39,6 +39,13 @@ public class PulsarConnectorConfig implements AutoCloseable { private int maxSplitEntryQueueSize = 1000; private String statsProvider = NullStatsProvider.class.getName(); private Map statsProviderConfigs = new HashMap<>(); + + /**** --- Ledger Offloading --- ****/ + private String managedLedgerOffloadDriver = null; + private int managedLedgerOffloadMaxThreads = 2; + private String offloadersDirectory = "./offloaders"; + private Map offloaderProperties = new HashMap<>(); + private PulsarAdmin pulsarAdmin; @NotNull @@ -129,6 +136,49 @@ public PulsarConnectorConfig setStatsProviderConfigs(String statsProviderConfigs return this; } + /**** --- Ledger Offloading --- ****/ + + public int getManagedLedgerOffloadMaxThreads() { + return this.managedLedgerOffloadMaxThreads; + } + + @Config("pulsar.managed-ledger-offload-max-threads") + public PulsarConnectorConfig setManagedLedgerOffloadMaxThreads(int managedLedgerOffloadMaxThreads) throws IOException { + this.managedLedgerOffloadMaxThreads = managedLedgerOffloadMaxThreads; + return this; + } + + public String getManagedLedgerOffloadDriver() { + return this.managedLedgerOffloadDriver; + } + + @Config("pulsar.managed-ledger-offload-driver") + public PulsarConnectorConfig setManagedLedgerOffloadDriver(String managedLedgerOffloadDriver) throws IOException { + this.managedLedgerOffloadDriver = managedLedgerOffloadDriver; + return this; + } + + public String getOffloadersDirectory() { + return this.offloadersDirectory; + } + + + @Config("pulsar.offloaders-directory") + public PulsarConnectorConfig setOffloadersDirectory(String offloadersDirectory) throws IOException { + this.offloadersDirectory = offloadersDirectory; + return this; + } + + public Map getOffloaderProperties() { + return this.offloaderProperties; + } + + @Config("pulsar.offloader-properties") + public PulsarConnectorConfig setOffloaderProperties(String offloaderProperties) throws IOException { + this.offloaderProperties = new ObjectMapper().readValue(offloaderProperties, Map.class); + return this; + } + @NotNull public PulsarAdmin getPulsarAdmin() throws PulsarClientException { if (this.pulsarAdmin == null) { diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorMetricsTracker.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorMetricsTracker.java index d62a78837e53b..34969d8035b71 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorMetricsTracker.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorMetricsTracker.java @@ -18,11 +18,11 @@ */ package org.apache.pulsar.sql.presto; -import org.apache.pulsar.shade.org.apache.bookkeeper.stats.Counter; -import org.apache.pulsar.shade.org.apache.bookkeeper.stats.OpStatsLogger; -import org.apache.pulsar.shade.org.apache.bookkeeper.stats.StatsProvider; -import org.apache.pulsar.shade.org.apache.bookkeeper.stats.NullStatsProvider; -import org.apache.pulsar.shade.org.apache.bookkeeper.stats.StatsLogger; +import org.apache.bookkeeper.stats.Counter; +import org.apache.bookkeeper.stats.OpStatsLogger; +import org.apache.bookkeeper.stats.StatsProvider; +import org.apache.bookkeeper.stats.NullStatsProvider; +import org.apache.bookkeeper.stats.StatsLogger; import java.util.concurrent.TimeUnit; diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorUtils.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorUtils.java index 520ee6804064c..ee256f6377332 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorUtils.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorUtils.java @@ -26,6 +26,8 @@ import java.lang.reflect.Constructor; import java.lang.reflect.InvocationTargetException; +import java.util.Map; +import java.util.Properties; public class PulsarConnectorUtils { @@ -73,4 +75,12 @@ public static T createInstance(String userClassName, throw new RuntimeException("User class constructor throws exception", e); } } + + public static Properties getProperties(Map configMap) { + Properties properties = new Properties(); + for (Map.Entry entry : configMap.entrySet()) { + properties.setProperty(entry.getKey(), entry.getValue()); + } + return properties; + } } diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java index 9427f20548c3d..6e8ff4d886318 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java @@ -110,19 +110,20 @@ public PulsarRecordCursor(List columnHandles, PulsarSplit pu throw new RuntimeException(e); } initialize(columnHandles, pulsarSplit, pulsarConnectorConfig, - pulsarConnectorCache.getManagedLedgerFactory(), + pulsarConnectorCache.getManagedLedgerFactory(), pulsarConnectorCache.getManagedLedgerConfig(), new PulsarConnectorMetricsTracker(pulsarConnectorCache.getStatsProvider())); } // Exposed for testing purposes PulsarRecordCursor(List columnHandles, PulsarSplit pulsarSplit, PulsarConnectorConfig - pulsarConnectorConfig, ManagedLedgerFactory managedLedgerFactory, PulsarConnectorMetricsTracker pulsarConnectorMetricsTracker) { + pulsarConnectorConfig, ManagedLedgerFactory managedLedgerFactory, ManagedLedgerConfig managedLedgerConfig, + PulsarConnectorMetricsTracker pulsarConnectorMetricsTracker) { this.splitSize = pulsarSplit.getSplitSize(); - initialize(columnHandles, pulsarSplit, pulsarConnectorConfig, managedLedgerFactory, pulsarConnectorMetricsTracker); + initialize(columnHandles, pulsarSplit, pulsarConnectorConfig, managedLedgerFactory, managedLedgerConfig, pulsarConnectorMetricsTracker); } private void initialize(List columnHandles, PulsarSplit pulsarSplit, PulsarConnectorConfig - pulsarConnectorConfig, ManagedLedgerFactory managedLedgerFactory, + pulsarConnectorConfig, ManagedLedgerFactory managedLedgerFactory, ManagedLedgerConfig managedLedgerConfig, PulsarConnectorMetricsTracker pulsarConnectorMetricsTracker) { this.columnHandles = columnHandles; this.pulsarSplit = pulsarSplit; @@ -143,7 +144,7 @@ private void initialize(List columnHandles, PulsarSplit puls try { this.cursor = getCursor(TopicName.get("persistent", NamespaceName.get(pulsarSplit.getSchemaName()), - pulsarSplit.getTableName()), pulsarSplit.getStartPosition(), managedLedgerFactory); + pulsarSplit.getTableName()), pulsarSplit.getStartPosition(), managedLedgerFactory, managedLedgerConfig); } catch (ManagedLedgerException | InterruptedException e) { log.error(e, "Failed to get read only cursor"); close(); @@ -168,11 +169,11 @@ private SchemaHandler getSchemaHandler(Schema schema, SchemaType schemaType, } private ReadOnlyCursor getCursor(TopicName topicName, Position startPosition, ManagedLedgerFactory - managedLedgerFactory) + managedLedgerFactory, ManagedLedgerConfig managedLedgerConfig) throws ManagedLedgerException, InterruptedException { ReadOnlyCursor cursor = managedLedgerFactory.openReadOnlyCursor(topicName.getPersistenceNamingEncoding(), - startPosition, new ManagedLedgerConfig()); + startPosition, managedLedgerConfig); return cursor; } @@ -507,8 +508,13 @@ public void close() { currentMessage.release(); } - messageQueue.drain(RawMessage::release); - entryQueue.drain(Entry::release); + if (messageQueue != null) { + messageQueue.drain(RawMessage::release); + } + + if (entryQueue != null) { + entryQueue.drain(Entry::release); + } if (deserializeEntries != null) { deserializeEntries.interrupt(); diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarSplitManager.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarSplitManager.java index a2962e52ce9ba..7ab2ab687d3dc 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarSplitManager.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarSplitManager.java @@ -47,8 +47,8 @@ import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.schema.SchemaInfo; -import org.apache.pulsar.shade.com.google.common.base.Predicate; -import org.apache.pulsar.shade.org.apache.bookkeeper.conf.ClientConfiguration; +import com.google.common.base.Predicate; +import org.apache.bookkeeper.conf.ClientConfiguration; import javax.inject.Inject; import java.util.Collection; @@ -60,7 +60,8 @@ import static java.util.Objects.requireNonNull; import static org.apache.bookkeeper.mledger.ManagedCursor.FindPositionConstraint.SearchAllAvailableEntries; -public class PulsarSplitManager implements ConnectorSplitManager { +public class +PulsarSplitManager implements ConnectorSplitManager { private final String connectorId; @@ -127,7 +128,7 @@ ManagedLedgerFactory getManagedLedgerFactory() throws Exception { ClientConfiguration bkClientConfiguration = new ClientConfiguration() .setZkServers(this.pulsarConnectorConfig.getZookeeperUri()) .setAllowShadedLedgerManagerFactoryClass(true) - .setShadedLedgerManagerFactoryClassPrefix("org.apache.pulsar.shade.") + .setShadedLedgerManagerFactoryClassPrefix("") .setClientTcpNoDelay(false) .setStickyReadsEnabled(true) .setUseV2WireProtocol(true); @@ -351,10 +352,10 @@ public static PredicatePushdownInfo getPredicatePushdownInfo(String connectorId, // Just use a close bound since presto can always filter out the extra entries even if // the bound // should be open or a mixture of open and closed - org.apache.pulsar.shade.com.google.common.collect.Range posRange - = org.apache.pulsar.shade.com.google.common.collect.Range.range(overallStartPos, - org.apache.pulsar.shade.com.google.common.collect.BoundType.CLOSED, - overallEndPos, org.apache.pulsar.shade.com.google.common.collect.BoundType.CLOSED); + com.google.common.collect.Range posRange + = com.google.common.collect.Range.range(overallStartPos, + com.google.common.collect.BoundType.CLOSED, + overallEndPos, com.google.common.collect.BoundType.CLOSED); long numOfEntries = readOnlyCursor.getNumberOfEntries(posRange) - 1; diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/SchemaHandler.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/SchemaHandler.java index 3cabc8a6d642e..1e02d2ba302d2 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/SchemaHandler.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/SchemaHandler.java @@ -18,7 +18,7 @@ */ package org.apache.pulsar.sql.presto; -import org.apache.pulsar.shade.io.netty.buffer.ByteBuf; +import io.netty.buffer.ByteBuf; public interface SchemaHandler { diff --git a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java index 7a62f675a91eb..ab591763a3972 100644 --- a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java +++ b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java @@ -29,6 +29,7 @@ import io.airlift.log.Logger; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; +import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ManagedLedgerFactory; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.ReadOnlyCursor; @@ -51,9 +52,9 @@ import org.apache.pulsar.common.partition.PartitionedTopicMetadata; import org.apache.pulsar.common.schema.SchemaInfo; import org.apache.pulsar.common.schema.SchemaType; -import org.apache.pulsar.shade.javax.ws.rs.ClientErrorException; -import org.apache.pulsar.shade.javax.ws.rs.core.Response; -import org.apache.pulsar.shade.org.apache.bookkeeper.stats.NullStatsProvider; +import javax.ws.rs.ClientErrorException; +import javax.ws.rs.core.Response; +import org.apache.bookkeeper.stats.NullStatsProvider; import org.mockito.Mockito; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; @@ -625,10 +626,10 @@ private static List getTopicEntries(String topicSchemaName) { Schema schema = topicsToSchemas.get(topicSchemaName).getType() == SchemaType.AVRO ? AvroSchema.of(SchemaDefinition.builder().withPojo(Foo.class).build()) : JSONSchema.of(SchemaDefinition.builder().withPojo(Foo.class).build()); - org.apache.pulsar.shade.io.netty.buffer.ByteBuf payload = org.apache.pulsar.shade.io.netty.buffer.Unpooled + io.netty.buffer.ByteBuf payload = io.netty.buffer.Unpooled .copiedBuffer(schema.encode(foo)); - org.apache.pulsar.shade.io.netty.buffer.ByteBuf byteBuf = serializeMetadataAndPayload( + io.netty.buffer.ByteBuf byteBuf = serializeMetadataAndPayload( Commands.ChecksumType.Crc32c, messageMetadata, payload); Entry entry = EntryImpl.create(0, i, byteBuf); @@ -875,10 +876,10 @@ public void run() { Schema schema = topicsToSchemas.get(schemaName).getType() == SchemaType.AVRO ? AvroSchema.of(Foo.class) : JSONSchema.of(Foo.class); - org.apache.pulsar.shade.io.netty.buffer.ByteBuf payload = org.apache.pulsar.shade.io.netty.buffer.Unpooled + io.netty.buffer.ByteBuf payload = io.netty.buffer.Unpooled .copiedBuffer(schema.encode(foo)); - org.apache.pulsar.shade.io.netty.buffer.ByteBuf byteBuf = serializeMetadataAndPayload( + io.netty.buffer.ByteBuf byteBuf = serializeMetadataAndPayload( Commands.ChecksumType.Crc32c, messageMetadata, payload); completedBytes += byteBuf.readableBytes(); @@ -907,8 +908,8 @@ public Boolean answer(InvocationOnMock invocationOnMock) throws Throwable { @Override public Position answer(InvocationOnMock invocationOnMock) throws Throwable { Object[] args = invocationOnMock.getArguments(); - org.apache.pulsar.shade.com.google.common.base.Predicate predicate - = (org.apache.pulsar.shade.com.google.common.base.Predicate) args[1]; + com.google.common.base.Predicate predicate + = (com.google.common.base.Predicate) args[1]; String schemaName = TopicName.get( TopicName.get( @@ -933,8 +934,8 @@ public Position answer(InvocationOnMock invocationOnMock) throws Throwable { @Override public Long answer(InvocationOnMock invocationOnMock) throws Throwable { Object[] args = invocationOnMock.getArguments(); - org.apache.pulsar.shade.com.google.common.collect.Range range - = (org.apache.pulsar.shade.com.google.common.collect.Range ) args[0]; + com.google.common.collect.Range range + = (com.google.common.collect.Range ) args[0]; return (range.upperEndpoint().getEntryId() + 1) - range.lowerEndpoint().getEntryId(); } @@ -948,8 +949,10 @@ public Long answer(InvocationOnMock invocationOnMock) throws Throwable { for (Map.Entry split : splits.entrySet()) { - PulsarRecordCursor pulsarRecordCursor = spy(new PulsarRecordCursor(fooColumnHandles, split.getValue(), - pulsarConnectorConfig, managedLedgerFactory, new PulsarConnectorMetricsTracker(new NullStatsProvider()))); + PulsarRecordCursor pulsarRecordCursor = spy(new PulsarRecordCursor( + fooColumnHandles, split.getValue(), + pulsarConnectorConfig, managedLedgerFactory, new ManagedLedgerConfig(), + new PulsarConnectorMetricsTracker(new NullStatsProvider()))); this.pulsarRecordCursors.put(split.getKey(), pulsarRecordCursor); } } diff --git a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarMetadata.java b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarMetadata.java index 0803cd67ea546..f03b34fe4b08e 100644 --- a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarMetadata.java +++ b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarMetadata.java @@ -32,8 +32,8 @@ import org.apache.pulsar.client.admin.PulsarAdminException; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.schema.SchemaInfo; -import org.apache.pulsar.shade.javax.ws.rs.ClientErrorException; -import org.apache.pulsar.shade.javax.ws.rs.core.Response; +import javax.ws.rs.ClientErrorException; +import javax.ws.rs.core.Response; import org.testng.Assert; import org.testng.annotations.Test; From 76abc0de8c7568c3777b132d9e0806e23c5fde01 Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Sat, 13 Apr 2019 23:22:24 -0700 Subject: [PATCH 2/7] cleaning up --- conf/presto/catalog/pulsar.properties | 18 ++++++++++--- conf/presto/config.properties | 2 +- .../bookkeeper/mledger/LedgerOffloader.java | 4 +++ .../mledger/offload/OffloaderUtils.java | 1 - .../apache/pulsar/broker/PulsarService.java | 9 ++----- .../sql/presto/PulsarConnectorCache.java | 27 +++++++------------ .../pulsar/sql/presto/PulsarSplitManager.java | 5 +--- 7 files changed, 33 insertions(+), 33 deletions(-) diff --git a/conf/presto/catalog/pulsar.properties b/conf/presto/catalog/pulsar.properties index 35d83500039f5..7f191e5e424b8 100644 --- a/conf/presto/catalog/pulsar.properties +++ b/conf/presto/catalog/pulsar.properties @@ -34,6 +34,18 @@ pulsar.max-split-entry-queue-size = 1000 ####### TIERED STORAGE OFFLOADER CONFIGS ####### -pulsar.managed-ledger-offload-driver = aws-s3 -pulsar.offloaders-directory = /Users/jerrypeng/workspace/incubator-pulsar/offloaders -pulsar.offloader-properties = {"s3ManagedLedgerOffloadBucket": "jerry-pulsar-test", "s3ManagedLedgerOffloadRegion": "us-west-2", "s3ManagedLedgerOffloadServiceEndpoint": "http://s3.amazonaws.com"} \ No newline at end of file + +## Driver to use to offload old data to long term storage +#pulsar.managed-ledger-offload-driver = aws-s3 + +## The directory to locate offloaders +#pulsar.offloaders-directory = /pulsar/offloaders + +## Maximum number of thread pool threads for ledger offloading +#pulsar.managed-ledger-offload-max-threads = 2 + +## Properties and configurations related to specific offloader implementation +#pulsar.offloader-properties = \ +# {"s3ManagedLedgerOffloadBucket": "offload-bucket", \ +# "s3ManagedLedgerOffloadRegion": "us-west-2", \ +# "s3ManagedLedgerOffloadServiceEndpoint": "http://s3.amazonaws.com"} \ No newline at end of file diff --git a/conf/presto/config.properties b/conf/presto/config.properties index 0d54d241dead6..9f17135523dfa 100644 --- a/conf/presto/config.properties +++ b/conf/presto/config.properties @@ -38,5 +38,5 @@ query.client.timeout=5m query.min-expire-age=30m presto.version=testversion -#distributed-joins-enabled=true +distributed-joins-enabled=true node-scheduler.include-coordinator=true diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/LedgerOffloader.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/LedgerOffloader.java index 8fc35cc0b6935..c85fe9fd79107 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/LedgerOffloader.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/LedgerOffloader.java @@ -33,6 +33,10 @@ @Beta public interface LedgerOffloader { + // TODO: improve the user metadata in subsequent changes + String METADATA_SOFTWARE_VERSION_KEY = "S3ManagedLedgerOffloaderSoftwareVersion"; + String METADATA_SOFTWARE_GITSHA_KEY = "S3ManagedLedgerOffloaderSoftwareGitSha"; + /** * Get offload driver name. * diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java index 42ecb6966693f..e10f528f454b3 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java @@ -72,7 +72,6 @@ static Pair getOffloaderFactory(String n Thread.currentThread().setContextClassLoader(ncl); try { Object offloader = factoryClass.newInstance(); - if (!(offloader instanceof LedgerOffloaderFactory)) { throw new IOException("Class " + conf.getOffloaderFactoryClass() + " does not implement interface " + LedgerOffloaderFactory.class.getName()); 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 3cf442a268ac5..52c43c40da5cc 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 @@ -692,11 +692,6 @@ public LedgerOffloader getManagedLedgerOffloader() { return offloader; } - // TODO: improve the user metadata in subsequent changes - static final String METADATA_SOFTWARE_VERSION_KEY = "S3ManagedLedgerOffloaderSoftwareVersion"; - static final String METADATA_SOFTWARE_GITSHA_KEY = "S3ManagedLedgerOffloaderSoftwareGitSha"; - - public synchronized LedgerOffloader createManagedLedgerOffloader(ServiceConfiguration conf) throws PulsarServerException { try { @@ -711,8 +706,8 @@ public synchronized LedgerOffloader createManagedLedgerOffloader(ServiceConfigur return offloaderFactory.create( conf.getProperties(), ImmutableMap.of( - METADATA_SOFTWARE_VERSION_KEY.toLowerCase(), PulsarVersion.getVersion(), - METADATA_SOFTWARE_GITSHA_KEY.toLowerCase(), PulsarVersion.getGitSha() + LedgerOffloader.METADATA_SOFTWARE_VERSION_KEY.toLowerCase(), PulsarVersion.getVersion(), + LedgerOffloader.METADATA_SOFTWARE_GITSHA_KEY.toLowerCase(), PulsarVersion.getGitSha() ), getOffloaderScheduler(conf)); } catch (IOException ioe) { diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java index d6d648fa0f99f..df9d5dd4759ba 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java @@ -55,8 +55,10 @@ public class PulsarConnectorCache { private Offloaders offloaderManager; private LedgerOffloader offloader; - static final String METADATA_SOFTWARE_VERSION_KEY = "S3ManagedLedgerOffloaderSoftwareVersion"; - static final String METADATA_SOFTWARE_GITSHA_KEY = "S3ManagedLedgerOffloaderSoftwareGitSha"; + private static final String OFFLOADERS_DIRECTOR = "offloadersDirectory"; + private static final String MANAGED_LEDGER_OFFLOAD_DRIVER = "managedLedgerOffloadDriver"; + private static final String MANAGED_LEDGER_OFFLOAD_MAX_THREADS = "managedLedgerOffloadMaxThreads"; + private PulsarConnectorCache(PulsarConnectorConfig pulsarConnectorConfig) throws Exception { this.managedLedgerFactory = initManagedLedgerFactory(pulsarConnectorConfig); @@ -85,8 +87,6 @@ public static PulsarConnectorCache getConnectorCache(PulsarConnectorConfig pulsa private static ManagedLedgerFactory initManagedLedgerFactory(PulsarConnectorConfig pulsarConnectorConfig) throws Exception { ClientConfiguration bkClientConfiguration = new ClientConfiguration() .setZkServers(pulsarConnectorConfig.getZookeeperUri()) - .setAllowShadedLedgerManagerFactoryClass(true) - .setShadedLedgerManagerFactoryClassPrefix("") .setClientTcpNoDelay(false) .setUseV2WireProtocol(true) .setStickyReadsEnabled(true) @@ -104,14 +104,13 @@ private synchronized OrderedScheduler getOffloaderScheduler(PulsarConnectorConfi if (this.offloaderScheduler == null) { this.offloaderScheduler = OrderedScheduler.newSchedulerBuilder() .numThreads(pulsarConnectorConfig.getManagedLedgerOffloadMaxThreads()) - .name("offloader").build(); + .name("pulsar-offloader").build(); } return this.offloaderScheduler; } private LedgerOffloader initManagedLedgerOffloader(PulsarConnectorConfig conf) { - log.info("driver: %s - %s", conf.getManagedLedgerOffloadDriver(), StringUtils.isNotBlank(conf.getManagedLedgerOffloadDriver())); try { if (StringUtils.isNotBlank(conf.getManagedLedgerOffloadDriver())) { checkNotNull(conf.getOffloadersDirectory(), @@ -121,23 +120,17 @@ private LedgerOffloader initManagedLedgerOffloader(PulsarConnectorConfig conf) { LedgerOffloaderFactory offloaderFactory = this.offloaderManager.getOffloaderFactory( conf.getManagedLedgerOffloadDriver()); - log.info("offloaderFactory: %s", offloaderFactory.getClass().getName()); - - log.info("supported: %s", offloaderFactory.isDriverSupported("aws-s3")); - - log.info("methods: %s", Arrays.toString(offloaderFactory.getClass().getDeclaredMethods())); - Map offloaderProperties = conf.getOffloaderProperties(); - offloaderProperties.put("offloadersDirectory", conf.getOffloadersDirectory()); - offloaderProperties.put("managedLedgerOffloadDriver", conf.getManagedLedgerOffloadDriver()); - offloaderProperties.put("managedLedgerOffloadMaxThreads", String.valueOf(conf.getManagedLedgerOffloadMaxThreads())); + offloaderProperties.put(OFFLOADERS_DIRECTOR, conf.getOffloadersDirectory()); + offloaderProperties.put(MANAGED_LEDGER_OFFLOAD_DRIVER, conf.getManagedLedgerOffloadDriver()); + offloaderProperties.put(MANAGED_LEDGER_OFFLOAD_MAX_THREADS, String.valueOf(conf.getManagedLedgerOffloadMaxThreads())); try { return offloaderFactory.create( PulsarConnectorUtils.getProperties(offloaderProperties), ImmutableMap.of( - METADATA_SOFTWARE_VERSION_KEY.toLowerCase(), PulsarVersion.getVersion(), - METADATA_SOFTWARE_GITSHA_KEY.toLowerCase(), PulsarVersion.getGitSha() + LedgerOffloader.METADATA_SOFTWARE_VERSION_KEY.toLowerCase(), PulsarVersion.getVersion(), + LedgerOffloader.METADATA_SOFTWARE_GITSHA_KEY.toLowerCase(), PulsarVersion.getGitSha() ), getOffloaderScheduler(conf)); } catch (IOException ioe) { diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarSplitManager.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarSplitManager.java index 7ab2ab687d3dc..aa11a73c012d6 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarSplitManager.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarSplitManager.java @@ -60,8 +60,7 @@ import static java.util.Objects.requireNonNull; import static org.apache.bookkeeper.mledger.ManagedCursor.FindPositionConstraint.SearchAllAvailableEntries; -public class -PulsarSplitManager implements ConnectorSplitManager { +public class PulsarSplitManager implements ConnectorSplitManager { private final String connectorId; @@ -127,8 +126,6 @@ public ConnectorSplitSource getSplits(ConnectorTransactionHandle transactionHand ManagedLedgerFactory getManagedLedgerFactory() throws Exception { ClientConfiguration bkClientConfiguration = new ClientConfiguration() .setZkServers(this.pulsarConnectorConfig.getZookeeperUri()) - .setAllowShadedLedgerManagerFactoryClass(true) - .setShadedLedgerManagerFactoryClassPrefix("") .setClientTcpNoDelay(false) .setStickyReadsEnabled(true) .setUseV2WireProtocol(true); From 12f4959503f78ee7961ec1a68d7a5cb861bbe953 Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Sat, 13 Apr 2019 23:24:31 -0700 Subject: [PATCH 3/7] cleaning up imports --- .../bookkeeper/mledger/offload/OffloaderUtils.java | 1 - pulsar-sql/presto-pulsar-plugin/pom.xml | 1 + .../apache/pulsar/sql/presto/PulsarConnectorCache.java | 9 +++------ 3 files changed, 4 insertions(+), 7 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java index e10f528f454b3..845b53f61ad07 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/offload/OffloaderUtils.java @@ -20,7 +20,6 @@ import java.io.File; import java.io.IOException; -import java.net.URLClassLoader; import java.nio.file.DirectoryStream; import java.nio.file.Files; import java.nio.file.Path; diff --git a/pulsar-sql/presto-pulsar-plugin/pom.xml b/pulsar-sql/presto-pulsar-plugin/pom.xml index ebd214f8212f1..8f65436a901cf 100644 --- a/pulsar-sql/presto-pulsar-plugin/pom.xml +++ b/pulsar-sql/presto-pulsar-plugin/pom.xml @@ -70,4 +70,5 @@ + \ No newline at end of file diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java index df9d5dd4759ba..a9cd070d3f13e 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarConnectorCache.java @@ -20,24 +20,21 @@ import com.google.common.collect.ImmutableMap; import io.airlift.log.Logger; +import org.apache.bookkeeper.common.util.OrderedScheduler; +import org.apache.bookkeeper.conf.ClientConfiguration; import org.apache.bookkeeper.mledger.LedgerOffloader; import org.apache.bookkeeper.mledger.LedgerOffloaderFactory; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; -import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.ManagedLedgerFactory; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; import org.apache.bookkeeper.mledger.impl.NullLedgerOffloader; import org.apache.bookkeeper.mledger.offload.OffloaderUtils; import org.apache.bookkeeper.mledger.offload.Offloaders; +import org.apache.bookkeeper.stats.StatsProvider; import org.apache.commons.lang3.StringUtils; import org.apache.pulsar.PulsarVersion; -import org.apache.bookkeeper.conf.ClientConfiguration; -import org.apache.bookkeeper.stats.StatsProvider; -import org.apache.bookkeeper.common.util.OrderedScheduler; import java.io.IOException; -import java.util.Arrays; -import java.util.HashMap; import java.util.Map; import static com.google.common.base.Preconditions.checkNotNull; From 4106d0647f47e79c8a5c16e1c08fa752fac3c8b1 Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Sat, 13 Apr 2019 23:31:11 -0700 Subject: [PATCH 4/7] cleaning up configs --- conf/standalone.conf | 6 ------ 1 file changed, 6 deletions(-) diff --git a/conf/standalone.conf b/conf/standalone.conf index 5542eb1704ac7..4cbf931c75748 100644 --- a/conf/standalone.conf +++ b/conf/standalone.conf @@ -17,12 +17,6 @@ # under the License. # -managedLedgerOffloadDriver=aws-s3 -s3ManagedLedgerOffloadBucket=jerry-pulsar-test -s3ManagedLedgerOffloadRegion=us-west-2 -//s3ManagedLedgerOffloadServiceEndpoint=https://apigateway.us-west-2.amazonaws.com -s3ManagedLedgerOffloadServiceEndpoint=http://s3.amazonaws.com - ### --- General broker settings --- ### # Zookeeper quorum connection string From 2715bfa1cd727dd8b71b8340a59b330738be32e2 Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Sun, 14 Apr 2019 12:58:29 -0700 Subject: [PATCH 5/7] fix imports --- .../apache/pulsar/sql/presto/TestPulsarConnector.java | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java index ab591763a3972..faa8958ed3567 100644 --- a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java +++ b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java @@ -27,6 +27,7 @@ import com.facebook.presto.spi.type.Type; import com.facebook.presto.spi.type.VarcharType; import io.airlift.log.Logger; +import io.netty.buffer.ByteBuf; import org.apache.bookkeeper.mledger.AsyncCallbacks; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; @@ -626,10 +627,10 @@ private static List getTopicEntries(String topicSchemaName) { Schema schema = topicsToSchemas.get(topicSchemaName).getType() == SchemaType.AVRO ? AvroSchema.of(SchemaDefinition.builder().withPojo(Foo.class).build()) : JSONSchema.of(SchemaDefinition.builder().withPojo(Foo.class).build()); - io.netty.buffer.ByteBuf payload = io.netty.buffer.Unpooled + ByteBuf payload = io.netty.buffer.Unpooled .copiedBuffer(schema.encode(foo)); - io.netty.buffer.ByteBuf byteBuf = serializeMetadataAndPayload( + ByteBuf byteBuf = serializeMetadataAndPayload( Commands.ChecksumType.Crc32c, messageMetadata, payload); Entry entry = EntryImpl.create(0, i, byteBuf); @@ -876,10 +877,10 @@ public void run() { Schema schema = topicsToSchemas.get(schemaName).getType() == SchemaType.AVRO ? AvroSchema.of(Foo.class) : JSONSchema.of(Foo.class); - io.netty.buffer.ByteBuf payload = io.netty.buffer.Unpooled + ByteBuf payload = io.netty.buffer.Unpooled .copiedBuffer(schema.encode(foo)); - io.netty.buffer.ByteBuf byteBuf = serializeMetadataAndPayload( + ByteBuf byteBuf = serializeMetadataAndPayload( Commands.ChecksumType.Crc32c, messageMetadata, payload); completedBytes += byteBuf.readableBytes(); From b64f0f30741244c1b867b0bc8f14bd58f96cd01b Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Mon, 15 Apr 2019 11:12:37 -0700 Subject: [PATCH 6/7] fix behavior when offloader not configured and fix license --- .../mledger/impl/ReadOnlyCursorImpl.java | 5 ++ pulsar-sql/presto-distribution/LICENSE | 73 ++++++++++++++++++- .../pulsar/sql/presto/PulsarRecordCursor.java | 22 +++++- 3 files changed, 94 insertions(+), 6 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyCursorImpl.java index c1a2216b4b6a6..8c9da1d8c8577 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ReadOnlyCursorImpl.java @@ -27,6 +27,7 @@ import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.ReadOnlyCursor; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl.PositionBound; +import org.apache.bookkeeper.mledger.proto.MLDataFormats; @Slf4j public class ReadOnlyCursorImpl extends ManagedCursorImpl implements ReadOnlyCursor { @@ -62,6 +63,10 @@ public void asyncClose(final AsyncCallbacks.CloseCallback callback, final Object callback.closeComplete(ctx); } + public MLDataFormats.ManagedLedgerInfo.LedgerInfo getCurrentLedgerInfo() { + return this.ledger.getLedgersInfo().get(this.readPosition.getLedgerId()); + } + public long getNumberOfEntries(Range range) { return this.ledger.getNumberOfEntries(range); } diff --git a/pulsar-sql/presto-distribution/LICENSE b/pulsar-sql/presto-distribution/LICENSE index b2a983085eccf..d61b2fa98afc2 100644 --- a/pulsar-sql/presto-distribution/LICENSE +++ b/pulsar-sql/presto-distribution/LICENSE @@ -235,6 +235,23 @@ The Apache Software License, Version 2.0 - commons-lang3-3.4.jar * Netty - netty-3.6.2.Final.jar + - netty-all-4.1.32.Final.jar + - netty-buffer-4.1.31.Final.jar + - netty-codec-4.1.31.Final.jar + - netty-codec-dns-4.1.33.Final.jar + - netty-codec-http-4.1.33.Final.jar + - netty-codec-socks-4.1.33.Final.jar + - netty-common-4.1.31.Final.jar + - netty-handler-4.1.31.Final.jar + - netty-handler-proxy-4.1.33.Final.jar + - netty-reactive-streams-2.0.0.jar + - netty-resolver-4.1.31.Final.jar + - netty-resolver-dns-4.1.33.Final.jar + - netty-tcnative-boringssl-static-2.0.20.Final.jar + - netty-transport-4.1.31.Final.jar + - netty-transport-native-epoll-4.1.31.Final.jar + - netty-transport-native-epoll-4.1.33.Final-linux-x86_64.jar + - netty-transport-native-unix-common-4.1.31.Final.jar * Joda Time - joda-time-2.9.9.jar * Jetty @@ -252,8 +269,6 @@ The Apache Software License, Version 2.0 - jetty-server-9.4.11.v20180605.jar - jetty-servlet-9.4.11.v20180605.jar - jetty-util-9.4.11.v20180605.jar - * Javassist - - javassist-3.22.0-CR2.jar * Asynchronous Http Client - async-http-client-1.6.5.jar * Apache BVal @@ -372,7 +387,6 @@ The Apache Software License, Version 2.0 - rocksdbjni-5.13.3.jar * SnakeYAML - snakeyaml-1.17.jar - - snakeyaml-1.23.jar * Snappy Java - snappy-java-1.1.1.3.jar * Bean Validation API @@ -392,10 +406,59 @@ The Apache Software License, Version 2.0 - lz4-java-1.5.0.jar * JCTools - jctools-core-2.1.2.jar + * Asynchronous Http Client + - async-http-client-2.7.0.jar + - async-http-client-netty-utils-2.7.0.jar + * Apache Bookkeeper + - bookkeeper-common-4.9.0.jar + - bookkeeper-common-allocator-4.9.0.jar + - bookkeeper-proto-4.9.0.jar + - bookkeeper-server-4.9.0.jar + - bookkeeper-stats-api-4.9.0.jar + - bookkeeper-tools-framework-4.9.0.jar + - circe-checksum-4.9.0.jar + - codahale-metrics-provider-4.9.0.jar + - cpu-affinity-4.9.0.jar + - http-server-4.9.0.jar + - prometheus-metrics-provider-4.9.0.jar + * Apache Commons + - commons-cli-1.2.jar + - commons-codec-1.10.jar + - commons-collections4-4.1.jar + - commons-configuration-1.10.jar + - commons-io-2.5.jar + - commons-lang-2.6.jar + - commons-logging-1.1.1.jar + * GSON + - gson-2.8.2.jar + * Jackson + - jackson-jaxrs-base-2.8.11.jar + - jackson-jaxrs-json-provider-2.8.11.jar + - jackson-module-jaxb-annotations-2.8.11.jar + - jackson-module-jsonSchema-2.8.11.jar + * Java Assist + - javassist-3.21.0-GA.jar + * Jetty + - jetty-http-9.4.12.v20180830.jar + - jetty-io-9.4.12.v20180830.jar + - jetty-security-9.4.12.v20180830.jar + - jetty-server-9.4.12.v20180830.jar + - jetty-servlet-9.4.12.v20180830.jar + - jetty-util-9.4.12.v20180830.jar + * Java Native Access + - jna-4.2.0.jar + * Yahoo Datasketches + - memory-0.8.3.jar + - sketches-core-0.8.3.jar + * Apache Zookeeper + - zookeeper-3.4.13.jar + * Apache Yetus Audience Annotations + - audience-annotations-0.5.0.jar Protocol Buffers License * Protocol Buffers - protobuf-shaded-2.1.0-incubating.jar + - protobuf-java-3.5.1.jar BSD 3-clause "New" or "Revised" License * RE2J TD -- re2j-td-1.4.jar @@ -476,6 +539,8 @@ CDDL-1.1 -- licenses/LICENSE-CDDL-1.1.txt - jgrapht-core-0.9.0.jar * Logback Core Module - logback-core-1.2.3.jar + * MIME Streaming Extension + - mimepull-1.9.6.jar Public Domain (CC0) -- licenses/LICENSE-CC0.txt * HdrHistogram @@ -484,6 +549,8 @@ Public Domain (CC0) -- licenses/LICENSE-CC0.txt - aopalliance-1.0.jar * XZ For Java - xz-1.5.jar + * Reactive Streams + - reactive-streams-1.0.2.jar Bouncy Castle License * Bouncy Castle -- licenses/LICENSE-bouncycastle.txt diff --git a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java index 6e8ff4d886318..4ffbd2f2e80a9 100644 --- a/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java +++ b/pulsar-sql/presto-pulsar/src/main/java/org/apache/pulsar/sql/presto/PulsarRecordCursor.java @@ -56,6 +56,7 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.ReadOnlyCursor; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.bookkeeper.mledger.impl.ReadOnlyCursorImpl; import org.apache.pulsar.common.api.raw.MessageParser; import org.apache.pulsar.common.api.raw.RawMessage; import org.apache.pulsar.common.naming.NamespaceName; @@ -83,6 +84,7 @@ public class PulsarRecordCursor implements RecordCursor { private DeserializeEntries deserializeEntries; private TopicName topicName; private PulsarConnectorMetricsTracker metricsTracker; + private boolean readOffloaded; // Stats total execution time of split private long startTime; @@ -135,6 +137,7 @@ private void initialize(List columnHandles, PulsarSplit puls NamespaceName.get(pulsarSplit.getSchemaName()), pulsarSplit.getTableName()); this.metricsTracker = pulsarConnectorMetricsTracker; + this.readOffloaded = pulsarConnectorConfig.getManagedLedgerOffloadDriver() != null; Schema schema = PulsarConnectorUtils.parseSchema(pulsarSplit.getSchema()); @@ -299,7 +302,6 @@ class ReadEntries implements AsyncCallbacks.ReadEntriesCallback { public void run() { if (outstandingReadsRequests.get() > 0) { - if (!cursor.hasMoreEntries() || ((PositionImpl) cursor.getReadPosition()) .compareTo(pulsarSplit.getEndPosition()) >= 0) { isDone = true; @@ -308,8 +310,22 @@ public void run() { int batchSize = Math.min(maxBatchSize, entryQueue.capacity() - entryQueue.size()); if (batchSize > 0) { - outstandingReadsRequests.decrementAndGet(); - cursor.asyncReadEntries(batchSize, this, System.nanoTime()); + + ReadOnlyCursorImpl readOnlyCursorImpl = ((ReadOnlyCursorImpl) cursor); + // check if ledger is offloaded + if (!readOffloaded && readOnlyCursorImpl.getCurrentLedgerInfo().hasOffloadContext()) { + log.warn("Ledger %s is offloaded for topic %s. Ignoring it because offloader is not configured", + readOnlyCursorImpl.getCurrentLedgerInfo().getLedgerId(), pulsarSplit.getTableName()); + + long numEntries = readOnlyCursorImpl.getCurrentLedgerInfo().getEntries(); + long entriesToSkip = (numEntries - ((PositionImpl) cursor.getReadPosition()).getEntryId()) + 1; + cursor.skipEntries(Math.toIntExact((entriesToSkip))); + + entriesProcessed += entriesToSkip; + } else { + outstandingReadsRequests.decrementAndGet(); + cursor.asyncReadEntries(batchSize, this, System.nanoTime()); + } // stats for successful read request metricsTracker.incr_READ_ATTEMPTS_SUCCESS(); From ef15c1d14a8daec75f618fb0a1fcb5a87b0584eb Mon Sep 17 00:00:00 2001 From: Jerry Peng Date: Mon, 15 Apr 2019 13:02:31 -0700 Subject: [PATCH 7/7] fix unit test --- .../org/apache/pulsar/sql/presto/TestPulsarConnector.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java index faa8958ed3567..cedde2078303a 100644 --- a/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java +++ b/pulsar-sql/presto-pulsar/src/test/java/org/apache/pulsar/sql/presto/TestPulsarConnector.java @@ -36,6 +36,8 @@ import org.apache.bookkeeper.mledger.ReadOnlyCursor; import org.apache.bookkeeper.mledger.impl.EntryImpl; import org.apache.bookkeeper.mledger.impl.PositionImpl; +import org.apache.bookkeeper.mledger.impl.ReadOnlyCursorImpl; +import org.apache.bookkeeper.mledger.proto.MLDataFormats; import org.apache.pulsar.client.admin.Namespaces; import org.apache.pulsar.client.admin.PulsarAdmin; import org.apache.pulsar.client.admin.PulsarAdminException; @@ -796,7 +798,7 @@ public ReadOnlyCursor answer(InvocationOnMock invocationOnMock) throws Throwable long entries = topicsToNumEntries.get(schemaName); - ReadOnlyCursor readOnlyCursor = mock(ReadOnlyCursor.class); + ReadOnlyCursorImpl readOnlyCursor = mock(ReadOnlyCursorImpl.class); doReturn(entries).when(readOnlyCursor).getNumberOfEntries(); doAnswer(new Answer() { @@ -942,6 +944,8 @@ public Long answer(InvocationOnMock invocationOnMock) throws Throwable { } }); + when(readOnlyCursor.getCurrentLedgerInfo()).thenReturn(MLDataFormats.ManagedLedgerInfo.LedgerInfo.newBuilder().setLedgerId(0).build()); + return readOnlyCursor; } });