From d4753a7d8cb6f17d4da72c87a2fabc160c62fbf8 Mon Sep 17 00:00:00 2001 From: tushengxia Date: Wed, 23 Dec 2020 17:20:49 +0800 Subject: [PATCH] improve split plan and get high concurrency to improve performance fix ut fix ut for hbase remove redundant parameters of split add hbase limit push down remove redundant code check update hbase docs for splitbychar fix problems for hbase create table --- hetu-docs/en/connector/HBase.md | 2 +- hetu-docs/zh/connector/HBase.md | 1 + hetu-hbase/pom.xml | 168 ++++------------ .../hbase/conf/HBaseTableProperties.java | 86 ++++---- .../hbase/connector/HBaseConnection.java | 11 +- .../connector/HBaseConnectorFactory.java | 5 +- .../hbase/connector/HBaseTableHandle.java | 30 ++- .../metadata/HBaseConnectorMetadata.java | 59 ++++-- .../plugin/hbase/metadata/HBaseTable.java | 39 +++- .../hbase/metadata/HetuHBaseMetastore.java | 4 +- .../plugin/hbase/query/HBasePageSink.java | 19 +- .../plugin/hbase/query/HBaseRecordCursor.java | 17 -- .../plugin/hbase/query/HBaseRecordSet.java | 30 +-- .../core/plugin/hbase/split/HBaseSplit.java | 25 +-- .../plugin/hbase/split/HBaseSplitManager.java | 184 ++++++++++++------ .../core/plugin/hbase/utils/Constants.java | 73 ++----- .../hetu/core/plugin/hbase/utils/Utils.java | 5 +- .../io/hetu/core/plugin/hbase/TestHBase.java | 6 +- .../io/hetu/core/plugin/hbase/TestQuery.java | 15 +- .../core/plugin/hbase/client/TestUtils.java | 8 +- .../hbase/connector/TestHBaseConnector.java | 8 +- .../hbase/connector/TestHBaseTableHandle.java | 7 +- .../plugin/hbase/metadata/TestHBaseTable.java | 1 + .../metadata/TestHetuHBaseMetastore.java | 9 +- .../hbase/metadata/TestingHetuMetastore.java | 3 +- .../plugin/hbase/split/TestHBaseSplit.java | 4 - .../hbase/split/TestHbaseSplitManager.java | 11 +- 27 files changed, 434 insertions(+), 396 deletions(-) diff --git a/hetu-docs/en/connector/HBase.md b/hetu-docs/en/connector/HBase.md index 0876de8f4..d48ef2a5e 100644 --- a/hetu-docs/en/connector/HBase.md +++ b/hetu-docs/en/connector/HBase.md @@ -96,7 +96,7 @@ debug=true; | row_id | String | The first column name | No | row_id is the column name corresponding to rowkey in the hbase table | | hbase_table_name | String | null | No | hbase_table_name specifies the tablespace and table name on the hbase data source to be linked, use ":" connect tablespace and table name, the default tablespace is "default". | | external | Boolean | true | No | If external is true, it means that the table is a mapping table of the table in the hbase data source. It does not support deleting the original table on the hbase data source; if external is false, the table on hbase data source will be deleted at the same time as the local hbase table is deleted. | - +| split\_by\_char| String| 0~9,a~z,A~Z| No| split\_by\_char is the basis for sharding, if the first character of rowkey consists of digits, sharding can be performed based on different digits to improve concurrency. Different types of symbols are separated by commas. If this parameter is set incorrectly, the query result may be incomplete. Set it based on the actual rowkey.| ## Data Types diff --git a/hetu-docs/zh/connector/HBase.md b/hetu-docs/zh/connector/HBase.md index 1875c7774..c7cfa5ecd 100644 --- a/hetu-docs/zh/connector/HBase.md +++ b/hetu-docs/zh/connector/HBase.md @@ -97,6 +97,7 @@ debug=true; | row\_id| String| 第一个列名| 否| row\_id为HBase表中RowKey对应的列名| | hbase\_table\_name| String| NULL| 否| hbase\_table\_name指定要链接的HBase数据源上的表空间和表名,使用“:”连接表空间和表名,默认表空间为“default”。| | external| Boolean| true| 否| 如果external为true,表示该表是HBase数据源中表的映射表。不支持删除HBase数据源上原有的表。当external为false时,删除本地HBase表的同时也会删除HBase数据源上的表。| +| split\_by\_char| String| 0~9,a~z,A~Z| 否| split\_by\_char为分片切割的依据,若RowKey的第一个字符由数字构成,则可以根据不同的数字进行分片切割,提高查询并发度。不同类型的符号用逗号隔开。如果设置不当,会导致查询数据结果不完整,请根据RowKey的实际情况进行配置。| ## 数据说明 diff --git a/hetu-hbase/pom.xml b/hetu-hbase/pom.xml index 70c6dfe50..3e1e5b090 100644 --- a/hetu-hbase/pom.xml +++ b/hetu-hbase/pom.xml @@ -37,29 +37,6 @@ - - org.weakref - jmxutils - ${dep.jmxutils.version} - compile - - - - io.prestosql.hadoop - hadoop-apache - ${dep.prestosql.apache.hadoop.version} - - - org.antlr - antlr4-runtime - - - com.fasterxml.jackson.core - jackson-databind - - - - commons-codec commons-codec @@ -71,6 +48,11 @@ hetu-common + + io.prestosql.hadoop + hadoop-apache + + io.hetu.core @@ -149,127 +131,23 @@ ${version.commons-lang3} - - org.apache.hbase - hbase-common - ${version.hbase} - - - commons-codec - commons-codec - - - jackson-databind - com.fasterxml.jackson.core - - - hadoop-common - org.apache.hadoop - - - guice-servlet - com.google.inject.extensions - - - jackson-annotations - com.fasterxml.jackson.core - - - commons-logging - commons-logging - - - jaxb-api - javax.xml.bind - - - jackson-jaxrs-json-provider - com.fasterxml.jackson.jaxrs - - - jackson-module-jaxb-annotations - com.fasterxml.jackson.module - - - jackson-core - com.fasterxml.jackson.core - - - hadoop-yarn-common - org.apache.hadoop - - - junit - junit - - - protobuf-java - com.google.protobuf - - - hadoop-annotations - org.apache.hadoop - - - jetty-util - org.mortbay.jetty - - - hadoop-mapreduce-client-core - org.apache.hadoop - - - error_prone_annotations - com.google.errorprone - - - - org.apache.hbase hbase-client ${version.hbase} - commons-codec - commons-codec - - - htrace-core4 - org.apache.htrace + error_prone_annotations + com.google.errorprone jcodings org.jruby.jcodings - - jackson-databind - com.fasterxml.jackson.core - - - hbase-common - org.apache.hbase - - - commons-logging - commons-logging - - - junit - junit - hadoop-common org.apache.hadoop - - netty-all - io.netty - - - hadoop-mapreduce-client-core - org.apache.hadoop - hadoop-auth org.apache.hadoop @@ -287,6 +165,19 @@ com.fasterxml.jackson.core runtime + + + org.apache.hbase + hbase-common + ${version.hbase} + + + hadoop-common + org.apache.hadoop + + + + org.testng @@ -294,10 +185,27 @@ test + + org.weakref + jmxutils + ${dep.jmxutils.version} + compile + + io.hetu.core presto-tests test + + + plexus-cipher + org.sonatype.plexus + + + plexus-classworlds + org.codehaus.plexus + + @@ -353,4 +261,4 @@ - + \ No newline at end of file diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/conf/HBaseTableProperties.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/conf/HBaseTableProperties.java index fa2a444c0..9e1f0bad1 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/conf/HBaseTableProperties.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/conf/HBaseTableProperties.java @@ -78,6 +78,10 @@ public final class HBaseTableProperties * hbase_table_name */ public static final String HBASE_TABLE_NAME = "hbase_table_name"; + /** + * split_by_char + */ + public static final String SPLIT_BY_CHAR = "split_by_char"; private static final String COLON_SPLITTER = ":"; private static final String COMMA_SPLITTER = ","; private static final String PIPE_SPLITTER = "\\|"; @@ -146,8 +150,14 @@ public final class HBaseTableProperties "Point to the table of hbase, establish the mapping of external access", null, false); + PropertyMetadata s8 = + stringProperty( + SPLIT_BY_CHAR, + "Create multi-splits by the first char of rowKey", + "0~9,a~z,A~Z", + false); - tableProperties = ImmutableList.of(s1, s2, s3, s4, s5, s6, s7); + tableProperties = ImmutableList.of(s1, s2, s3, s4, s5, s6, s7, s8); } public List> getTableProperties() @@ -165,16 +175,14 @@ public final class HBaseTableProperties */ public static Optional>> getColumnMapping(Map tableProperties) { - requireNonNull(tableProperties); - - String strMapping = String.valueOf(tableProperties.get(COLUMN_MAPPING)); - if (Constants.S_NULL.equals(strMapping)) { + Optional strMapping = getStringFromProperties(COLUMN_MAPPING, tableProperties); + if (!strMapping.isPresent()) { return Optional.empty(); } // Parse out the column mapping list : "column_name:family:qualifier" ImmutableMap.Builder> mapping = ImmutableMap.builder(); - for (String m1 : strMapping.split(COMMA_SPLITTER)) { + for (String m1 : strMapping.get().split(COMMA_SPLITTER)) { String[] tokens = m1.trim().split(COLON_SPLITTER); checkState( tokens.length == Constants.NUMBER3, @@ -194,14 +202,13 @@ public final class HBaseTableProperties */ public static Optional> getIndexColumns(Map tableProperties) { - requireNonNull(tableProperties); - - String indexColumns = String.valueOf(tableProperties.get(INDEX_COLUMNS)); - if (Constants.S_NULL.equals(indexColumns)) { + Optional indexColumns = getStringFromProperties(INDEX_COLUMNS, tableProperties); + if (!indexColumns.isPresent()) { return Optional.empty(); } - - return Optional.of(Arrays.asList(StringUtils.split(indexColumns, ','))); + else { + return Optional.of(Arrays.asList(StringUtils.split(indexColumns.get(), ','))); + } } /** @@ -212,14 +219,7 @@ public final class HBaseTableProperties */ public static Optional getIndexColumnsAsStr(Map tableProperties) { - requireNonNull(tableProperties); - - String indexColumns = String.valueOf(tableProperties.get(INDEX_COLUMNS)); - if (Constants.S_NULL.equals(indexColumns)) { - return Optional.empty(); - } - - return Optional.of(indexColumns); + return getStringFromProperties(INDEX_COLUMNS, tableProperties); } /** @@ -230,17 +230,15 @@ public final class HBaseTableProperties */ public static Optional>> getLocalityGroups(Map tableProperties) { - requireNonNull(tableProperties); - - String groupStr = String.valueOf(tableProperties.get(LOCALITY_GROUPS)); - if (Constants.S_NULL.equals(groupStr)) { + Optional groupStr = getStringFromProperties(LOCALITY_GROUPS, tableProperties); + if (!groupStr.isPresent()) { return Optional.empty(); } ImmutableMap.Builder> groups = ImmutableMap.builder(); // Split all configured locality groups - for (String group : groupStr.split(PIPE_SPLITTER)) { + for (String group : groupStr.get().split(PIPE_SPLITTER)) { String[] locGroups = group.trim().split(COLON_SPLITTER); if (locGroups.length != Constants.NUMBER2) { throw new PrestoException( @@ -269,13 +267,7 @@ public final class HBaseTableProperties */ public static Optional getRowId(Map tableProperties) { - requireNonNull(tableProperties); - - String rowId = String.valueOf(tableProperties.get(ROW_ID)); - if (Constants.S_NULL.equals(rowId)) { - return Optional.empty(); - } - return Optional.ofNullable(rowId); + return getStringFromProperties(ROW_ID, tableProperties); } /** @@ -286,13 +278,7 @@ public final class HBaseTableProperties */ public static Optional getSerializerClass(Map tableProperties) { - requireNonNull(tableProperties); - - String serializerClass = String.valueOf(tableProperties.get(SERIALIZER)); - if (Constants.S_NULL.equals(serializerClass)) { - return Optional.empty(); - } - return Optional.ofNullable(serializerClass); + return getStringFromProperties(SERIALIZER, tableProperties); } /** @@ -316,13 +302,29 @@ public final class HBaseTableProperties * @return Optional */ public static Optional getHBaseTableName(Map tableProperties) + { + return getStringFromProperties(HBASE_TABLE_NAME, tableProperties); + } + + /** + * getSplitByChar + * + * @param tableProperties tableProperties + * @return Optional + */ + public static Optional getSplitByChar(Map tableProperties) + { + return getStringFromProperties(SPLIT_BY_CHAR, tableProperties); + } + + private static Optional getStringFromProperties(String key, Map tableProperties) { requireNonNull(tableProperties); - String hbaseTableName = String.valueOf(tableProperties.get(HBASE_TABLE_NAME)); - if (Constants.S_NULL.equals(hbaseTableName)) { + String value = String.valueOf(tableProperties.get(key)); + if (Constants.S_NULL.equals(value)) { return Optional.empty(); } - return Optional.ofNullable(hbaseTableName); + return Optional.ofNullable(value); } } diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseConnection.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseConnection.java index c3fbd94cd..85a696a34 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseConnection.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseConnection.java @@ -542,7 +542,8 @@ public class HBaseConnection HBaseTableProperties.isExternal(tableProperties), HBaseTableProperties.getSerializerClass(tableProperties), HBaseTableProperties.getIndexColumnsAsStr(meta.getProperties()), - HBaseTableProperties.getHBaseTableName(meta.getProperties())); + HBaseTableProperties.getHBaseTableName(meta.getProperties()), + HBaseTableProperties.getSplitByChar(meta.getProperties())); table.setColumnsToMap(columnHandleMap); @@ -564,7 +565,9 @@ public class HBaseConnection } } else { - if (table.isExternal()) { + String hbaseTableName = table.getFullTableName().replace(Constants.POINT, Constants.SEPARATOR); + table.setHbaseTableName(Optional.of(hbaseTableName)); + if (table.isExternal() && !existTable(hbaseTableName)) { throw new PrestoException( HBaseErrorCode.HBASE_CREATE_ERROR, format("Use lk creating new HBase table [%s], we must specify 'with(external=false)'. ", table.getTable())); @@ -572,12 +575,8 @@ public class HBaseConnection // create namespace if not exist createNamespaceIfNotExist(this.getHbaseAdmin(), table.getSchema()); - - String hbaseTableName = table.getFullTableName().replace(Constants.POINT, Constants.SEPARATOR); - table.setHbaseTableName(Optional.ofNullable(hbaseTableName)); // create hbase table createHBaseTable(table); - // save tableCatalog to memory and file hBaseMetastore.addHBaseTable(table); return table; diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseConnectorFactory.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseConnectorFactory.java index 7356269c0..e1c0b130d 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseConnectorFactory.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseConnectorFactory.java @@ -75,12 +75,13 @@ public class HBaseConnectorFactory public Connector create(String catalogName, Map config, ConnectorContext context) { requireNonNull(config, "config is null"); - HBaseConnectorId.setConnectorId(catalogName); try (ThreadContextClassLoader ignored = new ThreadContextClassLoader(classLoader)) { hetuMetastore = context.getHetuMetastore(); - Bootstrap app = new Bootstrap(new HBaseModule(), this.module); + Bootstrap app = new Bootstrap(binder -> + binder.bind(HetuMetastore.class).toInstance(context.getHetuMetastore()), + new HBaseModule(), this.module); Injector injector = app.strictConfig().doNotInitializeLogging().quiet().setRequiredConfigurationProperties(config).initialize(); return injector.getInstance(HBaseConnector.class); diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseTableHandle.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseTableHandle.java index 89c3091b5..f73e517a3 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseTableHandle.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/connector/HBaseTableHandle.java @@ -29,6 +29,7 @@ import io.prestosql.spi.predicate.TupleDomain; import java.util.List; import java.util.Objects; import java.util.Optional; +import java.util.OptionalLong; import static com.google.common.base.MoreObjects.toStringHelper; import static java.util.Objects.requireNonNull; @@ -51,6 +52,7 @@ public class HBaseTableHandle private final TupleDomain constraint; private final List columns; private final int rowIdOrdinal; + private final OptionalLong limit; /** * constructor @@ -65,6 +67,7 @@ public class HBaseTableHandle * @param constraint constraint * @param columns columns * @param rowIdOrdinal rowIdOrdinal + * @param limit limit */ @JsonCreator public HBaseTableHandle( @@ -77,7 +80,8 @@ public class HBaseTableHandle @JsonProperty("indexColumns") String indexColumns, @JsonProperty("constraint") TupleDomain constraint, @JsonProperty("columns") List columns, - @JsonProperty("rowIdOrdinal") int rowIdOrdinal) + @JsonProperty("rowIdOrdinal") int rowIdOrdinal, + @JsonProperty("limit") OptionalLong limit) { this.external = external; this.rowId = requireNonNull(rowId, "rowId is null"); @@ -89,6 +93,7 @@ public class HBaseTableHandle this.constraint = constraint; this.columns = columns; this.rowIdOrdinal = rowIdOrdinal; + this.limit = requireNonNull(limit, "limit is null"); } /** @@ -100,6 +105,7 @@ public class HBaseTableHandle * @param columns columns * @param serializerClassName serializerClassName * @param hbaseTableName hbaseTableName + * @param limit limit */ public HBaseTableHandle( String schema, @@ -107,7 +113,8 @@ public class HBaseTableHandle int rowIdOrdinal, List columns, String serializerClassName, - Optional hbaseTableName) + Optional hbaseTableName, + OptionalLong limit) { this( schema, @@ -119,7 +126,8 @@ public class HBaseTableHandle "", TupleDomain.all(), columns, - rowIdOrdinal); + rowIdOrdinal, + limit); } /** @@ -132,6 +140,7 @@ public class HBaseTableHandle * @param serializerClassName serializerClassName * @param hbaseTableName hbaseTableName * @param indexColumns indexColumns + * @param limit limit */ public HBaseTableHandle( String schema, @@ -140,7 +149,8 @@ public class HBaseTableHandle boolean external, String serializerClassName, Optional hbaseTableName, - String indexColumns) + String indexColumns, + OptionalLong limit) { this( schema, @@ -152,7 +162,8 @@ public class HBaseTableHandle indexColumns, TupleDomain.all(), ImmutableList.of(), - 0); + 0, + limit); } @JsonProperty @@ -215,6 +226,12 @@ public class HBaseTableHandle return indexColumns; } + @JsonProperty + public OptionalLong getLimit() + { + return limit; + } + /** * toSchemaTableName * @@ -262,7 +279,8 @@ public class HBaseTableHandle oldHBaseConnectorTableHandle.getIndexColumns(), oldHBaseConnectorTableHandle.getConstraint(), oldHBaseConnectorTableHandle.getColumns(), - oldHBaseConnectorTableHandle.getRowIdOrdinal()); + oldHBaseConnectorTableHandle.getRowIdOrdinal(), + oldHBaseConnectorTableHandle.getLimit()); } @Override diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HBaseConnectorMetadata.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HBaseConnectorMetadata.java index 156968fae..209c869d5 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HBaseConnectorMetadata.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HBaseConnectorMetadata.java @@ -20,12 +20,9 @@ import com.google.common.collect.ImmutableSet; import com.google.inject.Inject; import io.airlift.log.Logger; import io.airlift.slice.Slice; -import io.hetu.core.plugin.hbase.conf.HBaseTableProperties; import io.hetu.core.plugin.hbase.connector.HBaseColumnHandle; import io.hetu.core.plugin.hbase.connector.HBaseConnection; import io.hetu.core.plugin.hbase.connector.HBaseTableHandle; -import io.hetu.core.plugin.hbase.utils.HBaseErrorCode; -import io.prestosql.spi.PrestoException; import io.prestosql.spi.connector.ColumnHandle; import io.prestosql.spi.connector.ColumnMetadata; import io.prestosql.spi.connector.ConnectorInsertTableHandle; @@ -40,6 +37,7 @@ import io.prestosql.spi.connector.ConnectorTableMetadata; import io.prestosql.spi.connector.ConnectorTableProperties; import io.prestosql.spi.connector.Constraint; import io.prestosql.spi.connector.ConstraintApplicationResult; +import io.prestosql.spi.connector.LimitApplicationResult; import io.prestosql.spi.connector.SchemaTableName; import io.prestosql.spi.connector.SchemaTablePrefix; import io.prestosql.spi.connector.TableNotFoundException; @@ -58,7 +56,6 @@ import java.util.Set; import java.util.concurrent.atomic.AtomicReference; import static com.google.common.base.Preconditions.checkState; -import static java.lang.String.format; import static java.util.Objects.requireNonNull; /** @@ -112,7 +109,8 @@ public class HBaseConnectorMetadata table.isExternal(), table.getSerializerClassName(), table.getHbaseTableName(), - table.getIndexColumns()); + table.getIndexColumns(), + OptionalLong.empty()); } @Override @@ -168,7 +166,7 @@ public class HBaseConnectorMetadata return Optional.ofNullable(tableMetadata).orElse(tableMetadata); } - return new ConnectorTableMetadata(tableName, table.getColumnMetadatas()); + return new ConnectorTableMetadata(tableName, table.getColumnMetadatas(), table.getTableProperties()); } @Override @@ -177,12 +175,6 @@ public class HBaseConnectorMetadata { checkNoRollback(); SchemaTableName tableName = tableMetadata.getTable(); - if (HBaseTableProperties.isExternal(tableMetadata.getProperties())) { - throw new PrestoException( - HBaseErrorCode.HBASE_CREATE_ERROR, - format("Use lk creating new HBase table [%s], we must specify 'with(external=false)'. ", tableName.toString())); - } - HBaseTable table = hbaseConn.createTable(tableMetadata); // support create table xxx as select * from yyy HBaseTableHandle handle = @@ -196,7 +188,8 @@ public class HBaseConnectorMetadata table.getIndexColumns(), TupleDomain.all(), table.getColumns(), - table.getRowIdOrdinal()); + table.getRowIdOrdinal(), + OptionalLong.empty()); setRollback(() -> rollbackCreateTable(table)); return handle; @@ -227,7 +220,8 @@ public class HBaseConnectorMetadata table.getRowIdOrdinal(), table.getColumns(), table.getSerializerClassName(), - table.getHbaseTableName()); + table.getHbaseTableName(), + OptionalLong.empty()); return handle; } @@ -406,10 +400,45 @@ public class HBaseConnectorMetadata tableHandle.getFullTableName(), newDomain, tableHandle.getColumns(), - tableHandle.getRowIdOrdinal()); + tableHandle.getRowIdOrdinal(), + tableHandle.getLimit()); return Optional.of(new ConstraintApplicationResult<>(newTableHandle, constraint.getSummary())); } + @Override + public Optional> applyLimit( + ConnectorSession session, ConnectorTableHandle connectorTableHandle, long limit) + { + HBaseTableHandle tableHandle; + + if (connectorTableHandle instanceof HBaseTableHandle) { + tableHandle = (HBaseTableHandle) connectorTableHandle; + } + else { + return Optional.empty(); + } + + if (tableHandle.getLimit().isPresent() && tableHandle.getLimit().getAsLong() <= limit) { + return Optional.empty(); + } + + HBaseTableHandle newTableHandle = + new HBaseTableHandle( + tableHandle.getSchema(), + tableHandle.getTable(), + tableHandle.getRowId(), + tableHandle.isExternal(), + tableHandle.getSerializerClassName(), + tableHandle.getHbaseTableName(), + tableHandle.getFullTableName(), + tableHandle.getConstraint(), + tableHandle.getColumns(), + tableHandle.getRowIdOrdinal(), + OptionalLong.of(limit)); + + return Optional.of(new LimitApplicationResult<>(newTableHandle, false)); + } + @Override public ColumnHandle getUpdateRowIdColumnHandle(ConnectorSession session, ConnectorTableHandle tableHandle) { diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HBaseTable.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HBaseTable.java index 89947b82d..0b02ed112 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HBaseTable.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HBaseTable.java @@ -18,6 +18,7 @@ import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonProperty; import io.airlift.log.Logger; +import io.hetu.core.plugin.hbase.conf.HBaseTableProperties; import io.hetu.core.plugin.hbase.connector.HBaseColumnHandle; import io.hetu.core.plugin.hbase.utils.Constants; import io.prestosql.spi.connector.ColumnHandle; @@ -28,6 +29,7 @@ import org.codehaus.jettison.json.JSONException; import org.codehaus.jettison.json.JSONObject; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Optional; @@ -55,6 +57,7 @@ public class HBaseTable private final SchemaTableName schemaTableName; private String indexColumns; private Optional hbaseTableName; + private Optional splitByChar; private Map columnsMap; /** @@ -68,6 +71,7 @@ public class HBaseTable * @param serializerClassName serializerClassName * @param indexColumns indexColumns * @param hbaseTableName hbaseTableName + * @param splitByChar splitByChar */ @JsonCreator public HBaseTable( @@ -78,7 +82,8 @@ public class HBaseTable @JsonProperty("external") boolean external, @JsonProperty("serializerClassName") Optional serializerClassName, @JsonProperty("indexColumns") Optional indexColumns, - @JsonProperty("hbaseTableName") Optional hbaseTableName) + @JsonProperty("hbaseTableName") Optional hbaseTableName, + @JsonProperty("splitByChar") Optional splitByChar) { this.external = external; this.rowId = requireNonNull(rowId, "rowId is null"); @@ -88,6 +93,7 @@ public class HBaseTable this.serializerClassName = serializerClassName.orElse(""); this.indexColumns = indexColumns.orElse(""); this.hbaseTableName = hbaseTableName; + this.splitByChar = splitByChar; boolean flag = false; int rowIdOrdinal0 = 0; @@ -175,6 +181,35 @@ public class HBaseTable return this.columnsMap; } + /** + * getSplitByChar + * + * @return all chars range + */ + @JsonProperty + public Optional getSplitByChar() + { + return splitByChar; + } + + /** + * get table properties + * + * @return all table properties + */ + @JsonProperty + public Map getTableProperties() + { + Map properties = new HashMap<>(); + properties.put(HBaseTableProperties.EXTERNAL, isExternal()); + properties.put(HBaseTableProperties.INDEX_COLUMNS, getIndexColumns()); + properties.put(HBaseTableProperties.SPLIT_BY_CHAR, getSplitByChar().orElse("")); + properties.put(HBaseTableProperties.HBASE_TABLE_NAME, getHbaseTableName().orElse("")); + properties.put(HBaseTableProperties.ROW_ID, getRowId()); + properties.put(HBaseTableProperties.SERIALIZER, getSerializerClassName()); + return properties; + } + /** * setHbaseTableName * @@ -290,6 +325,7 @@ public class HBaseTable .add("external", external) .add("serializerClassName", serializerClassName) .add("hbaseTableName", hbaseTableName) + .add("splitByChar", splitByChar) .toString(); } @@ -312,6 +348,7 @@ public class HBaseTable jo.put("table", table); jo.put("indexColumns", indexColumns); jo.put("hbaseTableName", hbaseTableName.orElse("")); + jo.put("splitByChar", splitByChar.orElse("")); JSONArray jac = new JSONArray(); JSONObject joc; diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HetuHBaseMetastore.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HetuHBaseMetastore.java index ffa284dd5..88025f954 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HetuHBaseMetastore.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/metadata/HetuHBaseMetastore.java @@ -204,7 +204,8 @@ public class HetuHBaseMetastore Constants.S_TRUE.equals(tableParam.get(Constants.S_EXTERNAL).toLowerCase(Locale.ENGLISH)), Optional.ofNullable(tableParam.get(Constants.S_SERIALIZER_CLASS_NAME)), Optional.ofNullable(tableParam.get(Constants.S_INDEX_COLUMNS)), - Optional.ofNullable(tableParam.get(Constants.S_HBASE_TABLE_NAME))); + Optional.ofNullable(tableParam.get(Constants.S_HBASE_TABLE_NAME)), + Optional.ofNullable(tableParam.get(Constants.S_SPLIT_BY_CHAR))); hBaseTable.setColumnsToMap(columnsMap); return hBaseTable; @@ -242,6 +243,7 @@ public class HetuHBaseMetastore tableParam.put(Constants.S_ROWID, hBaseTable.getRowId()); tableParam.put(Constants.S_INDEX_COLUMNS, hBaseTable.getIndexColumns()); tableParam.put(Constants.S_HBASE_TABLE_NAME, hBaseTable.getHbaseTableName().orElse(null)); + tableParam.put(Constants.S_SPLIT_BY_CHAR, hBaseTable.getSplitByChar().orElse(null)); tableEntity.setParameters(tableParam); diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBasePageSink.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBasePageSink.java index 8a8e9cb20..965816a7d 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBasePageSink.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBasePageSink.java @@ -20,6 +20,7 @@ import io.airlift.slice.Slice; import io.hetu.core.plugin.hbase.connector.HBaseColumnHandle; import io.hetu.core.plugin.hbase.connector.HBaseConnection; import io.hetu.core.plugin.hbase.connector.HBaseTableHandle; +import io.hetu.core.plugin.hbase.utils.Constants; import io.hetu.core.plugin.hbase.utils.HBaseErrorCode; import io.hetu.core.plugin.hbase.utils.serializers.HBaseRowSerializer; import io.prestosql.spi.Page; @@ -86,13 +87,19 @@ public class HBasePageSink // For each position within the page List puts = new ArrayList<>(); - for (int position = 0; position < page.getPositionCount(); ++position) { - // Convert Page to a Put, writing and indexing it - Put put = pageToPut(page, position); - puts.add(put); - } try { - hbaseConn.getConn().getTable(TableName.valueOf(tablename)).put(puts); + for (int position = 0; position < page.getPositionCount(); ++position) { + // Convert Page to a Put, writing and indexing it + Put put = pageToPut(page, position); + puts.add(put); + if (puts.size() >= Constants.PUT_BATCH_SIZE) { + hbaseConn.getConn().getTable(TableName.valueOf(tablename)).put(puts); + puts.clear(); + } + } + if (!puts.isEmpty()) { + hbaseConn.getConn().getTable(TableName.valueOf(tablename)).put(puts); + } } catch (IOException e) { LOG.error("appendPage PUT rejected by server... cause by %s", e.getMessage()); diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBaseRecordCursor.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBaseRecordCursor.java index 212aafae5..4052fb903 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBaseRecordCursor.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBaseRecordCursor.java @@ -26,7 +26,6 @@ import org.apache.commons.lang3.StringUtils; import org.apache.hadoop.hbase.client.Result; import org.apache.hadoop.hbase.client.ResultScanner; -import java.nio.charset.Charset; import java.util.Iterator; import java.util.List; @@ -187,22 +186,6 @@ public class HBaseRecordCursor if (iterator.hasNext()) { serializer.reset(); Result row = iterator.next(); - - bytesRead += row.getRow().length; - - byte[] bytes; - for (HBaseColumnHandle hc : columnHandles) { - if (!hc.getName().equals(rowIdName)) { - bytes = - row.getValue( - hc.getFamily().get().getBytes(Charset.forName("UTF-8")), - hc.getQualifier().get().getBytes(Charset.forName("UTF-8"))); - if (bytes != null) { - bytesRead += bytes.length; - } - } - } - serializer.deserialize(row, this.defaultValue); return true; } diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBaseRecordSet.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBaseRecordSet.java index 0540f432b..5c88e0124 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBaseRecordSet.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/query/HBaseRecordSet.java @@ -107,7 +107,6 @@ public class HBaseRecordSet List columnHandles) { requireNonNull(session, "session is null"); - rowIdName = table.getRowId(); this.split = split; this.conn = hbaseConn; @@ -135,6 +134,10 @@ public class HBaseRecordSet this.serializer.setMapping(hc.getName(), hc.getFamily().get(), hc.getQualifier().get()); } } + if (table.getLimit().isPresent() && table.getLimit().getAsLong() <= Integer.MAX_VALUE) { + scan.setLimit((int) table.getLimit().getAsLong()); + } + LOG.info("Worker handle split:" + split.toString()); } @Override @@ -177,14 +180,16 @@ public class HBaseRecordSet if (filters.getFilters().size() != 0) { scan.setFilter(filters); } - if (split.getStartRow() != null && !split.getStartRow().isEmpty()) { - scan.setStartRow(Bytes.toBytes(split.getStartRow())); + scan.withStartRow(Bytes.toBytes(split.getStartRow())); + } + if (split.getEndRow() != null && !split.getEndRow().isEmpty()) { + scan.withStopRow(Bytes.toBytes(split.getEndRow())); } - if (split.getEndRow() != null && !split.getEndRow().isEmpty()) { - scan.setStopRow(Bytes.toBytes(split.getEndRow())); - } + scan.setCaching(Constants.SCAN_CACHING_SIZE); + scan.setLoadColumnFamiliesOnDemand(true); + scan.setCacheBlocks(true); scanner = conn.getConn().getTable(TableName.valueOf(table.getHbaseTableName().get())).getScanner(scan); } catch (IOException e) { @@ -225,11 +230,10 @@ public class HBaseRecordSet { FilterList andFilters = new FilterList(FilterList.Operator.MUST_PASS_ALL); - // select count(rowKey) / rowKey from table_xxx; - if (this.columnHandles.size() == 1 - && this.columnHandles.get(0).getColumnName().equals(this.split.getRowKeyName())) { - scan.setCaching(Constants.SCAN_CACHING_SIZE); - scan.setCacheBlocks(false); + // select count(*) / count(rowKey) / rowKey from table_xxx; + if ((this.columnHandles.size() == 1 + && this.columnHandles.get(0).getColumnName().equals(this.split.getRowKeyName())) + || this.columnHandles.size() == 0) { andFilters.addFilter(new FirstKeyOnlyFilter()); andFilters.addFilter(new KeyOnlyFilter()); } @@ -398,11 +402,13 @@ public class HBaseRecordSet return new RowFilter(operator, new BinaryComparator(values)); } else { - return new SingleColumnValueFilter( + SingleColumnValueFilter filter = new SingleColumnValueFilter( Bytes.toBytes(columnHandle.getFamily().get()), Bytes.toBytes(columnHandle.getQualifier().get()), operator, values); + filter.setFilterIfMissing(true); + return filter; } } diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/split/HBaseSplit.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/split/HBaseSplit.java index d7937b18b..8baf95a41 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/split/HBaseSplit.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/split/HBaseSplit.java @@ -20,7 +20,6 @@ import io.hetu.core.plugin.hbase.connector.HBaseTableHandle; import io.prestosql.spi.HostAddress; import io.prestosql.spi.connector.ConnectorSplit; import io.prestosql.spi.predicate.Range; -import org.apache.hadoop.hbase.HRegionInfo; import java.util.List; import java.util.Map; @@ -43,8 +42,6 @@ public class HBaseSplit private final Map> ranges; - private final HRegionInfo regionInfo; - private final HBaseTableHandle tableHandle; private final boolean randomSplit; @@ -58,7 +55,6 @@ public class HBaseSplit * @param startRow startRow * @param endRow endRow * @param ranges search ranges - * @param regionInfo regionInfo * @param randomSplit randomSplit */ @JsonCreator @@ -69,7 +65,6 @@ public class HBaseSplit @JsonProperty("startRow") String startRow, @JsonProperty("endRow") String endRow, @JsonProperty("ranges") Map> ranges, - @JsonProperty("regionInfo") HRegionInfo regionInfo, @JsonProperty("randomSplit") boolean randomSplit) { this.rowKeyName = rowKeyName; @@ -78,7 +73,6 @@ public class HBaseSplit this.startRow = startRow; this.endRow = endRow; this.ranges = ranges; - this.regionInfo = regionInfo; this.randomSplit = randomSplit; } @@ -124,12 +118,6 @@ public class HBaseSplit return ranges; } - @JsonProperty - public HRegionInfo getRegionInfo() - { - return regionInfo; - } - @JsonProperty public HBaseTableHandle getTableHandle() { @@ -141,4 +129,17 @@ public class HBaseSplit { return randomSplit; } + + @Override + public String toString() + { + return "HBaseSplit{" + + "addresses='" + addresses + '\'' + + ", rowKeyName='" + rowKeyName + '\'' + + ", startRow='" + startRow + '\'' + + ", endRow='" + endRow + '\'' + + ", ranges='" + ranges + '\'' + + ", randomSplit='" + randomSplit + '\'' + + '}'; + } } diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/split/HBaseSplitManager.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/split/HBaseSplitManager.java index 43e0dd806..d0a9ec601 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/split/HBaseSplitManager.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/split/HBaseSplitManager.java @@ -20,10 +20,8 @@ import io.hetu.core.plugin.hbase.connector.HBaseColumnHandle; import io.hetu.core.plugin.hbase.connector.HBaseConnection; import io.hetu.core.plugin.hbase.connector.HBaseTableHandle; import io.hetu.core.plugin.hbase.utils.Constants; -import io.hetu.core.plugin.hbase.utils.HBaseErrorCode; import io.hetu.core.plugin.hbase.utils.Utils; import io.prestosql.spi.HostAddress; -import io.prestosql.spi.PrestoException; import io.prestosql.spi.connector.ColumnHandle; import io.prestosql.spi.connector.ConnectorSession; import io.prestosql.spi.connector.ConnectorSplitManager; @@ -35,18 +33,14 @@ import io.prestosql.spi.predicate.Domain; import io.prestosql.spi.predicate.Range; import io.prestosql.spi.predicate.TupleDomain; import org.apache.hadoop.hbase.TableName; -import org.apache.hadoop.hbase.TableNotFoundException; -import org.apache.hadoop.hbase.client.RegionLocator; -import org.apache.hadoop.hbase.util.Pair; -import java.io.IOException; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; - -import static java.lang.String.format; +import java.util.stream.Collectors; /** * HBaseSplitManager @@ -58,12 +52,12 @@ public class HBaseSplitManager { private static final Logger LOG = Logger.get(HBaseSplitManager.class); - private final HBaseConnection hBaseConnection; + private final HBaseConnection hbaseConnection; @Inject - public HBaseSplitManager(HBaseConnection hBaseConnection) + public HBaseSplitManager(HBaseConnection hbaseConnection) { - this.hBaseConnection = hBaseConnection; + this.hbaseConnection = hbaseConnection; } @Override @@ -89,59 +83,100 @@ public class HBaseSplitManager return new FixedSplitSource(splits); } + /** + * Get splits by slicing the rowKeys, according to the first character of rowKey (user can specify it when create + * table, the default value is "0~9,a~z,A~Z"), generate many startAndEndKey pairs. + * + * @param tupleDomain tupleDomain + * @param tableHandle tableHandle + * @return splits + */ private List getSplitsForScan(TupleDomain tupleDomain, HBaseTableHandle tableHandle) { List splits = new ArrayList<>(); - Pair startEndKeys = null; TableName hbaseTableName = TableName.valueOf(tableHandle.getHbaseTableName().get()); - - try { - if (hBaseConnection.getHbaseAdmin().getTableDescriptor(hbaseTableName) != null) { - RegionLocator regionLocator = - hBaseConnection.getConn().getRegionLocator(hbaseTableName); - startEndKeys = regionLocator.getStartEndKeys(); - } - } - catch (TableNotFoundException e) { - throw new PrestoException( - HBaseErrorCode.HBASE_TABLE_DNE, - format( - "table %s not found, maybe deleted by other user", tableHandle.getHbaseTableName().get())); - } - catch (IOException e) { - LOG.error(e.getMessage()); - } - Map> ranges = new HashMap<>(); Map predicates = tupleDomain.getDomains().get(); - predicates - .entrySet() - .forEach( - entry -> { - ColumnHandle handle = entry.getKey(); - if (handle instanceof HBaseColumnHandle) { - ranges.put( - ((HBaseColumnHandle) handle).getOrdinal(), - entry.getValue().getValues().getRanges().getOrderedRanges()); - } - }); + predicates.entrySet().forEach( + entry -> { + ColumnHandle handle = entry.getKey(); + if (handle instanceof HBaseColumnHandle) { + ranges.put( + ((HBaseColumnHandle) handle).getOrdinal(), + entry.getValue().getValues().getRanges().getOrderedRanges()); + } + }); List hostAddresses = new ArrayList<>(); - if (startEndKeys == null) { - throw new NullPointerException("null pointer found when getting splits for scan"); - } - for (int i = 0; i < startEndKeys.getFirst().length; i++) { - String startRow = new String(startEndKeys.getFirst()[i]); - String endRow = new String(startEndKeys.getSecond()[i]); + // splitByChar read from hetu metastore, the default value is "0~9" + String splitByChar = hbaseConnection.getTable(tableHandle.getTableName()).getSplitByChar().get(); + LOG.info("Create multi-splits by the first char of rowKey, table is " + hbaseTableName.getName() + + ", the range of first char is : " + splitByChar); + + List startAndEndRowKeys = + getStartAndEndKeys(splitByChar, Constants.START_END_KEYS_COUNT); + for (StartAndEndKey startAndEndRowKey : startAndEndRowKeys) { splits.add( new HBaseSplit( - tableHandle.getRowId(), tableHandle, hostAddresses, startRow, endRow, ranges, null, false)); + tableHandle.getRowId(), + tableHandle, + hostAddresses, + String.valueOf(startAndEndRowKey.start), + startAndEndRowKey.end + Constants.ROWKEY_TAIL, + ranges, + false)); } + printSplits("Scan", splits); return splits; } - private List getSplitsForBatchGet(TupleDomain tupleDomain, HBaseTableHandle table) + /** + * In order to get more splits to improve concurrency of tableScan, we slice the split by different character. + * HBase server support to use startRow and EndRow to get scanner. + * for example, splitByChar is 0~2, we will generate multi-pairs + * the size of pairs will less than startEndKeysCount, so we will calculate the pair gap first. + * (startKey = 0, endKey = 0),(startKey = 1, endKey = 1),(startKey = 2, endKey = 2) + * splitByChar is 0~9, a~z + * (startKey = 0, endKey = 1),(startKey = 2, endKey = 3)……(startKey = y, endKey = z) + * + * @param splitByChar range of the rowKey, value is like 0~9,A~Z,a~z or a~z,0~9 .. + * @param startEndKeysCount max number of key pairs + * @return start and end rowKeys + */ + private List getStartAndEndKeys(String splitByChar, int startEndKeysCount) + { + List allRanges = Arrays.stream(splitByChar.split(",")) + .map(StartAndEndKey::new).collect(Collectors.toList()); + int rangeLength = 0; + for (StartAndEndKey range : allRanges) { + rangeLength += (Math.abs(range.end - range.start) + 1); + } + + List startAndEndKeys = new ArrayList<>(); + // rounding step value + int gap = (int) Math.rint((rangeLength + 0.0) / startEndKeysCount); + // generate start and end keys + allRanges.forEach(range -> { + int realGap = gap == 0 ? 1 : gap; + for (char index = range.start; index <= range.end; index += realGap) { + char end = (index + realGap > range.end) ? range.end : (char) (index + realGap - 1); + startAndEndKeys.add(new StartAndEndKey(index, end)); + } + }); + + return startAndEndKeys; + } + + /** + * If the predicate of sql includes "rowKey='xxx'" or "rowKey in ('xxx','xxx')", + * we can specify rowkey values in each split, then performance will be good. + * + * @param tupleDomain tupleDomain + * @param tableHandle tableHandle + * @return splits + */ + private List getSplitsForBatchGet(TupleDomain tupleDomain, HBaseTableHandle tableHandle) { List splits = new ArrayList<>(); Domain rowIdDomain = null; @@ -150,13 +185,13 @@ public class HBaseSplitManager ColumnHandle handle = entry.getKey(); if (handle instanceof HBaseColumnHandle) { HBaseColumnHandle columnHandle = (HBaseColumnHandle) handle; - if (columnHandle.getOrdinal() == table.getRowIdOrdinal()) { + if (columnHandle.getOrdinal() == tableHandle.getRowIdOrdinal()) { rowIdDomain = entry.getValue(); } } } - List rowIds = rowIdDomain != null ? rowIdDomain.getValues().getRanges().getOrderedRanges() : new ArrayList<>(); + List rowIds = rowIdDomain != null ? rowIdDomain.getValues().getRanges().getOrderedRanges() : new ArrayList<>(); int maxSplitSize; // Each split has at least 20 pieces of data, and the maximum number of splits is 30. if (rowIds.size() / Constants.BATCHGET_SPLIT_RECORD_COUNT > Constants.BATCHGET_SPLIT_MAX_COUNT) { @@ -172,17 +207,48 @@ public class HBaseSplitManager while (currentIndex < rangeSize) { int endIndex = rangeSize - currentIndex > maxSplitSize ? (currentIndex + maxSplitSize) : rangeSize; Map> splitRange = new HashMap<>(); - splitRange.put(table.getRowIdOrdinal(), rowIds.subList(currentIndex, endIndex)); - splits.add(new HBaseSplit(table.getRowId(), table, hostAddresses, null, null, splitRange, null, false)); + splitRange.put(tableHandle.getRowIdOrdinal(), rowIds.subList(currentIndex, endIndex)); + splits.add(new HBaseSplit(tableHandle.getRowId(), tableHandle, hostAddresses, null, null, splitRange, false)); currentIndex = endIndex; } - for (HBaseSplit split : splits) { - if (LOG.isInfoEnabled()) { - LOG.info("Print Split: " + split.toString()); - } - } - + printSplits("Batch Get", splits); return splits; } + + private void printSplits(String scanType, List splits) + { + LOG.info("The final split count is " + splits.size() + "."); + for (HBaseSplit split : splits) { + LOG.info(scanType + ", Print Split: " + split.toString()); + } + } + + class StartAndEndKey + { + private final char start; + private final char end; + + StartAndEndKey(char start, char end) + { + this.start = start; + this.end = end; + } + + StartAndEndKey(String startAndEnd) + { + String[] pairs = startAndEnd.split(Constants.SEPARATOR_START_END_KEY); + this.start = pairs[0].charAt(0); + this.end = pairs[1].charAt(0); + } + + @Override + public String toString() + { + return "StartAndEndKey{" + + "start=" + start + + ", end=" + end + + '}'; + } + } } diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/utils/Constants.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/utils/Constants.java index 1f8d8635c..d6151f78d 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/utils/Constants.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/utils/Constants.java @@ -61,21 +61,26 @@ public class Constants */ public static final String SEPARATOR = ":"; + /** + * SEPARATOR_START_END_KEY + */ + public static final String SEPARATOR_START_END_KEY = "~"; + + /** + * START_END_ROW_KEYS_COUNT + */ + public static final int START_END_KEYS_COUNT = 100; + + /** + * ROWKEY_TAIL + */ + public static final String ROWKEY_TAIL = "|"; + /** * DEFAULT */ public static final String DEFAULT = "default"; - /** - * ARRAY - */ - public static final String ARRAY = "ARRYA [ "; - - /** - * apostrophe - */ - public static final String APOSTROPHE = "'"; - /** * 2 */ @@ -86,21 +91,6 @@ public class Constants */ public static final int NUMBER3 = 3; - /** - * 4 - */ - public static final int NUMBER4 = 4; - - /** - * 1024 - */ - public static final int NUMBER1024 = 1024; - - /** - * -1 - */ - public static final int NUMBER_NEGATIVE_1 = -1; - /** * "no such transaction: %s" */ @@ -114,22 +104,12 @@ public class Constants /** * SCAN_CACHING_SIZE */ - public static final int SCAN_CACHING_SIZE = 5000; + public static final int SCAN_CACHING_SIZE = 10000; /** - * constant string + * PUT_BATCH_SIZE */ - public static final String S_SCHEMA = "schema"; - - /** - * constant string - */ - public static final String S_TABLE = "table"; - - /** - * constant string - */ - public static final String S_COLUMNS = "columns"; + public static final int PUT_BATCH_SIZE = 10000; /** * constant string @@ -159,7 +139,7 @@ public class Constants /** * constant string */ - public static final String S_NAME = "name"; + public static final String S_SPLIT_BY_CHAR = "splitByChar"; /** * constant string @@ -171,21 +151,11 @@ public class Constants */ public static final String S_QUALIFIER = "qualifier"; - /** - * constant string - */ - public static final String S_TYPE = "type"; - /** * constant string */ public static final String S_ORDINAL = "ordinal"; - /** - * constant string - */ - public static final String S_COMMENT = "comment"; - /** * constant string */ @@ -196,11 +166,6 @@ public class Constants */ public static final String S_TRUE = "true"; - /** - * constant string - */ - public static final String S_FLASE = "false"; - /** * constant string */ diff --git a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/utils/Utils.java b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/utils/Utils.java index 244ada377..7104a1e7e 100644 --- a/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/utils/Utils.java +++ b/hetu-hbase/src/main/java/io/hetu/core/plugin/hbase/utils/Utils.java @@ -94,10 +94,7 @@ public class Utils } } } - catch (ClassNotFoundException e) { - LOG.error("createTypeByName failed... cause by : %s", e); - } - catch (IllegalAccessException e) { + catch (ClassNotFoundException | IllegalAccessException e) { LOG.error("createTypeByName failed... cause by : %s", e); } diff --git a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/TestHBase.java b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/TestHBase.java index 725607fa5..a50da4f3a 100644 --- a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/TestHBase.java +++ b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/TestHBase.java @@ -62,6 +62,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.OptionalLong; import static io.prestosql.spi.connector.ConnectorPageSink.NOT_BLOCKED; import static java.lang.String.format; @@ -379,7 +380,7 @@ public class TestHBase Map> ranges = new HashMap<>(); HBaseSplit split = new HBaseSplit( - "rowkey", TestUtils.createHBaseTableHandle(), hostAddressList, null, null, ranges, null, true); + "rowkey", TestUtils.createHBaseTableHandle(), hostAddressList, null, null, ranges, false); HBaseRecordSetProvider hrsp = new HBaseRecordSetProvider(hconn); RecordSet rs = @@ -406,7 +407,8 @@ public class TestHBase 0, hconn.getTable("hbase.test_table").getColumns(), hconn.getTable("hbase.test_table").getSerializerClassName(), - Optional.of("test_table")); + Optional.of("test_table"), + OptionalLong.empty()); if (insertHandler instanceof ConnectorInsertTableHandle) { ConnectorPageSink cps = hpsp.createPageSink( diff --git a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/TestQuery.java b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/TestQuery.java index c697bd86c..66191c351 100644 --- a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/TestQuery.java +++ b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/TestQuery.java @@ -44,6 +44,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.OptionalLong; import java.util.UUID; import static io.prestosql.spi.type.BigintType.BIGINT; @@ -90,8 +91,7 @@ public class TestQuery "startrow", "endrow", new HashMap<>(), - null, - true); + false); recordSet = new HBaseRecordSet( hconn, session, split, TestUtils.createHBaseTableHandle(), TestUtils.createColumnList()); @@ -123,7 +123,8 @@ public class TestQuery "", TestUtils.createTupleDomain(1), TestUtils.createColumnList(), - 0); + 0, + OptionalLong.empty()); HBaseSplit hBasesplit = new HBaseSplit( "rowKey", @@ -132,8 +133,7 @@ public class TestQuery "startrow", "endrow", new HashMap<>(), - null, - true); + false); HBaseRecordSet rSet = new HBaseRecordSet(hconn, session, hBasesplit, tableHandle, TestUtils.createColumnList()); rSet.cursor(); } @@ -157,7 +157,8 @@ public class TestQuery "", null, list, - 0); + 0, + OptionalLong.empty()); // case ABOVE Map> ranges = new HashMap<>(); @@ -168,7 +169,7 @@ public class TestQuery ranges.put(0, range); HBaseSplit hBasesplit = new HBaseSplit( - "rowkey", tableHandle, new ArrayList(1), "1", "12345678", ranges, null, true); + "rowkey", tableHandle, new ArrayList(1), "1", "12345678", ranges, false); HBaseRecordSet rSet = new HBaseRecordSet(hconn, session, hBasesplit, tableHandle, list); rSet.getFiltersFromDomains(ranges); diff --git a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/client/TestUtils.java b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/client/TestUtils.java index b09094d11..734fa19ae 100644 --- a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/client/TestUtils.java +++ b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/client/TestUtils.java @@ -35,6 +35,7 @@ import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.OptionalLong; import static io.hetu.core.plugin.hbase.utils.ValueSetUtils.createValueSet; import static io.prestosql.spi.type.VarcharType.VARCHAR; @@ -139,7 +140,8 @@ public class TestUtils false, "io.hetu.core.plugin.hbase.utils.serializers.StringRowSerializer", Optional.of("test_table"), - ""); + "", + OptionalLong.empty()); } /** @@ -154,7 +156,8 @@ public class TestUtils false, "StringRowSerializer", Optional.of(table), - ""); + "", + OptionalLong.empty()); } /** @@ -256,6 +259,7 @@ public class TestUtils external, Optional.of(serializerClassName), Optional.of(table), + Optional.of(""), Optional.of("")); hTableMetaMemory.put(key, hBaseTable); } diff --git a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/connector/TestHBaseConnector.java b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/connector/TestHBaseConnector.java index 2c7a0999f..6acaf5dae 100644 --- a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/connector/TestHBaseConnector.java +++ b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/connector/TestHBaseConnector.java @@ -427,7 +427,8 @@ public class TestHBaseConnector false, Optional.of("StringRowSerializer"), Optional.of(""), - Optional.of("table")); + Optional.of("table"), + Optional.of("")); // table is null try { hconn.dropTable(hBaseTable); @@ -492,7 +493,8 @@ public class TestHBaseConnector false, Optional.of("StringRowSerializer"), Optional.of(""), - Optional.of("test_table")); + Optional.of("test_table"), + Optional.of("")); try { method.invoke(hconn, hBaseTable, "rowkey"); } @@ -581,7 +583,7 @@ public class TestHBaseConnector hConnector.getPageSourceProvider(); hConnector.getPageSinkProvider(); hConnector.getPageSinkProvider(); - assertEquals(hConnector.getTableProperties().size(), 7); + assertEquals(hConnector.getTableProperties().size(), 8); assertEquals(hConnector.getColumnProperties().size(), 2); } } diff --git a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/connector/TestHBaseTableHandle.java b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/connector/TestHBaseTableHandle.java index 936a8f492..91a7a85aa 100644 --- a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/connector/TestHBaseTableHandle.java +++ b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/connector/TestHBaseTableHandle.java @@ -17,6 +17,7 @@ package io.hetu.core.plugin.hbase.connector; import org.testng.annotations.Test; import java.util.Optional; +import java.util.OptionalLong; import static org.testng.Assert.assertEquals; @@ -47,7 +48,8 @@ public class TestHBaseTableHandle false, serializerClassName, Optional.of(table), - ""); + "", + OptionalLong.empty()); assertEquals( "HBaseTableHandle{schema=hbase, table=table-1, rowId=rowkey, internal=false, " + "serializerClassName=io.hetu.core.plugin.hbase." @@ -88,7 +90,8 @@ public class TestHBaseTableHandle false, serializerClassName, Optional.of(table), - ""); + "", + OptionalLong.empty()); assertEquals(true, hBaseTableHandle.equals(hBaseTableHandle2)); } diff --git a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestHBaseTable.java b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestHBaseTable.java index 1bd43ca19..a037813c8 100644 --- a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestHBaseTable.java +++ b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestHBaseTable.java @@ -47,6 +47,7 @@ public class TestHBaseTable false, Optional.of("StringRowSerializer"), Optional.of("table"), + Optional.of(""), Optional.of("")); } diff --git a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestHetuHBaseMetastore.java b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestHetuHBaseMetastore.java index 197c69de5..11efb2d43 100644 --- a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestHetuHBaseMetastore.java +++ b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestHetuHBaseMetastore.java @@ -66,7 +66,8 @@ public class TestHetuHBaseMetastore false, serializerClassName, Optional.of(""), - Optional.of("testTable")); + Optional.of("testTable"), + Optional.of("")); metaStore.addHBaseTable(testHBaseTable1); @@ -79,7 +80,8 @@ public class TestHetuHBaseMetastore false, serializerClassName, Optional.of(""), - Optional.of("testTable2")); + Optional.of("testTable2"), + Optional.of("")); metaStore.addHBaseTable(testHBaseTable2); // test get all tables @@ -100,7 +102,8 @@ public class TestHetuHBaseMetastore false, serializerClassName, Optional.of(""), - Optional.of("testTable3")); + Optional.of("testTable3"), + Optional.of("")); metaStore.renameHBaseTable(testHBaseTable3, "hbase.testTable1"); hBaseTables = metaStore.getAllHBaseTables(); assertFalse(hBaseTables.containsKey("hbase.testTable1")); diff --git a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestingHetuMetastore.java b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestingHetuMetastore.java index 4f5142302..d1f70d331 100644 --- a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestingHetuMetastore.java +++ b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/metadata/TestingHetuMetastore.java @@ -61,7 +61,8 @@ public class TestingHetuMetastore false, Optional.of("io.hetu.core.plugin.hbase.utils.serializers.StringRowSerializer"), Optional.empty(), - Optional.of("hbase:test_table")); + Optional.of("hbase:test_table"), + Optional.empty()); private final TestingMySqlServer mySqlServer; private final HetuHBaseMetastore metaStore; diff --git a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/split/TestHBaseSplit.java b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/split/TestHBaseSplit.java index 56e4b661c..6d73b2e00 100644 --- a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/split/TestHBaseSplit.java +++ b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/split/TestHBaseSplit.java @@ -44,18 +44,14 @@ public class TestHBaseSplit "startrow", "endrow", new HashMap<>(), - null, true); assertEquals(0, split.getRanges().size()); assertEquals(TestUtils.createHBaseTableHandle(), split.getTableHandle()); - String className = "io.hetu.core.plugin.hbase.split.HBaseSplit"; - assertEquals(className, split.getInfo().toString().substring(0, className.length())); assertEquals(true, split.isRemotelyAccessible()); assertEquals("rowKey", split.getRowKeyName()); assertEquals("endrow", split.getEndRow()); assertEquals(true, split.isRandomSplit()); - assertEquals(null, split.getRegionInfo()); assertEquals(0, split.getAddresses().size()); } } diff --git a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/split/TestHbaseSplitManager.java b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/split/TestHbaseSplitManager.java index 0acd66a36..02524385c 100644 --- a/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/split/TestHbaseSplitManager.java +++ b/hetu-hbase/src/test/java/io/hetu/core/plugin/hbase/split/TestHbaseSplitManager.java @@ -24,7 +24,9 @@ import org.testng.annotations.AfterClass; import org.testng.annotations.BeforeClass; import org.testng.annotations.Test; +import java.util.NoSuchElementException; import java.util.Optional; +import java.util.OptionalLong; import static org.testng.Assert.assertEquals; @@ -70,10 +72,10 @@ public class TestHbaseSplitManager { try { hsm.getSplits(null, null, TestUtils.createHBaseTableHandle(), null); - throw new NullPointerException("testSplitManager : failed"); + throw new NoSuchElementException("No value present"); } - catch (NullPointerException e) { - assertEquals(e.toString(), "java.lang.NullPointerException: null pointer found when getting splits for scan"); + catch (NoSuchElementException e) { + assertEquals(e.toString(), "java.util.NoSuchElementException: No value present"); } } @@ -94,7 +96,8 @@ public class TestHbaseSplitManager "", TestUtils.createTupleDomain(1), TestUtils.createColumnList(), - 0); + 0, + OptionalLong.empty()); hsm.getSplits(null, null, tableHandle, null); }