Repository: metamodel Updated Branches: refs/heads/master a2701136f -> 4cdc4e8c1
METAMODEL-161: Fixed Fixes #33 Project: http://git-wip-us.apache.org/repos/asf/metamodel/repo Commit: http://git-wip-us.apache.org/repos/asf/metamodel/commit/4cdc4e8c Tree: http://git-wip-us.apache.org/repos/asf/metamodel/tree/4cdc4e8c Diff: http://git-wip-us.apache.org/repos/asf/metamodel/diff/4cdc4e8c Branch: refs/heads/master Commit: 4cdc4e8c129c4143319fff4348aba35ce62123bd Parents: a270113 Author: Kasper Sørensen <[email protected]> Authored: Mon Jul 20 09:03:09 2015 +0200 Committer: Kasper Sørensen <[email protected]> Committed: Mon Jul 20 09:03:09 2015 +0200 ---------------------------------------------------------------------- CHANGES.md | 1 + hbase/pom.xml | 2 +- .../metamodel/hbase/HBaseConfiguration.java | 45 ++++++++++- .../metamodel/hbase/HBaseDataContext.java | 84 +++++++++++--------- .../apache/metamodel/hbase/HBaseDataSet.java | 5 +- .../org/apache/metamodel/hbase/HBaseTable.java | 22 +++-- .../metamodel/hbase/HBaseDataContextTest.java | 31 ++++---- pom.xml | 2 +- 8 files changed, 121 insertions(+), 71 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/metamodel/blob/4cdc4e8c/CHANGES.md ---------------------------------------------------------------------- diff --git a/CHANGES.md b/CHANGES.md index 6017760..9491a11 100644 --- a/CHANGES.md +++ b/CHANGES.md @@ -1,5 +1,6 @@ ### Apache MetaModel (work in progress) + * [METAMODEL-161] - Upgraded HBase client API to version 1.1.1 * [METAMODEL-160] - Added support for Apache Hive via the JDBC module of MetaModel. * [METAMODEL-162] - Made HdfsResource Serializable and added property getters http://git-wip-us.apache.org/repos/asf/metamodel/blob/4cdc4e8c/hbase/pom.xml ---------------------------------------------------------------------- diff --git a/hbase/pom.xml b/hbase/pom.xml index 5e0b697..2af0944 100644 --- a/hbase/pom.xml +++ b/hbase/pom.xml @@ -21,7 +21,7 @@ <name>MetaModel module for Apache HBase</name> <properties> - <hbase.version>1.0.0</hbase.version> + <hbase.version>1.1.1</hbase.version> </properties> <dependencies> http://git-wip-us.apache.org/repos/asf/metamodel/blob/4cdc4e8c/hbase/src/main/java/org/apache/metamodel/hbase/HBaseConfiguration.java ---------------------------------------------------------------------- diff --git a/hbase/src/main/java/org/apache/metamodel/hbase/HBaseConfiguration.java b/hbase/src/main/java/org/apache/metamodel/hbase/HBaseConfiguration.java index b3782d5..f2ec2ec 100644 --- a/hbase/src/main/java/org/apache/metamodel/hbase/HBaseConfiguration.java +++ b/hbase/src/main/java/org/apache/metamodel/hbase/HBaseConfiguration.java @@ -35,12 +35,18 @@ public class HBaseConfiguration extends BaseObject implements Serializable { public static final String DEFAULT_SCHEMA_NAME = "HBase"; public static final String DEFAULT_ZOOKEEPER_HOSTNAME = "127.0.0.1"; public static final int DEFAULT_ZOOKEEPER_PORT = 2181; + public static final int DEFAULT_HBASE_CLIENT_RETRIES = 1; + public static final int DEFAULT_ZOOKEEPER_SESSION_TIMEOUT = 5000; + public static final int DEFAULT_ZOOKEEPER_RECOVERY_RETRIES = 1; private final String _schemaName; private final int _zookeeperPort; private final String _zookeeperHostname; private final SimpleTableDef[] _tableDefinitions; private final ColumnType _defaultRowKeyType; + private final int _hbaseClientRetries; + private final int _zookeeperSessionTimeout; + private final int _zookeeperRecoveryRetries; /** * Creates a {@link HBaseConfiguration} using default values. @@ -52,7 +58,7 @@ public class HBaseConfiguration extends BaseObject implements Serializable { public HBaseConfiguration(String zookeeperHostname, int zookeeperPort) { this(DEFAULT_SCHEMA_NAME, zookeeperHostname, zookeeperPort, null, ColumnType.BINARY); } - + public HBaseConfiguration(String zookeeperHostname, int zookeeperPort, ColumnType defaultRowKeyType) { this(DEFAULT_SCHEMA_NAME, zookeeperHostname, zookeeperPort, null, defaultRowKeyType); } @@ -69,11 +75,34 @@ public class HBaseConfiguration extends BaseObject implements Serializable { */ public HBaseConfiguration(String schemaName, String zookeeperHostname, int zookeeperPort, SimpleTableDef[] tableDefinitions, ColumnType defaultRowKeyType) { + this(schemaName, zookeeperHostname, zookeeperPort, tableDefinitions, defaultRowKeyType, + DEFAULT_HBASE_CLIENT_RETRIES, DEFAULT_ZOOKEEPER_SESSION_TIMEOUT, DEFAULT_ZOOKEEPER_RECOVERY_RETRIES); + } + + /** + * Creates a {@link HBaseConfiguration} using detailed configuration + * properties. + * + * @param schemaName + * @param zookeeperHostname + * @param zookeeperPort + * @param tableDefinitions + * @param defaultRowKeyType + * @param hbaseClientRetries + * @param zookeeperSessionTimeout + * @param zookeeperRecoveryRetries + */ + public HBaseConfiguration(String schemaName, String zookeeperHostname, int zookeeperPort, + SimpleTableDef[] tableDefinitions, ColumnType defaultRowKeyType, int hbaseClientRetries, + int zookeeperSessionTimeout, int zookeeperRecoveryRetries) { _schemaName = schemaName; _zookeeperHostname = zookeeperHostname; _zookeeperPort = zookeeperPort; _tableDefinitions = tableDefinitions; _defaultRowKeyType = defaultRowKeyType; + _hbaseClientRetries = hbaseClientRetries; + _zookeeperSessionTimeout = zookeeperSessionTimeout; + _zookeeperRecoveryRetries = zookeeperRecoveryRetries; } public String getSchemaName() { @@ -91,7 +120,7 @@ public class HBaseConfiguration extends BaseObject implements Serializable { public SimpleTableDef[] getTableDefinitions() { return _tableDefinitions; } - + public ColumnType getDefaultRowKeyType() { return _defaultRowKeyType; } @@ -104,4 +133,16 @@ public class HBaseConfiguration extends BaseObject implements Serializable { list.add(_tableDefinitions); list.add(_defaultRowKeyType); } + + public int getHBaseClientRetries() { + return _hbaseClientRetries; + } + + public int getZookeeperSessionTimeout() { + return _zookeeperSessionTimeout; + } + + public int getZookeeperRecoveryRetries() { + return _zookeeperRecoveryRetries; + } } http://git-wip-us.apache.org/repos/asf/metamodel/blob/4cdc4e8c/hbase/src/main/java/org/apache/metamodel/hbase/HBaseDataContext.java ---------------------------------------------------------------------- diff --git a/hbase/src/main/java/org/apache/metamodel/hbase/HBaseDataContext.java b/hbase/src/main/java/org/apache/metamodel/hbase/HBaseDataContext.java index 6cc1bc1..e49076e 100644 --- a/hbase/src/main/java/org/apache/metamodel/hbase/HBaseDataContext.java +++ b/hbase/src/main/java/org/apache/metamodel/hbase/HBaseDataContext.java @@ -24,11 +24,11 @@ import java.util.List; import org.apache.hadoop.conf.Configuration; 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.ConnectionFactory; import org.apache.hadoop.hbase.client.Get; -import org.apache.hadoop.hbase.client.HBaseAdmin; -import org.apache.hadoop.hbase.client.HTableInterface; -import org.apache.hadoop.hbase.client.HTablePool; import org.apache.hadoop.hbase.client.Result; import org.apache.hadoop.hbase.client.ResultScanner; import org.apache.hadoop.hbase.client.Scan; @@ -61,8 +61,7 @@ public class HBaseDataContext extends QueryPostprocessDataContext { public static final String FIELD_ID = "_id"; private final HBaseConfiguration _configuration; - private final HBaseAdmin _admin; - private final HTablePool _tablePool; + private final Connection _connection; /** * Creates a {@link HBaseDataContext}. @@ -72,30 +71,24 @@ public class HBaseDataContext extends QueryPostprocessDataContext { public HBaseDataContext(HBaseConfiguration configuration) { Configuration config = createConfig(configuration); _configuration = configuration; - _admin = createHbaseAdmin(config); - _tablePool = new HTablePool(config, 100); + _connection = createConnection(config); } /** * Creates a {@link HBaseDataContext}. * * @param configuration - * @param admin - * @param hTablePool + * @param connection */ - public HBaseDataContext(HBaseConfiguration configuration, HBaseAdmin admin, HTablePool hTablePool) { + public HBaseDataContext(HBaseConfiguration configuration, Connection connection) { _configuration = configuration; - _tablePool = hTablePool; - _admin = admin; + _connection = connection; } - private HBaseAdmin createHbaseAdmin(Configuration config) { + private Connection createConnection(Configuration config) { try { - return new HBaseAdmin(config); - } catch (Exception e) { - if (e instanceof RuntimeException) { - throw (RuntimeException) e; - } + return ConnectionFactory.createConnection(config); + } catch (IOException e) { throw new MetaModelException(e); } } @@ -104,20 +97,27 @@ public class HBaseDataContext extends QueryPostprocessDataContext { Configuration config = org.apache.hadoop.hbase.HBaseConfiguration.create(); config.set("hbase.zookeeper.quorum", configuration.getZookeeperHostname()); config.set("hbase.zookeeper.property.clientPort", Integer.toString(configuration.getZookeeperPort())); + config.set("hbase.client.retries.number", Integer.toString(configuration.getHBaseClientRetries())); + config.set("zookeeper.session.timeout", Integer.toString(configuration.getZookeeperSessionTimeout())); + config.set("zookeeper.recovery.retry", Integer.toString(configuration.getZookeeperRecoveryRetries())); return config; } - public HTablePool getTablePool() { - return _tablePool; - } - /** - * Gets the HBaseAdmin used by this {@link DataContext} + * Gets the {@link Admin} used by this {@link DataContext} * * @return */ - public HBaseAdmin getHBaseAdmin() { - return _admin; + public Admin getAdmin() { + try { + return _connection.getAdmin(); + } catch (IOException e) { + throw new MetaModelException(e); + } + } + + public Connection getConnection() { + return _connection; } @Override @@ -127,7 +127,7 @@ public class HBaseDataContext extends QueryPostprocessDataContext { try { SimpleTableDef[] tableDefinitions = _configuration.getTableDefinitions(); if (tableDefinitions == null) { - final HTableDescriptor[] tables = _admin.listTables(); + final HTableDescriptor[] tables = getAdmin().listTables(); tableDefinitions = new SimpleTableDef[tables.length]; for (int i = 0; i < tables.length; i++) { SimpleTableDef emptyTableDef = new SimpleTableDef(tables[i].getNameAsString(), new String[0]); @@ -136,7 +136,7 @@ public class HBaseDataContext extends QueryPostprocessDataContext { } for (SimpleTableDef tableDef : tableDefinitions) { - schema.addTable(new HBaseTable(tableDef, schema, _admin, _configuration.getDefaultRowKeyType())); + schema.addTable(new HBaseTable(this, tableDef, schema, _configuration.getDefaultRowKeyType())); } return schema; @@ -166,7 +166,7 @@ public class HBaseDataContext extends QueryPostprocessDataContext { } long result = 0; - final HTableInterface hTable = _tablePool.getTable(table.getName()); + final org.apache.hadoop.hbase.client.Table hTable = getHTable(table.getName()); try { ResultScanner scanner = hTable.getScanner(new Scan()); try { @@ -182,17 +182,29 @@ public class HBaseDataContext extends QueryPostprocessDataContext { } } + protected org.apache.hadoop.hbase.client.Table getHTable(String name) { + try { + final TableName tableName = TableName.valueOf(name); + final org.apache.hadoop.hbase.client.Table hTable = _connection.getTable(tableName); + return hTable; + } catch (IOException e) { + throw new MetaModelException(e); + } + } + @Override - protected Row executePrimaryKeyLookupQuery(Table table, List<SelectItem> selectItems, Column primaryKeyColumn, Object keyValue) { - HTableInterface hTable = _tablePool.getTable(table.getName()); - Get get = new Get(ByteUtils.toBytes(keyValue)); + protected Row executePrimaryKeyLookupQuery(Table table, List<SelectItem> selectItems, Column primaryKeyColumn, + Object keyValue) { + final org.apache.hadoop.hbase.client.Table hTable = getHTable(table.getName()); + final Get get = new Get(ByteUtils.toBytes(keyValue)); try { - Result result = hTable.get(get); - DataSetHeader header = new SimpleDataSetHeader(selectItems); - Row row = new HBaseRow(header, result); + final Result result = hTable.get(get); + final DataSetHeader header = new SimpleDataSetHeader(selectItems); + final Row row = new HBaseRow(header, result); return row; } catch (IOException e) { - throw new IllegalStateException("Failed to execute HBase get operation with " + primaryKeyColumn.getName() + " = " + keyValue, e); + throw new IllegalStateException("Failed to execute HBase get operation with " + primaryKeyColumn.getName() + + " = " + keyValue, e); } finally { FileHelper.safeClose(hTable); } @@ -217,7 +229,7 @@ public class HBaseDataContext extends QueryPostprocessDataContext { setMaxRows(scan, maxRows); } - final HTableInterface hTable = _tablePool.getTable(table.getName()); + final org.apache.hadoop.hbase.client.Table hTable = getHTable(table.getName()); try { final ResultScanner scanner = hTable.getScanner(scan); return new HBaseDataSet(columns, scanner, hTable); http://git-wip-us.apache.org/repos/asf/metamodel/blob/4cdc4e8c/hbase/src/main/java/org/apache/metamodel/hbase/HBaseDataSet.java ---------------------------------------------------------------------- diff --git a/hbase/src/main/java/org/apache/metamodel/hbase/HBaseDataSet.java b/hbase/src/main/java/org/apache/metamodel/hbase/HBaseDataSet.java index 699774d..a4158d5 100644 --- a/hbase/src/main/java/org/apache/metamodel/hbase/HBaseDataSet.java +++ b/hbase/src/main/java/org/apache/metamodel/hbase/HBaseDataSet.java @@ -20,7 +20,6 @@ package org.apache.metamodel.hbase; import java.io.IOException; -import org.apache.hadoop.hbase.client.HTableInterface; import org.apache.hadoop.hbase.client.Result; import org.apache.hadoop.hbase.client.ResultScanner; import org.apache.metamodel.MetaModelException; @@ -35,10 +34,10 @@ final class HBaseDataSet extends AbstractDataSet { private static final Logger logger = LoggerFactory.getLogger(HBaseDataSet.class); private final ResultScanner _scanner; - private final HTableInterface _hTable; + private final org.apache.hadoop.hbase.client.Table _hTable; private volatile Result _nextResult; - public HBaseDataSet(Column[] columns, ResultScanner scanner, HTableInterface hTable) { + public HBaseDataSet(Column[] columns, ResultScanner scanner, org.apache.hadoop.hbase.client.Table hTable) { super(columns); _scanner = scanner; _hTable = hTable; http://git-wip-us.apache.org/repos/asf/metamodel/blob/4cdc4e8c/hbase/src/main/java/org/apache/metamodel/hbase/HBaseTable.java ---------------------------------------------------------------------- diff --git a/hbase/src/main/java/org/apache/metamodel/hbase/HBaseTable.java b/hbase/src/main/java/org/apache/metamodel/hbase/HBaseTable.java index b79acac..03c3263 100644 --- a/hbase/src/main/java/org/apache/metamodel/hbase/HBaseTable.java +++ b/hbase/src/main/java/org/apache/metamodel/hbase/HBaseTable.java @@ -21,8 +21,6 @@ package org.apache.metamodel.hbase; import java.util.List; import org.apache.hadoop.hbase.HColumnDescriptor; -import org.apache.hadoop.hbase.HTableDescriptor; -import org.apache.hadoop.hbase.client.HBaseAdmin; import org.apache.metamodel.MetaModelException; import org.apache.metamodel.schema.Column; import org.apache.metamodel.schema.ColumnType; @@ -42,13 +40,13 @@ final class HBaseTable extends MutableTable { private static final long serialVersionUID = 1L; private static final Logger logger = LoggerFactory.getLogger(HBaseTable.class); - private final transient HBaseAdmin _admin; + private final transient HBaseDataContext _dataContext; private final transient ColumnType _defaultRowKeyColumnType; - public HBaseTable(SimpleTableDef tableDef, MutableSchema schema, HBaseAdmin admin, + public HBaseTable(HBaseDataContext dataContext, SimpleTableDef tableDef, MutableSchema schema, ColumnType defaultRowKeyColumnType) { super(tableDef.getName(), TableType.TABLE, schema); - _admin = admin; + _dataContext = dataContext; _defaultRowKeyColumnType = defaultRowKeyColumnType; final String[] columnNames = tableDef.getColumnNames(); @@ -72,8 +70,8 @@ final class HBaseTable extends MutableTable { if (columnNumber == 1) { // insert a default definition of the id column - final MutableColumn idColumn = new MutableColumn(HBaseDataContext.FIELD_ID, - defaultRowKeyColumnType).setPrimaryKey(true).setColumnNumber(columnNumber).setTable(this); + final MutableColumn idColumn = new MutableColumn(HBaseDataContext.FIELD_ID, defaultRowKeyColumnType) + .setPrimaryKey(true).setColumnNumber(columnNumber).setTable(this); addColumn(idColumn); columnNumber++; } @@ -96,19 +94,19 @@ final class HBaseTable extends MutableTable { @Override protected List<Column> getColumnsInternal() { final List<Column> columnsInternal = super.getColumnsInternal(); - if (columnsInternal.isEmpty() && _admin != null) { + if (columnsInternal.isEmpty() && _dataContext != null) { try { - HTableDescriptor tableDescriptor = _admin.getTableDescriptor(getName().getBytes()); + final org.apache.hadoop.hbase.client.Table table = _dataContext.getHTable(getName()); int columnNumber = 1; - final MutableColumn idColumn = new MutableColumn(HBaseDataContext.FIELD_ID, - _defaultRowKeyColumnType).setPrimaryKey(true).setColumnNumber(columnNumber).setTable(this); + final MutableColumn idColumn = new MutableColumn(HBaseDataContext.FIELD_ID, _defaultRowKeyColumnType) + .setPrimaryKey(true).setColumnNumber(columnNumber).setTable(this); addColumn(idColumn); columnNumber++; // What about timestamp? - final HColumnDescriptor[] columnFamilies = tableDescriptor.getColumnFamilies(); + final HColumnDescriptor[] columnFamilies = table.getTableDescriptor().getColumnFamilies(); for (int i = 0; i < columnFamilies.length; i++) { final HColumnDescriptor columnDescriptor = columnFamilies[i]; final String columnFamilyName = columnDescriptor.getNameAsString(); http://git-wip-us.apache.org/repos/asf/metamodel/blob/4cdc4e8c/hbase/src/test/java/org/apache/metamodel/hbase/HBaseDataContextTest.java ---------------------------------------------------------------------- diff --git a/hbase/src/test/java/org/apache/metamodel/hbase/HBaseDataContextTest.java b/hbase/src/test/java/org/apache/metamodel/hbase/HBaseDataContextTest.java index 1441b7d..eb0ba48 100644 --- a/hbase/src/test/java/org/apache/metamodel/hbase/HBaseDataContextTest.java +++ b/hbase/src/test/java/org/apache/metamodel/hbase/HBaseDataContextTest.java @@ -22,9 +22,8 @@ import java.util.Arrays; import org.apache.hadoop.hbase.HColumnDescriptor; import org.apache.hadoop.hbase.HTableDescriptor; -import org.apache.hadoop.hbase.client.HBaseAdmin; -import org.apache.hadoop.hbase.client.HTableInterface; -import org.apache.hadoop.hbase.client.HTablePool; +import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.client.Admin; import org.apache.hadoop.hbase.client.Put; import org.apache.metamodel.data.DataSet; import org.apache.metamodel.schema.ColumnType; @@ -140,37 +139,37 @@ public class HBaseDataContextTest extends HBaseTestCase { } private void insertRecordsNatively() throws Exception { - final HTablePool tablePool = _dataContext.getTablePool(); - final HTableInterface hTable = tablePool.getTable(EXAMPLE_TABLE_NAME); + final org.apache.hadoop.hbase.client.Table hTable = _dataContext.getHTable(EXAMPLE_TABLE_NAME); try { final Put put1 = new Put("junit1".getBytes()); - put1.add("foo".getBytes(), "hello".getBytes(), "world".getBytes()); - put1.add("bar".getBytes(), "hi".getBytes(), "there".getBytes()); - put1.add("bar".getBytes(), "hey".getBytes(), "yo".getBytes()); + put1.addColumn("foo".getBytes(), "hello".getBytes(), "world".getBytes()); + put1.addColumn("bar".getBytes(), "hi".getBytes(), "there".getBytes()); + put1.addColumn("bar".getBytes(), "hey".getBytes(), "yo".getBytes()); final Put put2 = new Put("junit2".getBytes()); - put2.add("bar".getBytes(), "bah".getBytes(), new byte[] { 1, 2, 3 }); - put2.add("bar".getBytes(), "hi".getBytes(), "you".getBytes()); + put2.addColumn("bar".getBytes(), "bah".getBytes(), new byte[] { 1, 2, 3 }); + put2.addColumn("bar".getBytes(), "hi".getBytes(), "you".getBytes()); - hTable.batch(Arrays.asList(put1, put2)); + final Object[] result = new Object[2]; + hTable.batch(Arrays.asList(put1, put2), result); } finally { hTable.close(); - tablePool.closeTablePool(EXAMPLE_TABLE_NAME); - tablePool.close(); } } private void createTableNatively() throws Exception { + final TableName tableName = TableName.valueOf(EXAMPLE_TABLE_NAME); + // check if the table exists - if (_dataContext.getHBaseAdmin().isTableAvailable(EXAMPLE_TABLE_NAME)) { + if (_dataContext.getAdmin().isTableAvailable(tableName)) { System.out.println("Unittest table already exists: " + EXAMPLE_TABLE_NAME); // table already exists return; } - HBaseAdmin admin = _dataContext.getHBaseAdmin(); + Admin admin = _dataContext.getAdmin(); System.out.println("Creating table"); - final HTableDescriptor tableDescriptor = new HTableDescriptor(EXAMPLE_TABLE_NAME.getBytes()); + final HTableDescriptor tableDescriptor = new HTableDescriptor(tableName); tableDescriptor.addFamily(new HColumnDescriptor("foo".getBytes())); tableDescriptor.addFamily(new HColumnDescriptor("bar".getBytes())); admin.createTable(tableDescriptor); http://git-wip-us.apache.org/repos/asf/metamodel/blob/4cdc4e8c/pom.xml ---------------------------------------------------------------------- diff --git a/pom.xml b/pom.xml index 80b6e12..31e9ada 100644 --- a/pom.xml +++ b/pom.xml @@ -26,7 +26,7 @@ under the License. <javadoc.version>2.9.1</javadoc.version> <slf4j.version>1.7.7</slf4j.version> <junit.version>4.11</junit.version> - <guava.version>16.0</guava.version> + <guava.version>16.0.1</guava.version> <hadoop.version>2.6.0</hadoop.version> <easymock.version>3.2</easymock.version> <httpcomponents.version>4.3.1</httpcomponents.version>
