KYLIN-2565, upgrade to Hadoop3.0
Project: http://git-wip-us.apache.org/repos/asf/kylin/repo Commit: http://git-wip-us.apache.org/repos/asf/kylin/commit/dcb7cf47 Tree: http://git-wip-us.apache.org/repos/asf/kylin/tree/dcb7cf47 Diff: http://git-wip-us.apache.org/repos/asf/kylin/diff/dcb7cf47 Branch: refs/heads/master-hadoop3.0 Commit: dcb7cf47ce97d7e65a1327cf803a520e976f95ad Parents: a0b537b Author: Cheng Wang <[email protected]> Authored: Tue Apr 25 18:45:57 2017 +0800 Committer: liyang-gmt8 <[email protected]> Committed: Wed Apr 26 15:45:56 2017 +0800 ---------------------------------------------------------------------- .../common/DefaultSslProtocolSocketFactory.java | 150 ------------------- .../engine/mr/common/HadoopStatusGetter.java | 70 ++++++--- .../kylin/engine/spark/SparkCountDemo.java | 4 +- .../apache/kylin/engine/spark/SparkCubing.java | 22 +-- .../kylin/engine/spark/SparkCubingByLayer.java | 25 ++-- pom.xml | 21 ++- server-base/pom.xml | 5 + .../apache/kylin/rest/security/MockHTable.java | 116 ++++++++++---- .../apache/kylin/rest/service/AclService.java | 2 +- .../apache/kylin/rest/service/AdminService.java | 12 +- storage-hbase/pom.xml | 4 + .../kylin/storage/hbase/HBaseConnection.java | 5 + .../hbase/cube/v2/CubeHBaseEndpointRPC.java | 4 +- .../storage/hbase/cube/v2/CubeHBaseScanRPC.java | 2 +- .../coprocessor/endpoint/CubeVisitService.java | 4 +- .../kylin/storage/hbase/steps/CubeHFileJob.java | 18 ++- .../storage/hbase/steps/HBaseCuboidWriter.java | 2 +- .../storage/hbase/util/CubeMigrationCLI.java | 2 +- .../hbase/util/DeployCoprocessorCLI.java | 5 +- .../hbase/util/ExtendCubeToHybridCLI.java | 2 +- .../hbase/util/GridTableHBaseBenchmark.java | 2 +- .../kylin/storage/hbase/util/PingHBaseCLI.java | 3 +- .../hbase/steps/CubeHFileMapperTest.java | 22 ++- .../storage/hbase/steps/TestHbaseClient.java | 14 +- .../org/apache/kylin/tool/CubeMigrationCLI.java | 25 ++-- .../kylin/tool/CubeMigrationCheckCLI.java | 17 ++- .../kylin/tool/ExtendCubeToHybridCLI.java | 2 +- .../apache/kylin/tool/StorageCleanupJob.java | 18 ++- .../org/apache/kylin/tool/util/ToolUtil.java | 9 +- 29 files changed, 293 insertions(+), 294 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/engine-mr/src/main/java/org/apache/kylin/engine/mr/common/DefaultSslProtocolSocketFactory.java ---------------------------------------------------------------------- diff --git a/engine-mr/src/main/java/org/apache/kylin/engine/mr/common/DefaultSslProtocolSocketFactory.java b/engine-mr/src/main/java/org/apache/kylin/engine/mr/common/DefaultSslProtocolSocketFactory.java deleted file mode 100644 index d66e4eb..0000000 --- a/engine-mr/src/main/java/org/apache/kylin/engine/mr/common/DefaultSslProtocolSocketFactory.java +++ /dev/null @@ -1,150 +0,0 @@ -/* - * 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.kylin.engine.mr.common; - -import java.io.IOException; -import java.net.InetAddress; -import java.net.Socket; -import java.net.UnknownHostException; - -import javax.net.ssl.SSLContext; -import javax.net.ssl.TrustManager; - -import org.apache.commons.httpclient.ConnectTimeoutException; -import org.apache.commons.httpclient.HttpClientError; -import org.apache.commons.httpclient.params.HttpConnectionParams; -import org.apache.commons.httpclient.protocol.ControllerThreadSocketFactory; -import org.apache.commons.httpclient.protocol.SecureProtocolSocketFactory; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * @author xduo - * - */ -public class DefaultSslProtocolSocketFactory implements SecureProtocolSocketFactory { - /** Log object for this class. */ - private static Logger logger = LoggerFactory.getLogger(DefaultSslProtocolSocketFactory.class); - private SSLContext sslcontext = null; - - /** - * Constructor for DefaultSslProtocolSocketFactory. - */ - public DefaultSslProtocolSocketFactory() { - super(); - } - - /** - * @see SecureProtocolSocketFactory#createSocket(java.lang.String,int,java.net.InetAddress,int) - */ - public Socket createSocket(String host, int port, InetAddress clientHost, int clientPort) throws IOException, UnknownHostException { - return getSSLContext().getSocketFactory().createSocket(host, port, clientHost, clientPort); - } - - /** - * Attempts to get a new socket connection to the given host within the - * given time limit. - * - * <p> - * To circumvent the limitations of older JREs that do not support connect - * timeout a controller thread is executed. The controller thread attempts - * to create a new socket within the given limit of time. If socket - * constructor does not return until the timeout expires, the controller - * terminates and throws an {@link ConnectTimeoutException} - * </p> - * - * @param host - * the host name/IP - * @param port - * the port on the host - * @param localAddress - * the local host name/IP to bind the socket to - * @param localPort - * the port on the local machine - * @param params - * {@link HttpConnectionParams Http connection parameters} - * - * @return Socket a new socket - * - * @throws IOException - * if an I/O error occurs while creating the socket - * @throws UnknownHostException - * if the IP address of the host cannot be determined - * @throws ConnectTimeoutException - * DOCUMENT ME! - * @throws IllegalArgumentException - * DOCUMENT ME! - */ - public Socket createSocket(final String host, final int port, final InetAddress localAddress, final int localPort, final HttpConnectionParams params) throws IOException, UnknownHostException, ConnectTimeoutException { - if (params == null) { - throw new IllegalArgumentException("Parameters may not be null"); - } - - int timeout = params.getConnectionTimeout(); - - if (timeout == 0) { - return createSocket(host, port, localAddress, localPort); - } else { - // To be eventually deprecated when migrated to Java 1.4 or above - return ControllerThreadSocketFactory.createSocket(this, host, port, localAddress, localPort, timeout); - } - } - - /** - * @see SecureProtocolSocketFactory#createSocket(java.lang.String,int) - */ - public Socket createSocket(String host, int port) throws IOException, UnknownHostException { - return getSSLContext().getSocketFactory().createSocket(host, port); - } - - /** - * @see SecureProtocolSocketFactory#createSocket(java.net.Socket,java.lang.String,int,boolean) - */ - public Socket createSocket(Socket socket, String host, int port, boolean autoClose) throws IOException, UnknownHostException { - return getSSLContext().getSocketFactory().createSocket(socket, host, port, autoClose); - } - - public boolean equals(Object obj) { - return ((obj != null) && obj.getClass().equals(DefaultX509TrustManager.class)); - } - - public int hashCode() { - return DefaultX509TrustManager.class.hashCode(); - } - - private static SSLContext createEasySSLContext() { - try { - SSLContext context = SSLContext.getInstance("TLS"); - context.init(null, new TrustManager[] { new DefaultX509TrustManager(null) }, null); - - return context; - } catch (Exception e) { - logger.error(e.getMessage(), e); - throw new HttpClientError(e.toString()); - } - } - - private SSLContext getSSLContext() { - if (this.sslcontext == null) { - this.sslcontext = createEasySSLContext(); - } - - return this.sslcontext; - } -} http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/engine-mr/src/main/java/org/apache/kylin/engine/mr/common/HadoopStatusGetter.java ---------------------------------------------------------------------- diff --git a/engine-mr/src/main/java/org/apache/kylin/engine/mr/common/HadoopStatusGetter.java b/engine-mr/src/main/java/org/apache/kylin/engine/mr/common/HadoopStatusGetter.java index f31369b..0245c1c 100644 --- a/engine-mr/src/main/java/org/apache/kylin/engine/mr/common/HadoopStatusGetter.java +++ b/engine-mr/src/main/java/org/apache/kylin/engine/mr/common/HadoopStatusGetter.java @@ -21,25 +21,31 @@ package org.apache.kylin.engine.mr.common; import java.io.IOException; import java.net.MalformedURLException; import java.nio.charset.Charset; +import java.security.KeyManagementException; import java.security.Principal; +import java.security.SecureRandom; +import java.security.cert.X509Certificate; + +import javax.net.ssl.SSLContext; +import javax.net.ssl.TrustManager; -import org.apache.commons.httpclient.Header; -import org.apache.commons.httpclient.HttpClient; -import org.apache.commons.httpclient.HttpMethod; -import org.apache.commons.httpclient.methods.GetMethod; -import org.apache.commons.httpclient.protocol.Protocol; -import org.apache.commons.httpclient.protocol.ProtocolSocketFactory; import org.apache.commons.io.IOUtils; import org.apache.commons.lang.StringUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.hadoop.yarn.api.records.FinalApplicationStatus; import org.apache.hadoop.yarn.server.resourcemanager.rmapp.RMAppState; +import org.apache.http.Header; import org.apache.http.HttpResponse; import org.apache.http.auth.AuthSchemeRegistry; import org.apache.http.auth.AuthScope; import org.apache.http.auth.Credentials; +import org.apache.http.client.HttpClient; import org.apache.http.client.methods.HttpGet; import org.apache.http.client.params.AuthPolicy; +import org.apache.http.conn.ClientConnectionManager; +import org.apache.http.conn.scheme.Scheme; +import org.apache.http.conn.scheme.SchemeRegistry; +import org.apache.http.conn.ssl.SSLSocketFactory; import org.apache.http.impl.auth.SPNegoSchemeFactory; import org.apache.http.impl.client.BasicCredentialsProvider; import org.apache.http.impl.client.DefaultHttpClient; @@ -107,7 +113,7 @@ public class HadoopStatusGetter { String response = null; while (response == null) { if (url.startsWith("https://")) { - registerEasyHttps(); + registerEasyHttps(client); } if (url.contains("anonymous=true") == false) { url += url.contains("?") ? "&" : "?"; @@ -162,26 +168,26 @@ public class HadoopStatusGetter { } private String getHttpResponse(String url) throws IOException { - HttpClient client = new HttpClient(); + HttpClient client = new DefaultHttpClient(); String response = null; while (response == null) { // follow redirects via 'refresh' if (url.startsWith("https://")) { - registerEasyHttps(); + registerEasyHttps(client); } if (url.contains("anonymous=true") == false) { url += url.contains("?") ? "&" : "?"; url += "anonymous=true"; } - HttpMethod get = new GetMethod(url); - get.addRequestHeader("accept", "application/json"); + HttpGet get = new HttpGet(url); + get.addHeader("accept", "application/json"); try { - client.executeMethod(get); + HttpResponse res = client.execute(get); String redirect = null; - Header h = get.getResponseHeader("Location"); + Header h = res.getFirstHeader("Location"); if (h != null) { redirect = h.getValue(); if (isValidURL(redirect) == false) { @@ -190,7 +196,7 @@ public class HadoopStatusGetter { continue; } } else { - h = get.getResponseHeader("Refresh"); + h = res.getFirstHeader("Refresh"); if (h != null) { String s = h.getValue(); int cut = s.indexOf("url="); @@ -207,7 +213,7 @@ public class HadoopStatusGetter { } if (redirect == null) { - response = get.getResponseBodyAsString(); + response = res.getStatusLine().toString(); logger.debug("Job " + mrJobId + " get status check result.\n"); } else { url = redirect; @@ -224,13 +230,35 @@ public class HadoopStatusGetter { return response; } - private static Protocol EASY_HTTPS = null; + private static void registerEasyHttps(HttpClient client) { + SSLContext sslContext; + try { + sslContext = SSLContext.getInstance("SSL"); + + // set up a TrustManager that trusts everything + try { + sslContext.init(null, new TrustManager[] { new DefaultX509TrustManager(null) { + public X509Certificate[] getAcceptedIssuers() { + logger.debug("getAcceptedIssuers"); + return null; + } + + public void checkClientTrusted(X509Certificate[] certs, String authType) { + logger.debug("checkClientTrusted"); + } - private static void registerEasyHttps() { - // by pass all https issue - if (EASY_HTTPS == null) { - EASY_HTTPS = new Protocol("https", (ProtocolSocketFactory) new DefaultSslProtocolSocketFactory(), 443); - Protocol.registerProtocol("https", EASY_HTTPS); + public void checkServerTrusted(X509Certificate[] certs, String authType) { + logger.debug("checkServerTrusted"); + } + } }, new SecureRandom()); + } catch (KeyManagementException e) { + } + SSLSocketFactory ssf = new SSLSocketFactory(sslContext, SSLSocketFactory.ALLOW_ALL_HOSTNAME_VERIFIER); + ClientConnectionManager ccm = client.getConnectionManager(); + SchemeRegistry sr = ccm.getSchemeRegistry(); + sr.register(new Scheme("https", 443, ssf)); + } catch (Exception e) { + logger.error(e.getMessage(), e); } } http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCountDemo.java ---------------------------------------------------------------------- diff --git a/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCountDemo.java b/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCountDemo.java index 6478c10..a079a57 100644 --- a/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCountDemo.java +++ b/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCountDemo.java @@ -22,7 +22,7 @@ import org.apache.commons.cli.OptionBuilder; import org.apache.commons.cli.Options; import org.apache.hadoop.hbase.KeyValue; import org.apache.hadoop.hbase.io.ImmutableBytesWritable; -import org.apache.hadoop.hbase.mapreduce.HFileOutputFormat; +import org.apache.hadoop.hbase.mapreduce.HFileOutputFormat2; import org.apache.kylin.common.util.AbstractApplication; import org.apache.kylin.common.util.OptionsHelper; import org.apache.spark.SparkConf; @@ -74,7 +74,7 @@ public class SparkCountDemo extends AbstractApplication { KeyValue value = new KeyValue(stringIntegerTuple2._1().getBytes(), "f".getBytes(), "c".getBytes(), String.valueOf(stringIntegerTuple2._2()).getBytes()); return new Tuple2(key, value); } - }).saveAsNewAPIHadoopFile("hdfs://10.249.65.231:8020/tmp/hfile", ImmutableBytesWritable.class, KeyValue.class, HFileOutputFormat.class); + }).saveAsNewAPIHadoopFile("hdfs://10.249.65.231:8020/tmp/hfile", ImmutableBytesWritable.class, KeyValue.class, HFileOutputFormat2.class); } } http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCubing.java ---------------------------------------------------------------------- diff --git a/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCubing.java b/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCubing.java index 2a0981a..a87d66b 100644 --- a/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCubing.java +++ b/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCubing.java @@ -41,9 +41,11 @@ import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.FsShell; import org.apache.hadoop.fs.Path; import org.apache.hadoop.hbase.KeyValue; -import org.apache.hadoop.hbase.client.HTable; +import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.client.Connection; +import org.apache.hadoop.hbase.client.Table; import org.apache.hadoop.hbase.io.ImmutableBytesWritable; -import org.apache.hadoop.hbase.mapreduce.HFileOutputFormat; +import org.apache.hadoop.hbase.mapreduce.HFileOutputFormat2; import org.apache.hadoop.hbase.mapreduce.LoadIncrementalHFiles; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.util.ToolRunner; @@ -367,7 +369,7 @@ public class SparkCubing extends AbstractApplication { final JavaPairRDD<byte[], byte[]> javaPairRDD = javaRDD.glom().mapPartitionsToPair(new PairFlatMapFunction<Iterator<List<List<String>>>, byte[], byte[]>() { @Override - public Iterable<Tuple2<byte[], byte[]>> call(Iterator<List<List<String>>> listIterator) throws Exception { + public Iterator<Tuple2<byte[], byte[]>> call(Iterator<List<List<String>>> listIterator) throws Exception { long t = System.currentTimeMillis(); prepare(); @@ -390,7 +392,7 @@ public class SparkCubing extends AbstractApplication { throw new RuntimeException(e); } System.out.println("build partition cost: " + (System.currentTimeMillis() - t) + "ms"); - return sparkCuboidWriter.getResult(); + return sparkCuboidWriter.getResult().iterator(); } }); @@ -430,8 +432,8 @@ public class SparkCubing extends AbstractApplication { } }, UnsignedBytes.lexicographicalComparator()).mapPartitions(new FlatMapFunction<Iterator<Tuple2<byte[], byte[]>>, Tuple2<byte[], byte[]>>() { @Override - public Iterable<Tuple2<byte[], byte[]>> call(final Iterator<Tuple2<byte[], byte[]>> tuple2Iterator) throws Exception { - return new Iterable<Tuple2<byte[], byte[]>>() { + public Iterator<Tuple2<byte[], byte[]>> call(final Iterator<Tuple2<byte[], byte[]>> tuple2Iterator) throws Exception { + Iterable<Tuple2<byte[], byte[]>> iterable = new Iterable<Tuple2<byte[], byte[]>>() { final BufferedMeasureCodec codec = new BufferedMeasureCodec(dataTypes); final Object[] input = new Object[measureSize]; final Object[] result = new Object[measureSize]; @@ -459,6 +461,7 @@ public class SparkCubing extends AbstractApplication { }); } }; + return iterable.iterator(); } }, true).mapToPair(new PairFunction<Tuple2<byte[], byte[]>, ImmutableBytesWritable, KeyValue>() { @Override @@ -467,7 +470,7 @@ public class SparkCubing extends AbstractApplication { KeyValue value = new KeyValue(tuple2._1(), "F1".getBytes(), "M".getBytes(), tuple2._2()); return new Tuple2(key, value); } - }).saveAsNewAPIHadoopFile(hFileLocation, ImmutableBytesWritable.class, KeyValue.class, HFileOutputFormat.class, conf); + }).saveAsNewAPIHadoopFile(hFileLocation, ImmutableBytesWritable.class, KeyValue.class, HFileOutputFormat2.class, conf); } public static void prepare() throws Exception { @@ -496,8 +499,9 @@ public class SparkCubing extends AbstractApplication { Job job = Job.getInstance(conf); job.setMapOutputKeyClass(ImmutableBytesWritable.class); job.setMapOutputValueClass(KeyValue.class); - HTable table = new HTable(conf, hTableName); - HFileOutputFormat.configureIncrementalLoad(job, table); + Connection connection = HBaseConnection.get(); + Table table = connection.getTable(TableName.valueOf(hTableName)); + HFileOutputFormat2.configureIncrementalLoad(job, table, connection.getRegionLocator(TableName.valueOf(hTableName))); return conf; } http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCubingByLayer.java ---------------------------------------------------------------------- diff --git a/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCubingByLayer.java b/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCubingByLayer.java index f70fd30..28591ef 100644 --- a/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCubingByLayer.java +++ b/engine-spark/src/main/java/org/apache/kylin/engine/spark/SparkCubingByLayer.java @@ -17,6 +17,15 @@ */ package org.apache.kylin.engine.spark; +import java.io.File; +import java.io.FileFilter; +import java.io.Serializable; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Iterator; +import java.util.List; + import org.apache.commons.cli.Option; import org.apache.commons.cli.OptionBuilder; import org.apache.commons.cli.Options; @@ -65,16 +74,8 @@ import org.apache.spark.sql.hive.HiveContext; import org.apache.spark.storage.StorageLevel; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import scala.Tuple2; - -import java.io.File; -import java.io.FileFilter; -import java.io.Serializable; -import java.nio.ByteBuffer; -import java.util.ArrayList; -import java.util.Collection; -import java.util.List; +import scala.Tuple2; /** * Spark application to build cube with the "by-layer" algorithm. Only support source data from Hive; Metadata in HBase. @@ -354,7 +355,7 @@ public class SparkCubingByLayer extends AbstractApplication implements Serializa } @Override - public Iterable<Tuple2<ByteArray, Object[]>> call(Tuple2<ByteArray, Object[]> tuple2) throws Exception { + public Iterator<Tuple2<ByteArray, Object[]>> call(Tuple2<ByteArray, Object[]> tuple2) throws Exception { if (initialized == false) { prepare(); initialized = true; @@ -368,7 +369,7 @@ public class SparkCubingByLayer extends AbstractApplication implements Serializa // if still empty or null if (myChildren == null || myChildren.size() == 0) { - return EMTPY_ITERATOR; + return EMTPY_ITERATOR.iterator(); } List<Tuple2<ByteArray, Object[]>> tuples = new ArrayList(myChildren.size()); @@ -382,7 +383,7 @@ public class SparkCubingByLayer extends AbstractApplication implements Serializa tuples.add(new Tuple2<>(new ByteArray(newKey), tuple2._2())); } - return tuples; + return tuples.iterator(); } } http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/pom.xml ---------------------------------------------------------------------- diff --git a/pom.xml b/pom.xml index 0e772ca..ea35550 100644 --- a/pom.xml +++ b/pom.xml @@ -41,29 +41,30 @@ <properties> <!-- General Properties --> - <javaVersion>1.7</javaVersion> + <javaVersion>1.8</javaVersion> <maven-model.version>3.3.9</maven-model.version> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding> <!-- Hadoop versions --> - <hadoop2.version>2.7.1</hadoop2.version> - <yarn.version>2.7.1</yarn.version> + <hadoop2.version>3.0.0-alpha2</hadoop2.version> + <yarn.version>3.0.0-alpha2</yarn.version> <!-- Hive versions --> - <hive.version>1.2.1</hive.version> - <hive-hcatalog.version>1.2.1</hive-hcatalog.version> + <hive.version>2.1.0</hive.version> + <hive-hcatalog.version>2.1.0</hive-hcatalog.version> <!-- HBase versions --> - <hbase-hadoop2.version>1.1.1</hbase-hadoop2.version> + <hbase-hadoop2.version>2.0.0-SNAPSHOT</hbase-hadoop2.version> <!-- Kafka versions --> <kafka.version>0.10.1.0</kafka.version> <!-- Spark versions --> - <spark.version>1.6.3</spark.version> + <spark.version>2.0.0-SNAPSHOT</spark.version> <kryo.version>4.0.0</kryo.version> + <commons-configuration.version>1.6</commons-configuration.version> <!-- <reflections.version>0.9.10</reflections.version> --> <!-- Calcite Version --> @@ -518,6 +519,12 @@ <version>${yarn.version}</version> </dependency> + <dependency> + <groupId>commons-configuration</groupId> + <artifactId>commons-configuration</artifactId> + <version>${commons-configuration.version}</version> + </dependency> + <!-- Calcite dependencies --> <dependency> <groupId>org.apache.calcite</groupId> http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/server-base/pom.xml ---------------------------------------------------------------------- diff --git a/server-base/pom.xml b/server-base/pom.xml index 627ae35..34d8ffb 100644 --- a/server-base/pom.xml +++ b/server-base/pom.xml @@ -169,6 +169,11 @@ <artifactId>junit</artifactId> <scope>test</scope> </dependency> + <dependency> + <groupId>commons-configuration</groupId> + <artifactId>commons-configuration</artifactId> + <scope>provided</scope> + </dependency> </dependencies> <repositories> http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/server-base/src/main/java/org/apache/kylin/rest/security/MockHTable.java ---------------------------------------------------------------------- diff --git a/server-base/src/main/java/org/apache/kylin/rest/security/MockHTable.java b/server-base/src/main/java/org/apache/kylin/rest/security/MockHTable.java index 972eea9..fd53b5b 100644 --- a/server-base/src/main/java/org/apache/kylin/rest/security/MockHTable.java +++ b/server-base/src/main/java/org/apache/kylin/rest/security/MockHTable.java @@ -43,6 +43,8 @@ import java.util.TreeMap; import org.apache.commons.lang.NotImplementedException; import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.hbase.Cell; +import org.apache.hadoop.hbase.CellUtil; import org.apache.hadoop.hbase.HColumnDescriptor; import org.apache.hadoop.hbase.HTableDescriptor; import org.apache.hadoop.hbase.KeyValue; @@ -51,7 +53,6 @@ import org.apache.hadoop.hbase.client.Append; import org.apache.hadoop.hbase.client.Delete; import org.apache.hadoop.hbase.client.Durability; import org.apache.hadoop.hbase.client.Get; -import org.apache.hadoop.hbase.client.Table; import org.apache.hadoop.hbase.client.Increment; import org.apache.hadoop.hbase.client.Mutation; import org.apache.hadoop.hbase.client.Put; @@ -60,7 +61,9 @@ import org.apache.hadoop.hbase.client.ResultScanner; import org.apache.hadoop.hbase.client.Row; import org.apache.hadoop.hbase.client.RowMutations; import org.apache.hadoop.hbase.client.Scan; +import org.apache.hadoop.hbase.client.Table; import org.apache.hadoop.hbase.client.coprocessor.Batch; +import org.apache.hadoop.hbase.client.metrics.ScanMetrics; import org.apache.hadoop.hbase.filter.CompareFilter; import org.apache.hadoop.hbase.filter.Filter; import org.apache.hadoop.hbase.ipc.CoprocessorRpcChannel; @@ -97,7 +100,7 @@ public class MockHTable implements Table { private NavigableMap<byte[], NavigableMap<byte[], NavigableMap<byte[], NavigableMap<Long, byte[]>>>> data = new TreeMap<>(Bytes.BYTES_COMPARATOR); - private static List<KeyValue> toKeyValue(byte[] row, NavigableMap<byte[], NavigableMap<byte[], NavigableMap<Long, byte[]>>> rowdata, int maxVersions) { + private static List<Cell> toKeyValue(byte[] row, NavigableMap<byte[], NavigableMap<byte[], NavigableMap<Long, byte[]>>> rowdata, int maxVersions) { return toKeyValue(row, rowdata, 0, Long.MAX_VALUE, maxVersions); } @@ -162,8 +165,8 @@ public class MockHTable implements Table { throw new RuntimeException(this.getClass() + " does NOT implement this method."); } - private static List<KeyValue> toKeyValue(byte[] row, NavigableMap<byte[], NavigableMap<byte[], NavigableMap<Long, byte[]>>> rowdata, long timestampStart, long timestampEnd, int maxVersions) { - List<KeyValue> ret = new ArrayList<KeyValue>(); + private static List<Cell> toKeyValue(byte[] row, NavigableMap<byte[], NavigableMap<byte[], NavigableMap<Long, byte[]>>> rowdata, long timestampStart, long timestampEnd, int maxVersions) { + List<Cell> ret = new ArrayList<>(); for (byte[] family : rowdata.keySet()) for (byte[] qualifier : rowdata.get(family).keySet()) { int versionsAdded = 0; @@ -207,7 +210,6 @@ public class MockHTable implements Table { /** * {@inheritDoc} */ - @Override public Object[] batch(List<? extends Row> actions) throws IOException, InterruptedException { Object[] results = new Object[actions.size()]; // same size. for (int i = 0; i < actions.size(); i++) { @@ -241,11 +243,6 @@ public class MockHTable implements Table { } - @Override - public <R> Object[] batchCallback(List<? extends Row> actions, Batch.Callback<R> callback) throws IOException, InterruptedException { - return new Object[0]; - } - /** * {@inheritDoc} */ @@ -254,7 +251,7 @@ public class MockHTable implements Table { if (!data.containsKey(get.getRow())) return new Result(); byte[] row = get.getRow(); - List<KeyValue> kvs = new ArrayList<KeyValue>(); + List<Cell> kvs = new ArrayList<>(); if (!get.hasFamilies()) { kvs = toKeyValue(row, data.get(row), get.getMaxVersions()); } else { @@ -279,7 +276,7 @@ public class MockHTable implements Table { kvs = filter(filter, kvs); } - return new Result(kvs); + return Result.create(kvs); } /** @@ -317,11 +314,11 @@ public class MockHTable implements Table { break; } - List<KeyValue> kvs = null; + List<Cell> kvs = null; if (!scan.hasFamilies()) { kvs = toKeyValue(row, data.get(row), scan.getTimeRange().getMin(), scan.getTimeRange().getMax(), scan.getMaxVersions()); } else { - kvs = new ArrayList<KeyValue>(); + kvs = new ArrayList<>(); for (byte[] family : scan.getFamilyMap().keySet()) { if (data.get(row).get(family) == null) continue; @@ -353,7 +350,7 @@ public class MockHTable implements Table { } } if (!kvs.isEmpty()) { - ret.add(new Result(kvs)); + ret.add(Result.create(kvs)); } } @@ -387,6 +384,16 @@ public class MockHTable implements Table { public void close() { } + + @Override + public boolean renewLease() { + return false; + } + + @Override + public ScanMetrics getScanMetrics() { + return null; + } }; } @@ -397,10 +404,10 @@ public class MockHTable implements Table { * @param kvs List of a row's KeyValues * @return List of KeyValues that were not filtered. */ - private List<KeyValue> filter(Filter filter, List<KeyValue> kvs) throws IOException { + private List<Cell> filter(Filter filter, List<Cell> kvs) throws IOException { filter.reset(); - List<KeyValue> tmp = new ArrayList<KeyValue>(kvs.size()); + List<Cell> tmp = new ArrayList<>(kvs.size()); tmp.addAll(kvs); /* @@ -409,9 +416,9 @@ public class MockHTable implements Table { * See Figure 4-2 on p. 163. */ boolean filteredOnRowKey = false; - List<KeyValue> nkvs = new ArrayList<KeyValue>(tmp.size()); - for (KeyValue kv : tmp) { - if (filter.filterRowKey(kv.getBuffer(), kv.getRowOffset(), kv.getRowLength())) { + List<Cell> nkvs = new ArrayList<>(tmp.size()); + for (Cell kv : tmp) { + if (filter.filterRowKey(kv)) { filteredOnRowKey = true; break; } @@ -474,16 +481,16 @@ public class MockHTable implements Table { public void put(Put put) throws IOException { byte[] row = put.getRow(); NavigableMap<byte[], NavigableMap<byte[], NavigableMap<Long, byte[]>>> rowData = forceFind(data, row, new TreeMap<byte[], NavigableMap<byte[], NavigableMap<Long, byte[]>>>(Bytes.BYTES_COMPARATOR)); - for (byte[] family : put.getFamilyMap().keySet()) { + for (byte[] family : put.getFamilyCellMap().keySet()) { if (columnFamilies.contains(new String(family)) == false) { throw new RuntimeException("Not Exists columnFamily : " + new String(family)); } NavigableMap<byte[], NavigableMap<Long, byte[]>> familyData = forceFind(rowData, family, new TreeMap<byte[], NavigableMap<Long, byte[]>>(Bytes.BYTES_COMPARATOR)); - for (KeyValue kv : put.getFamilyMap().get(family)) { - kv.updateLatestStamp(Bytes.toBytes(System.currentTimeMillis())); - byte[] qualifier = kv.getQualifier(); + for (Cell kv : put.getFamilyCellMap().get(family)) { + CellUtil.updateLatestStamp(kv, System.currentTimeMillis()); + byte[] qualifier = kv.getQualifierArray(); NavigableMap<Long, byte[]> qualifierData = forceFind(familyData, qualifier, new TreeMap<Long, byte[]>()); - qualifierData.put(kv.getTimestamp(), kv.getValue()); + qualifierData.put(kv.getTimestamp(), kv.getValueArray()); } } } @@ -531,22 +538,22 @@ public class MockHTable implements Table { byte[] row = delete.getRow(); if (data.get(row) == null) return; - if (delete.getFamilyMap().size() == 0) { + if (delete.getFamilyCellMap().size() == 0) { data.remove(row); return; } - for (byte[] family : delete.getFamilyMap().keySet()) { + for (byte[] family : delete.getFamilyCellMap().keySet()) { if (data.get(row).get(family) == null) continue; - if (delete.getFamilyMap().get(family).isEmpty()) { + if (delete.getFamilyCellMap().get(family).isEmpty()) { data.get(row).remove(family); continue; } - for (KeyValue kv : delete.getFamilyMap().get(family)) { - if (kv.isDelete()) { - data.get(row).get(kv.getFamily()).clear(); + for (Cell kv : delete.getFamilyCellMap().get(family)) { + if (CellUtil.isDelete(kv)) { + data.get(row).get(kv.getFamilyArray()).clear(); } else { - data.get(row).get(kv.getFamily()).remove(kv.getQualifier()); + data.get(row).get(kv.getFamilyArray()).remove(kv.getQualifierArray()); } } if (data.get(row).get(family).isEmpty()) { @@ -665,4 +672,49 @@ public class MockHTable implements Table { throw new NotImplementedException(); } + + /*** + * + * All values are default + * + * **/ + @Override + public void setOperationTimeout(int i) { + + } + + @Override + public int getOperationTimeout() { + return 0; + } + + @Override + public int getRpcTimeout() { + return 0; + } + + @Override + public void setRpcTimeout(int i) { + + } + + @Override + public int getReadRpcTimeout() { + return 0; + } + + @Override + public void setReadRpcTimeout(int i) { + + } + + @Override + public int getWriteRpcTimeout() { + return 0; + } + + @Override + public void setWriteRpcTimeout(int i) { + + } } \ No newline at end of file http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/server-base/src/main/java/org/apache/kylin/rest/service/AclService.java ---------------------------------------------------------------------- diff --git a/server-base/src/main/java/org/apache/kylin/rest/service/AclService.java b/server-base/src/main/java/org/apache/kylin/rest/service/AclService.java index b80d97d..184e4b0 100644 --- a/server-base/src/main/java/org/apache/kylin/rest/service/AclService.java +++ b/server-base/src/main/java/org/apache/kylin/rest/service/AclService.java @@ -290,7 +290,7 @@ public class AclService implements MutableAclService { htable = aclHBaseStorage.getTable(aclTableName); Delete delete = new Delete(Bytes.toBytes(String.valueOf(acl.getObjectIdentity().getIdentifier()))); - delete.deleteFamily(Bytes.toBytes(AclHBaseStorage.ACL_ACES_FAMILY)); + delete.addFamily(Bytes.toBytes(AclHBaseStorage.ACL_ACES_FAMILY)); htable.delete(delete); Put put = new Put(Bytes.toBytes(String.valueOf(acl.getObjectIdentity().getIdentifier()))); http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/server-base/src/main/java/org/apache/kylin/rest/service/AdminService.java ---------------------------------------------------------------------- diff --git a/server-base/src/main/java/org/apache/kylin/rest/service/AdminService.java b/server-base/src/main/java/org/apache/kylin/rest/service/AdminService.java index 66725dc..e06157d 100644 --- a/server-base/src/main/java/org/apache/kylin/rest/service/AdminService.java +++ b/server-base/src/main/java/org/apache/kylin/rest/service/AdminService.java @@ -18,12 +18,6 @@ package org.apache.kylin.rest.service; -import java.io.ByteArrayOutputStream; -import java.io.IOException; -import java.util.Map; -import java.util.Properties; -import java.util.TreeMap; - import org.apache.commons.configuration.ConfigurationException; import org.apache.commons.configuration.PropertiesConfiguration; import org.apache.kylin.common.KylinConfig; @@ -36,6 +30,12 @@ import org.slf4j.LoggerFactory; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.stereotype.Component; +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.util.Map; +import java.util.Properties; +import java.util.TreeMap; + /** * @author jianliu */ http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/pom.xml ---------------------------------------------------------------------- diff --git a/storage-hbase/pom.xml b/storage-hbase/pom.xml index 29ca7e5..118a946 100644 --- a/storage-hbase/pom.xml +++ b/storage-hbase/pom.xml @@ -96,6 +96,10 @@ <artifactId>junit</artifactId> <scope>test</scope> </dependency> + <dependency> + <groupId>com.google.code.findbugs</groupId> + <artifactId>jsr305</artifactId> + </dependency> </dependencies> <build> http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/HBaseConnection.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/HBaseConnection.java b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/HBaseConnection.java index 7e2cefc..714c669 100644 --- a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/HBaseConnection.java +++ b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/HBaseConnection.java @@ -210,6 +210,11 @@ public class HBaseConnection { // ============================================================================ + public static Connection get() { + String url = KylinConfig.getInstanceFromEnv().getStorageUrl(); + return get(url); + } + // returned Connection can be shared by multiple threads and does not require close() @SuppressWarnings("resource") public static Connection get(String url) { http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/CubeHBaseEndpointRPC.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/CubeHBaseEndpointRPC.java b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/CubeHBaseEndpointRPC.java index e822ada..aad5234 100644 --- a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/CubeHBaseEndpointRPC.java +++ b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/CubeHBaseEndpointRPC.java @@ -30,7 +30,7 @@ import org.apache.hadoop.hbase.TableName; import org.apache.hadoop.hbase.client.Connection; import org.apache.hadoop.hbase.client.Table; import org.apache.hadoop.hbase.client.coprocessor.Batch; -import org.apache.hadoop.hbase.ipc.BlockingRpcCallback; +import org.apache.hadoop.hbase.ipc.CoprocessorRpcUtils; import org.apache.hadoop.hbase.ipc.ServerRpcController; import org.apache.kylin.common.KylinConfig; import org.apache.kylin.common.exceptions.KylinTimeoutException; @@ -183,7 +183,7 @@ public class CubeHBaseEndpointRPC extends CubeHBaseRPC { new Batch.Call<CubeVisitService, CubeVisitResponse>() { public CubeVisitResponse call(CubeVisitService rowsService) throws IOException { ServerRpcController controller = new ServerRpcController(); - BlockingRpcCallback<CubeVisitResponse> rpcCallback = new BlockingRpcCallback<>(); + CoprocessorRpcUtils.BlockingRpcCallback<CubeVisitResponse> rpcCallback = new CoprocessorRpcUtils.BlockingRpcCallback<>(); rowsService.visitCube(controller, request, rpcCallback); CubeVisitResponse response = rpcCallback.get(); if (controller.failedOnException()) { http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/CubeHBaseScanRPC.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/CubeHBaseScanRPC.java b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/CubeHBaseScanRPC.java index 951e2ef..e116935 100644 --- a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/CubeHBaseScanRPC.java +++ b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/CubeHBaseScanRPC.java @@ -181,7 +181,7 @@ public class CubeHBaseScanRPC extends CubeHBaseRPC { public List<Cell> next() { List<Cell> result = allResultsIterator.next().listCells(); for (Cell cell : result) { - scannedBytes += CellUtil.estimatedSizeOf(cell); + scannedBytes += CellUtil.estimatedSerializedSizeOf(cell); } scannedRows++; return result; http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/coprocessor/endpoint/CubeVisitService.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/coprocessor/endpoint/CubeVisitService.java b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/coprocessor/endpoint/CubeVisitService.java index cde127e..6142025 100644 --- a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/coprocessor/endpoint/CubeVisitService.java +++ b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/cube/v2/coprocessor/endpoint/CubeVisitService.java @@ -37,9 +37,9 @@ import org.apache.hadoop.hbase.client.Scan; import org.apache.hadoop.hbase.coprocessor.CoprocessorException; import org.apache.hadoop.hbase.coprocessor.CoprocessorService; import org.apache.hadoop.hbase.coprocessor.RegionCoprocessorEnvironment; -import org.apache.hadoop.hbase.protobuf.ResponseConverter; import org.apache.hadoop.hbase.regionserver.HRegion; import org.apache.hadoop.hbase.regionserver.RegionScanner; +import org.apache.hadoop.hbase.shaded.protobuf.ResponseConverter; import org.apache.kylin.common.KylinConfig; import org.apache.kylin.common.exceptions.KylinTimeoutException; import org.apache.kylin.common.exceptions.ResourceLimitExceededException; @@ -177,7 +177,7 @@ public class CubeVisitService extends CubeVisitProtos.CubeVisitService implement List<Cell> result = delegate.next(); rowCount++; for (Cell cell : result) { - rowBytes += CellUtil.estimatedSizeOf(cell); + rowBytes += CellUtil.estimatedSerializedSizeOf(cell); } return result; } http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/steps/CubeHFileJob.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/steps/CubeHFileJob.java b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/steps/CubeHFileJob.java index 1a624c4..0b92c64 100644 --- a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/steps/CubeHFileJob.java +++ b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/steps/CubeHFileJob.java @@ -25,8 +25,12 @@ import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.hadoop.hbase.HBaseConfiguration; -import org.apache.hadoop.hbase.client.HTable; -import org.apache.hadoop.hbase.mapreduce.HFileOutputFormat; +import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.client.Connection; +import org.apache.hadoop.hbase.client.ConnectionFactory; +import org.apache.hadoop.hbase.client.RegionLocator; +import org.apache.hadoop.hbase.client.Table; +import org.apache.hadoop.hbase.mapreduce.HFileOutputFormat2; import org.apache.hadoop.hbase.mapreduce.KeyValueSortReducer; import org.apache.hadoop.hdfs.DFSConfigKeys; import org.apache.hadoop.io.SequenceFile; @@ -56,6 +60,7 @@ public class CubeHFileJob extends AbstractHadoopJob { public int run(String[] args) throws Exception { Options options = new Options(); + Connection connection = null; try { options.addOption(OPTION_JOB_NAME); options.addOption(OPTION_CUBE_NAME); @@ -92,10 +97,13 @@ public class CubeHFileJob extends AbstractHadoopJob { attachCubeMetadata(cube, job.getConfiguration()); Configuration hbaseConf = HBaseConfiguration.create(getConf()); - HTable htable = new HTable(hbaseConf, getOptionValue(OPTION_HTABLE_NAME).toUpperCase()); + String hTableName = getOptionValue(OPTION_HTABLE_NAME).toUpperCase(); + connection = ConnectionFactory.createConnection(hbaseConf); + Table table = connection.getTable(TableName.valueOf(hTableName)); + RegionLocator regionLocator = connection.getRegionLocator(TableName.valueOf(hTableName)); // Automatic config ! - HFileOutputFormat.configureIncrementalLoad(job, htable); + HFileOutputFormat2.configureIncrementalLoad(job, table, regionLocator); reconfigurePartitions(hbaseConf, partitionFilePath); // set block replication to 3 for hfiles @@ -107,6 +115,8 @@ public class CubeHFileJob extends AbstractHadoopJob { } finally { if (job != null) cleanupTempConfFile(job.getConfiguration()); + if (null != connection) + connection.close(); } } http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/steps/HBaseCuboidWriter.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/steps/HBaseCuboidWriter.java b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/steps/HBaseCuboidWriter.java index 6587d4e..7bb5ecd 100644 --- a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/steps/HBaseCuboidWriter.java +++ b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/steps/HBaseCuboidWriter.java @@ -103,7 +103,7 @@ public class HBaseCuboidWriter implements ICuboidWriter { byte[] family = copy(keyValue.getFamilyArray(), keyValue.getFamilyOffset(), keyValue.getFamilyLength()); byte[] qualifier = copy(keyValue.getQualifierArray(), keyValue.getQualifierOffset(), keyValue.getQualifierLength()); byte[] value = copy(keyValue.getValueArray(), keyValue.getValueOffset(), keyValue.getValueLength()); - put.add(family, qualifier, value); + put.addColumn(family, qualifier, value); puts.add(put); } if (puts.size() >= BATCH_PUT_THRESHOLD) { http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/CubeMigrationCLI.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/CubeMigrationCLI.java b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/CubeMigrationCLI.java index 581de38..58fba8d 100644 --- a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/CubeMigrationCLI.java +++ b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/CubeMigrationCLI.java @@ -470,7 +470,7 @@ public class CubeMigrationCLI { value = Bytes.toBytes(valueString); } Put put = new Put(Bytes.toBytes(cubeId)); - put.add(family, column, value); + put.addColumn(family, column, value); destAclHtable.put(put); } } http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/DeployCoprocessorCLI.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/DeployCoprocessorCLI.java b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/DeployCoprocessorCLI.java index d51b71e..cee0aa9 100644 --- a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/DeployCoprocessorCLI.java +++ b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/DeployCoprocessorCLI.java @@ -46,7 +46,6 @@ import org.apache.hadoop.hbase.TableName; import org.apache.hadoop.hbase.TableNotFoundException; import org.apache.hadoop.hbase.client.Admin; import org.apache.hadoop.hbase.client.Connection; -import org.apache.hadoop.hbase.io.ImmutableBytesWritable; import org.apache.kylin.common.KylinConfig; import org.apache.kylin.common.KylinVersion; import org.apache.kylin.common.util.Bytes; @@ -156,7 +155,7 @@ public class DeployCoprocessorCLI { ProjectInstance projectInstance = projectManager.getProject(p); List<RealizationEntry> cubeList = projectInstance.getRealizationEntries(RealizationType.CUBE); - for (RealizationEntry cube: cubeList) { + for (RealizationEntry cube : cubeList) { CubeInstance cubeInstance = cubeManager.getCube(cube.getRealization()); for (CubeSegment segment : cubeInstance.getSegments()) { String tableName = segment.getStorageLocationIdentifier(); @@ -440,7 +439,7 @@ public class DeployCoprocessorCLI { Matcher keyMatcher; Matcher valueMatcher; - for (Map.Entry<ImmutableBytesWritable, ImmutableBytesWritable> e : tableDescriptor.getValues().entrySet()) { + for (Map.Entry<org.apache.hadoop.hbase.util.Bytes, org.apache.hadoop.hbase.util.Bytes> e : tableDescriptor.getValues().entrySet()) { keyMatcher = HConstants.CP_HTD_ATTR_KEY_PATTERN.matcher(Bytes.toString(e.getKey().get())); if (!keyMatcher.matches()) { continue; http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/ExtendCubeToHybridCLI.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/ExtendCubeToHybridCLI.java b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/ExtendCubeToHybridCLI.java index 1cdb2f8..18b894a 100644 --- a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/ExtendCubeToHybridCLI.java +++ b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/ExtendCubeToHybridCLI.java @@ -254,7 +254,7 @@ public class ExtendCubeToHybridCLI { value = Bytes.toBytes(valueString); } Put put = new Put(Bytes.toBytes(newCubeId)); - put.add(family, column, value); + put.addColumn(family, column, value); aclHtable.put(put); } } http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/GridTableHBaseBenchmark.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/GridTableHBaseBenchmark.java b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/GridTableHBaseBenchmark.java index dd5f8fa..218841c 100644 --- a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/GridTableHBaseBenchmark.java +++ b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/GridTableHBaseBenchmark.java @@ -232,7 +232,7 @@ public class GridTableHBaseBenchmark { byte[] rowkey = Bytes.toBytes(i); Put put = new Put(rowkey); byte[] cell = randomBytes(); - put.add(CF, QN, cell); + put.addColumn(CF, QN, cell); table.put(put); nBytes += cell.length; dot(i, N_ROWS); http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/PingHBaseCLI.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/PingHBaseCLI.java b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/PingHBaseCLI.java index bba6745..ff038d1 100644 --- a/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/PingHBaseCLI.java +++ b/storage-hbase/src/main/java/org/apache/kylin/storage/hbase/util/PingHBaseCLI.java @@ -50,7 +50,8 @@ public class PingHBaseCLI { if (User.isHBaseSecurityEnabled(hconf)) { try { System.out.println("--------------Getting kerberos credential for user " + UserGroupInformation.getCurrentUser().getUserName()); - TokenUtil.obtainAndCacheToken(hconf, UserGroupInformation.getCurrentUser()); + Connection connection = HBaseConnection.get(); + TokenUtil.obtainAndCacheToken(connection, User.create(UserGroupInformation.getCurrentUser())); } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("--------------Error while getting kerberos credential for user " + UserGroupInformation.getCurrentUser().getUserName()); http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/test/java/org/apache/kylin/storage/hbase/steps/CubeHFileMapperTest.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/test/java/org/apache/kylin/storage/hbase/steps/CubeHFileMapperTest.java b/storage-hbase/src/test/java/org/apache/kylin/storage/hbase/steps/CubeHFileMapperTest.java index f8282d3..d019188 100644 --- a/storage-hbase/src/test/java/org/apache/kylin/storage/hbase/steps/CubeHFileMapperTest.java +++ b/storage-hbase/src/test/java/org/apache/kylin/storage/hbase/steps/CubeHFileMapperTest.java @@ -68,13 +68,23 @@ public class CubeHFileMapperTest { Pair<ImmutableBytesWritable, KeyValue> p2 = result.get(1); assertEquals(key, p1.getFirst()); - assertEquals("cf1", new String(p1.getSecond().getFamily())); - assertEquals("usd_amt", new String(p1.getSecond().getQualifier())); - assertEquals("35.43", new String(p1.getSecond().getValue())); + assertEquals("cf1", new String(copy(p1.getSecond()))); + assertEquals("usd_amt", new String(copy(p1.getSecond()))); + assertEquals("35.43", new String(copy(p1.getSecond()))); assertEquals(key, p2.getFirst()); - assertEquals("cf1", new String(p2.getSecond().getFamily())); - assertEquals("item_count", new String(p2.getSecond().getQualifier())); - assertEquals("2", new String(p2.getSecond().getValue())); + assertEquals("cf1", new String(copy(p2.getSecond()))); + assertEquals("item_count", new String(copy(p2.getSecond()))); + assertEquals("2", new String(copy(p2.getSecond()))); + } + + private byte[] copy(KeyValue kv) { + return copy(kv.getFamilyArray(), kv.getFamilyOffset(), kv.getFamilyLength()); + } + + private byte[] copy(byte[] array, int offset, int length) { + byte[] result = new byte[length]; + System.arraycopy(array, offset, result, 0, length); + return result; } } http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/storage-hbase/src/test/java/org/apache/kylin/storage/hbase/steps/TestHbaseClient.java ---------------------------------------------------------------------- diff --git a/storage-hbase/src/test/java/org/apache/kylin/storage/hbase/steps/TestHbaseClient.java b/storage-hbase/src/test/java/org/apache/kylin/storage/hbase/steps/TestHbaseClient.java index 2b8ecae..b77d2cb 100644 --- a/storage-hbase/src/test/java/org/apache/kylin/storage/hbase/steps/TestHbaseClient.java +++ b/storage-hbase/src/test/java/org/apache/kylin/storage/hbase/steps/TestHbaseClient.java @@ -22,8 +22,11 @@ import java.io.IOException; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.hbase.HBaseConfiguration; -import org.apache.hadoop.hbase.client.HTable; +import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.client.Connection; +import org.apache.hadoop.hbase.client.ConnectionFactory; import org.apache.hadoop.hbase.client.Put; +import org.apache.hadoop.hbase.client.Table; import org.apache.kylin.common.util.Bytes; /** @@ -89,13 +92,16 @@ public class TestHbaseClient { conf.set("hbase.zookeeper.quorum", "hbase_host"); conf.set("zookeeper.znode.parent", "/hbase-unsecure"); - HTable table = new HTable(conf, "test1"); + Connection connection = ConnectionFactory.createConnection(conf); + + Table table = connection.getTable(TableName.valueOf("test1")); Put put = new Put(Bytes.toBytes("row1")); - put.add(Bytes.toBytes("colfam1"), Bytes.toBytes("qual1"), Bytes.toBytes("val1")); - put.add(Bytes.toBytes("colfam1"), Bytes.toBytes("qual2"), Bytes.toBytes("val2")); + put.addColumn(Bytes.toBytes("colfam1"), Bytes.toBytes("qual1"), Bytes.toBytes("val1")); + put.addColumn(Bytes.toBytes("colfam1"), Bytes.toBytes("qual2"), Bytes.toBytes("val2")); table.put(put); table.close(); + connection.close(); } } http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/tool/src/main/java/org/apache/kylin/tool/CubeMigrationCLI.java ---------------------------------------------------------------------- diff --git a/tool/src/main/java/org/apache/kylin/tool/CubeMigrationCLI.java b/tool/src/main/java/org/apache/kylin/tool/CubeMigrationCLI.java index c162a76..eb4f492 100644 --- a/tool/src/main/java/org/apache/kylin/tool/CubeMigrationCLI.java +++ b/tool/src/main/java/org/apache/kylin/tool/CubeMigrationCLI.java @@ -26,16 +26,16 @@ import java.util.Map; import java.util.Set; import org.apache.commons.io.IOUtils; -import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.hadoop.hbase.Cell; import org.apache.hadoop.hbase.CellUtil; import org.apache.hadoop.hbase.HTableDescriptor; import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.client.Admin; +import org.apache.hadoop.hbase.client.Connection; import org.apache.hadoop.hbase.client.Delete; import org.apache.hadoop.hbase.client.Get; -import org.apache.hadoop.hbase.client.HBaseAdmin; import org.apache.hadoop.hbase.client.Put; import org.apache.hadoop.hbase.client.Result; import org.apache.hadoop.hbase.client.Table; @@ -90,7 +90,8 @@ public class CubeMigrationCLI { private ResourceStore srcStore; private ResourceStore dstStore; private FileSystem hdfsFS; - private HBaseAdmin hbaseAdmin; + private Admin hbaseAdmin; + private Connection connection; public static final String ACL_INFO_FAMILY = "i"; private static final String ACL_TABLE_NAME = "_acl"; @@ -134,8 +135,8 @@ public class CubeMigrationCLI { checkAndGetHbaseUrl(); - Configuration conf = HBaseConnection.getCurrentHBaseConfiguration(); - hbaseAdmin = new HBaseAdmin(conf); + connection = HBaseConnection.get(); + hbaseAdmin = connection.getAdmin(); hdfsFS = HadoopUtil.getWorkingFileSystem(); @@ -337,10 +338,10 @@ public class CubeMigrationCLI { String tableName = (String) opt.params[0]; System.out.println("CHANGE_HTABLE_HOST, table name: " + tableName); HTableDescriptor desc = hbaseAdmin.getTableDescriptor(TableName.valueOf(tableName)); - hbaseAdmin.disableTable(tableName); + hbaseAdmin.disableTable(TableName.valueOf(tableName)); desc.setValue(IRealizationConstants.HTableTag, dstConfig.getMetadataUrlPrefix()); - hbaseAdmin.modifyTable(tableName, desc); - hbaseAdmin.enableTable(tableName); + hbaseAdmin.modifyTable(TableName.valueOf(tableName), desc); + hbaseAdmin.enableTable(TableName.valueOf(tableName)); logger.info("CHANGE_HTABLE_HOST is completed"); break; } @@ -482,7 +483,7 @@ public class CubeMigrationCLI { value = Bytes.toBytes(valueString); } Put put = new Put(Bytes.toBytes(cubeId)); - put.add(family, column, value); + put.addColumn(family, column, value); destAclHtable.put(put); } } @@ -518,10 +519,10 @@ public class CubeMigrationCLI { case CHANGE_HTABLE_HOST: { String tableName = (String) opt.params[0]; HTableDescriptor desc = hbaseAdmin.getTableDescriptor(TableName.valueOf(tableName)); - hbaseAdmin.disableTable(tableName); + hbaseAdmin.disableTable(TableName.valueOf(tableName)); desc.setValue(IRealizationConstants.HTableTag, srcConfig.getMetadataUrlPrefix()); - hbaseAdmin.modifyTable(tableName, desc); - hbaseAdmin.enableTable(tableName); + hbaseAdmin.modifyTable(TableName.valueOf(tableName), desc); + hbaseAdmin.enableTable(TableName.valueOf(tableName)); break; } case COPY_FILE_IN_META: { http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/tool/src/main/java/org/apache/kylin/tool/CubeMigrationCheckCLI.java ---------------------------------------------------------------------- diff --git a/tool/src/main/java/org/apache/kylin/tool/CubeMigrationCheckCLI.java b/tool/src/main/java/org/apache/kylin/tool/CubeMigrationCheckCLI.java index 54fbbc0..52bad9d 100644 --- a/tool/src/main/java/org/apache/kylin/tool/CubeMigrationCheckCLI.java +++ b/tool/src/main/java/org/apache/kylin/tool/CubeMigrationCheckCLI.java @@ -29,7 +29,9 @@ import org.apache.commons.cli.ParseException; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.hbase.HTableDescriptor; import org.apache.hadoop.hbase.TableName; -import org.apache.hadoop.hbase.client.HBaseAdmin; +import org.apache.hadoop.hbase.client.Admin; +import org.apache.hadoop.hbase.client.Connection; +import org.apache.hadoop.hbase.client.ConnectionFactory; import org.apache.kylin.common.KylinConfig; import org.apache.kylin.common.util.OptionsHelper; import org.apache.kylin.cube.CubeInstance; @@ -61,7 +63,8 @@ public class CubeMigrationCheckCLI { private static final Option OPTION_CUBE = OptionBuilder.withArgName("cube").hasArg().isRequired(false).withDescription("The name of cube migrated").create("cube"); private KylinConfig dstCfg; - private HBaseAdmin hbaseAdmin; + private Admin hbaseAdmin; + private Connection connection; private List<String> issueExistHTables; private List<String> inconsistentHTables; @@ -123,6 +126,7 @@ public class CubeMigrationCheckCLI { } fixInconsistent(); printIssueExistingHTables(); + connection.close(); } public CubeMigrationCheckCLI(KylinConfig kylinConfig, Boolean isFix) throws IOException { @@ -130,7 +134,8 @@ public class CubeMigrationCheckCLI { this.ifFix = isFix; Configuration conf = HBaseConnection.getCurrentHBaseConfiguration(); - hbaseAdmin = new HBaseAdmin(conf); + connection = ConnectionFactory.createConnection(conf); + hbaseAdmin = connection.getAdmin(); issueExistHTables = Lists.newArrayList(); inconsistentHTables = Lists.newArrayList(); @@ -189,10 +194,10 @@ public class CubeMigrationCheckCLI { String[] sepNameList = segFullName.split(","); HTableDescriptor desc = hbaseAdmin.getTableDescriptor(TableName.valueOf(sepNameList[0])); logger.info("Change the host of htable " + sepNameList[0] + "belonging to cube " + sepNameList[1] + " from " + desc.getValue(IRealizationConstants.HTableTag) + " to " + dstCfg.getMetadataUrlPrefix()); - hbaseAdmin.disableTable(sepNameList[0]); + hbaseAdmin.disableTable(TableName.valueOf(sepNameList[0])); desc.setValue(IRealizationConstants.HTableTag, dstCfg.getMetadataUrlPrefix()); - hbaseAdmin.modifyTable(sepNameList[0], desc); - hbaseAdmin.enableTable(sepNameList[0]); + hbaseAdmin.modifyTable(TableName.valueOf(sepNameList[0]), desc); + hbaseAdmin.enableTable(TableName.valueOf(sepNameList[0])); } } else { logger.info("------ Inconsistent HTables Needed To Be Fixed ------"); http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/tool/src/main/java/org/apache/kylin/tool/ExtendCubeToHybridCLI.java ---------------------------------------------------------------------- diff --git a/tool/src/main/java/org/apache/kylin/tool/ExtendCubeToHybridCLI.java b/tool/src/main/java/org/apache/kylin/tool/ExtendCubeToHybridCLI.java index f52fc3e..d628ca2 100644 --- a/tool/src/main/java/org/apache/kylin/tool/ExtendCubeToHybridCLI.java +++ b/tool/src/main/java/org/apache/kylin/tool/ExtendCubeToHybridCLI.java @@ -250,7 +250,7 @@ public class ExtendCubeToHybridCLI { value = Bytes.toBytes(valueString); } Put put = new Put(Bytes.toBytes(newCubeId)); - put.add(family, column, value); + put.addColumn(family, column, value); aclHtable.put(put); } } http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/tool/src/main/java/org/apache/kylin/tool/StorageCleanupJob.java ---------------------------------------------------------------------- diff --git a/tool/src/main/java/org/apache/kylin/tool/StorageCleanupJob.java b/tool/src/main/java/org/apache/kylin/tool/StorageCleanupJob.java index f1a3ebe..a725f23 100644 --- a/tool/src/main/java/org/apache/kylin/tool/StorageCleanupJob.java +++ b/tool/src/main/java/org/apache/kylin/tool/StorageCleanupJob.java @@ -41,7 +41,10 @@ import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.hadoop.hbase.HBaseConfiguration; import org.apache.hadoop.hbase.HTableDescriptor; -import org.apache.hadoop.hbase.client.HBaseAdmin; +import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.client.Admin; +import org.apache.hadoop.hbase.client.Connection; +import org.apache.hadoop.hbase.client.ConnectionFactory; import org.apache.kylin.common.KylinConfig; import org.apache.kylin.common.util.AbstractApplication; import org.apache.kylin.common.util.CliCommandExecutor; @@ -82,7 +85,8 @@ public class StorageCleanupJob extends AbstractApplication { private void cleanUnusedHBaseTables(Configuration conf) throws IOException { CubeManager cubeMgr = CubeManager.getInstance(KylinConfig.getInstanceFromEnv()); // get all kylin hbase tables - try (HBaseAdmin hbaseAdmin = new HBaseAdmin(conf)) { + Connection connection = ConnectionFactory.createConnection(conf); + try (Admin hbaseAdmin = connection.getAdmin()) { String tableNamePrefix = IRealizationConstants.SharedHbaseStorageLocationPrefix; HTableDescriptor[] tableDescriptors = hbaseAdmin.listTables(tableNamePrefix + ".*"); List<String> allTablesNeedToBeDropped = new ArrayList<String>(); @@ -129,6 +133,8 @@ public class StorageCleanupJob extends AbstractApplication { } System.out.println("----------------------------------------------------"); } + } finally { + connection.close(); } } @@ -157,12 +163,12 @@ public class StorageCleanupJob extends AbstractApplication { } class DeleteHTableRunnable implements Callable { - HBaseAdmin hbaseAdmin; - String htableName; + Admin hbaseAdmin; + TableName htableName; - DeleteHTableRunnable(HBaseAdmin hbaseAdmin, String htableName) { + DeleteHTableRunnable(Admin hbaseAdmin, String htableName) { this.hbaseAdmin = hbaseAdmin; - this.htableName = htableName; + this.htableName = TableName.valueOf(htableName); } public Object call() throws Exception { http://git-wip-us.apache.org/repos/asf/kylin/blob/dcb7cf47/tool/src/main/java/org/apache/kylin/tool/util/ToolUtil.java ---------------------------------------------------------------------- diff --git a/tool/src/main/java/org/apache/kylin/tool/util/ToolUtil.java b/tool/src/main/java/org/apache/kylin/tool/util/ToolUtil.java index c41d6a8..7eae46c 100644 --- a/tool/src/main/java/org/apache/kylin/tool/util/ToolUtil.java +++ b/tool/src/main/java/org/apache/kylin/tool/util/ToolUtil.java @@ -29,7 +29,9 @@ import org.apache.commons.lang.StringUtils; import org.apache.hadoop.hbase.HBaseConfiguration; import org.apache.hadoop.hbase.HTableDescriptor; import org.apache.hadoop.hbase.TableName; -import org.apache.hadoop.hbase.client.HBaseAdmin; +import org.apache.hadoop.hbase.client.Admin; +import org.apache.hadoop.hbase.client.Connection; +import org.apache.hadoop.hbase.client.ConnectionFactory; import org.apache.kylin.common.KylinConfig; import org.apache.kylin.common.util.HadoopUtil; import org.apache.kylin.storage.hbase.HBaseConnection; @@ -51,12 +53,15 @@ public class ToolUtil { } public static String getHBaseMetaStoreId() throws IOException { - try (final HBaseAdmin hbaseAdmin = new HBaseAdmin(HBaseConfiguration.create(HadoopUtil.getCurrentConfiguration()))) { + Connection connection = ConnectionFactory.createConnection(HBaseConfiguration.create(HadoopUtil.getCurrentConfiguration())); + try (final Admin hbaseAdmin = connection.getAdmin()) { final String metaStoreName = KylinConfig.getInstanceFromEnv().getMetadataUrlPrefix(); final HTableDescriptor desc = hbaseAdmin.getTableDescriptor(TableName.valueOf(metaStoreName)); return desc.getValue(HBaseConnection.HTABLE_UUID_TAG); } catch (Exception e) { return null; + } finally { + connection.close(); } }
