Repository: incubator-blur Updated Branches: refs/heads/master d65f7eb29 -> e3854f06e
Adding support for multi values in blur-hive project and updating the testsuite to test for that as well as running a full hive driven data load. NOTE: non hadoop1 tests for blur-hive will be broken until testsuite is fixed for other profiles. Project: http://git-wip-us.apache.org/repos/asf/incubator-blur/repo Commit: http://git-wip-us.apache.org/repos/asf/incubator-blur/commit/e3854f06 Tree: http://git-wip-us.apache.org/repos/asf/incubator-blur/tree/e3854f06 Diff: http://git-wip-us.apache.org/repos/asf/incubator-blur/diff/e3854f06 Branch: refs/heads/master Commit: e3854f06e4293598eb910a8df439061303a69fbf Parents: d65f7eb Author: Aaron McCurry <[email protected]> Authored: Thu Jan 22 08:49:04 2015 -0500 Committer: Aaron McCurry <[email protected]> Committed: Thu Jan 22 08:49:04 2015 -0500 ---------------------------------------------------------------------- .../manager/writer/BlurIndexSimpleWriter.java | 7 +- blur-hive/pom.xml | 6 + .../blur/hive/BlurColumnNameResolver.java | 71 +++++ .../blur/hive/BlurObjectInspectorGenerator.java | 8 +- .../java/org/apache/blur/hive/BlurSerDe.java | 8 +- .../org/apache/blur/hive/BlurSerializer.java | 46 ++- .../org/apache/blur/hive/BlurSerDeTest.java | 291 +++++++++++++++++-- blur-hive/src/test/java/test.hive | 2 +- 8 files changed, 399 insertions(+), 40 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/e3854f06/blur-core/src/main/java/org/apache/blur/manager/writer/BlurIndexSimpleWriter.java ---------------------------------------------------------------------- diff --git a/blur-core/src/main/java/org/apache/blur/manager/writer/BlurIndexSimpleWriter.java b/blur-core/src/main/java/org/apache/blur/manager/writer/BlurIndexSimpleWriter.java index 7595663..12ad54b 100644 --- a/blur-core/src/main/java/org/apache/blur/manager/writer/BlurIndexSimpleWriter.java +++ b/blur-core/src/main/java/org/apache/blur/manager/writer/BlurIndexSimpleWriter.java @@ -605,8 +605,13 @@ public class BlurIndexSimpleWriter extends BlurIndex { final String table = _tableContext.getTable(); final String shard = _shardContext.getShard(); - LOG.info("Shard [{2}/{3}] Id [{0}] Finishing bulk mutate apply [{1}]", bulkId, apply, table, shard); + final BulkEntry bulkEntry = _bulkWriters.get(bulkId); + if (bulkEntry == null) { + LOG.info("Shard [{2}/{3}] Id [{0}] Nothing to apply.", bulkId, apply, table, shard); + return; + } + LOG.info("Shard [{2}/{3}] Id [{0}] Finishing bulk mutate apply [{1}]", bulkId, apply, table, shard); bulkEntry._writer.close(); Configuration configuration = _tableContext.getConfiguration(); http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/e3854f06/blur-hive/pom.xml ---------------------------------------------------------------------- diff --git a/blur-hive/pom.xml b/blur-hive/pom.xml index d78e7e4..b0c12e2 100644 --- a/blur-hive/pom.xml +++ b/blur-hive/pom.xml @@ -57,6 +57,12 @@ <scope>test</scope> </dependency> <dependency> + <groupId>org.apache.hive</groupId> + <artifactId>hive-jdbc</artifactId> + <version>${hive.version}</version> + <scope>test</scope> + </dependency> + <dependency> <groupId>org.apache.blur</groupId> <artifactId>blur-core</artifactId> <version>${project.version}</version> http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/e3854f06/blur-hive/src/main/java/org/apache/blur/hive/BlurColumnNameResolver.java ---------------------------------------------------------------------- diff --git a/blur-hive/src/main/java/org/apache/blur/hive/BlurColumnNameResolver.java b/blur-hive/src/main/java/org/apache/blur/hive/BlurColumnNameResolver.java new file mode 100644 index 0000000..7b6ac30 --- /dev/null +++ b/blur-hive/src/main/java/org/apache/blur/hive/BlurColumnNameResolver.java @@ -0,0 +1,71 @@ +/** + * 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.blur.hive; + +import java.util.Collection; +import java.util.HashMap; +import java.util.Map; + +import org.apache.blur.thrift.generated.ColumnDefinition; +import org.apache.blur.utils.BlurConstants; +import org.apache.hadoop.hive.serde2.SerDeException; + +public class BlurColumnNameResolver { + + private Map<String, String> _blurToHive = new HashMap<String, String>(); + private Map<String, String> _hiveToBlur = new HashMap<String, String>(); + + public BlurColumnNameResolver(Collection<ColumnDefinition> values) throws SerDeException { + for (ColumnDefinition columnDefinition : values) { + String blurName = columnDefinition.getColumnName(); + if (blurName.contains("-")) { + String derivedHiveName = blurName.replace("-", "_"); + addMapping(blurName, derivedHiveName); + } else { + addMapping(blurName); + } + } + addMapping(BlurConstants.ROW_ID); + addMapping(BlurConstants.RECORD_ID); + } + + private void addMapping(String blurName) throws SerDeException { + addMapping(blurName, blurName); + } + + private void addMapping(String blurName, String derivedHiveName) throws SerDeException { + if (_blurToHive.containsKey(derivedHiveName)) { + throw new SerDeException("Name collision while trying to map blur name [" + blurName + "] to hive name [" + + derivedHiveName + "]. Name already exists."); + } else if (_hiveToBlur.containsKey(derivedHiveName)) { + throw new SerDeException("Name collision while trying to map blur name [" + blurName + "] to hive name [" + + derivedHiveName + "]. Name already exists."); + } else { + _blurToHive.put(blurName, derivedHiveName); + _hiveToBlur.put(derivedHiveName, blurName); + } + } + + public String fromHiveToBlur(String hive) { + return _hiveToBlur.get(hive); + } + + public String fromBlurToHive(String blur) { + return _blurToHive.get(blur); + } + +} http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/e3854f06/blur-hive/src/main/java/org/apache/blur/hive/BlurObjectInspectorGenerator.java ---------------------------------------------------------------------- diff --git a/blur-hive/src/main/java/org/apache/blur/hive/BlurObjectInspectorGenerator.java b/blur-hive/src/main/java/org/apache/blur/hive/BlurObjectInspectorGenerator.java index 30cb9bf..7fb9f08 100644 --- a/blur-hive/src/main/java/org/apache/blur/hive/BlurObjectInspectorGenerator.java +++ b/blur-hive/src/main/java/org/apache/blur/hive/BlurObjectInspectorGenerator.java @@ -62,8 +62,11 @@ public class BlurObjectInspectorGenerator { private ObjectInspector _objectInspector; private List<String> _columnNames = new ArrayList<String>(); private List<TypeInfo> _columnTypes = new ArrayList<TypeInfo>(); + private BlurColumnNameResolver _columnNameResolver; - public BlurObjectInspectorGenerator(Collection<ColumnDefinition> colDefs) throws SerDeException { + public BlurObjectInspectorGenerator(Collection<ColumnDefinition> colDefs, BlurColumnNameResolver columnNameResolver) + throws SerDeException { + _columnNameResolver = columnNameResolver; List<ColumnDefinition> colDefList = new ArrayList<ColumnDefinition>(colDefs); Collections.sort(colDefList, COMPARATOR); @@ -74,7 +77,8 @@ public class BlurObjectInspectorGenerator { _columnTypes.add(TypeInfoFactory.stringTypeInfo); for (ColumnDefinition columnDefinition : colDefList) { - _columnNames.add(columnDefinition.getColumnName()); + String hiveColumnName = _columnNameResolver.fromBlurToHive(columnDefinition.getColumnName()); + _columnNames.add(hiveColumnName); _columnTypes.add(getTypeInfo(columnDefinition)); } _objectInspector = createObjectInspector(); http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/e3854f06/blur-hive/src/main/java/org/apache/blur/hive/BlurSerDe.java ---------------------------------------------------------------------- diff --git a/blur-hive/src/main/java/org/apache/blur/hive/BlurSerDe.java b/blur-hive/src/main/java/org/apache/blur/hive/BlurSerDe.java index ea1a675..055c6f4 100644 --- a/blur-hive/src/main/java/org/apache/blur/hive/BlurSerDe.java +++ b/blur-hive/src/main/java/org/apache/blur/hive/BlurSerDe.java @@ -53,6 +53,7 @@ public class BlurSerDe extends AbstractSerDe { private List<String> _columnNames; private List<TypeInfo> _columnTypes; private BlurSerializer _serializer; + private BlurColumnNameResolver _columnNameResolver; @Override public void initialize(Configuration conf, Properties tbl) throws SerDeException { @@ -104,12 +105,15 @@ public class BlurSerDe extends AbstractSerDe { } } - BlurObjectInspectorGenerator blurObjectInspectorGenerator = new BlurObjectInspectorGenerator(_schema.values()); + _columnNameResolver = new BlurColumnNameResolver(_schema.values()); + + BlurObjectInspectorGenerator blurObjectInspectorGenerator = new BlurObjectInspectorGenerator(_schema.values(), + _columnNameResolver); _objectInspector = blurObjectInspectorGenerator.getObjectInspector(); _columnNames = blurObjectInspectorGenerator.getColumnNames(); _columnTypes = blurObjectInspectorGenerator.getColumnTypes(); - _serializer = new BlurSerializer(_schema); + _serializer = new BlurSerializer(_schema, _columnNameResolver); } private void nullCheck(String name, String value) throws SerDeException { http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/e3854f06/blur-hive/src/main/java/org/apache/blur/hive/BlurSerializer.java ---------------------------------------------------------------------- diff --git a/blur-hive/src/main/java/org/apache/blur/hive/BlurSerializer.java b/blur-hive/src/main/java/org/apache/blur/hive/BlurSerializer.java index b7685b0..460c7f9 100644 --- a/blur-hive/src/main/java/org/apache/blur/hive/BlurSerializer.java +++ b/blur-hive/src/main/java/org/apache/blur/hive/BlurSerializer.java @@ -40,8 +40,10 @@ public class BlurSerializer { private static final String DATE_FORMAT = "dateFormat"; private static final String DATE = "date"; private Map<String, ThreadLocal<SimpleDateFormat>> _dateFormat = new HashMap<String, ThreadLocal<SimpleDateFormat>>(); + private BlurColumnNameResolver _columnNameResolver; - public BlurSerializer(Map<String, ColumnDefinition> colDefs) { + public BlurSerializer(Map<String, ColumnDefinition> colDefs, BlurColumnNameResolver columnNameResolver) { + _columnNameResolver = columnNameResolver; Set<Entry<String, ColumnDefinition>> entrySet = colDefs.entrySet(); for (Entry<String, ColumnDefinition> e : entrySet) { String columnName = e.getKey(); @@ -83,7 +85,7 @@ public class BlurSerializer { } for (int i = 0; i < size; i++) { - String columnName = columnNames.get(i); + String columnName = _columnNameResolver.fromHiveToBlur(columnNames.get(i)); StructField structFieldRef = outputFieldRefs.get(i); ObjectInspector fieldOI = structFieldRef.getFieldObjectInspector(); Object structFieldData = structFieldsDataAsList.get(i); @@ -99,8 +101,7 @@ public class BlurSerializer { } if (objectInspector instanceof PrimitiveObjectInspector) { PrimitiveObjectInspector primitiveObjectInspector = (PrimitiveObjectInspector) objectInspector; - Object primitiveJavaObject = primitiveObjectInspector.getPrimitiveJavaObject(data); - String strValue = toString(columnName, primitiveJavaObject); + String strValue = toString(columnName, data, primitiveObjectInspector); if (columnName.equals(BlurObjectInspectorGenerator.ROWID)) { blurRecord.setRowId(strValue); } else if (columnName.equals(BlurObjectInspectorGenerator.RECORDID)) { @@ -111,11 +112,11 @@ public class BlurSerializer { } else if (objectInspector instanceof StructObjectInspector) { StructObjectInspector structObjectInspector = (StructObjectInspector) objectInspector; Map<String, StructField> allStructFieldRefs = toMap(structObjectInspector.getAllStructFieldRefs()); - StructField latStructField = allStructFieldRefs.get(BlurObjectInspectorGenerator.LATITUDE); - StructField longStructField = allStructFieldRefs.get(BlurObjectInspectorGenerator.LONGITUDE); - Object latStructFieldData = structObjectInspector.getStructFieldData(data, latStructField); - Object longStructFieldData = structObjectInspector.getStructFieldData(data, longStructField); - blurRecord.addColumn(columnName, toLatLong(latStructFieldData, longStructFieldData)); + String latitude = getFieldData(columnName, data, structObjectInspector, allStructFieldRefs, + BlurObjectInspectorGenerator.LATITUDE); + String longitude = getFieldData(columnName, data, structObjectInspector, allStructFieldRefs, + BlurObjectInspectorGenerator.LONGITUDE); + blurRecord.addColumn(columnName, toLatLong(latitude, longitude)); } else if (objectInspector instanceof ListObjectInspector) { ListObjectInspector listObjectInspector = (ListObjectInspector) objectInspector; List<?> list = listObjectInspector.getList(data); @@ -129,9 +130,28 @@ public class BlurSerializer { } } - private String toLatLong(Object latStructFieldData, Object longStructFieldData) throws SerDeException { - return toString(BlurObjectInspectorGenerator.LATITUDE, latStructFieldData) + "," - + toString(BlurObjectInspectorGenerator.LONGITUDE, longStructFieldData); + private String getFieldData(String columnName, Object data, StructObjectInspector structObjectInspector, + Map<String, StructField> allStructFieldRefs, String name) throws SerDeException { + StructField structField = allStructFieldRefs.get(name); + ObjectInspector fieldObjectInspector = structField.getFieldObjectInspector(); + Object structFieldData = structObjectInspector.getStructFieldData(data, structField); + if (fieldObjectInspector instanceof PrimitiveObjectInspector) { + return toString(columnName, structFieldData, (PrimitiveObjectInspector) fieldObjectInspector); + } else { + throw new SerDeException("Embedded non-primitive type is not supported columnName [" + columnName + + "] objectInspector [" + fieldObjectInspector + "]."); + } + } + + private String toString(String columnName, Object data, PrimitiveObjectInspector primitiveObjectInspector) + throws SerDeException { + Object primitiveJavaObject = primitiveObjectInspector.getPrimitiveJavaObject(data); + return toString(columnName, primitiveJavaObject); + } + + private String toLatLong(String latitude, String longitude) throws SerDeException { + return toString(BlurObjectInspectorGenerator.LATITUDE, latitude) + "," + + toString(BlurObjectInspectorGenerator.LONGITUDE, longitude); } private Map<String, StructField> toMap(List<? extends StructField> allStructFieldRefs) { @@ -159,7 +179,7 @@ public class BlurSerializer { SimpleDateFormat simpleDateFormat = getSimpleDateFormat(columnName); return simpleDateFormat.format((Date) o); } else { - throw new SerDeException("Unknown type [" + o + "]"); + throw new SerDeException("Unknown type [" + o + "] with class [" + o == null ? "unknown" : o.getClass() + "]"); } } http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/e3854f06/blur-hive/src/test/java/org/apache/blur/hive/BlurSerDeTest.java ---------------------------------------------------------------------- diff --git a/blur-hive/src/test/java/org/apache/blur/hive/BlurSerDeTest.java b/blur-hive/src/test/java/org/apache/blur/hive/BlurSerDeTest.java index cf70bab..4c11089 100644 --- a/blur-hive/src/test/java/org/apache/blur/hive/BlurSerDeTest.java +++ b/blur-hive/src/test/java/org/apache/blur/hive/BlurSerDeTest.java @@ -16,11 +16,19 @@ */ package org.apache.blur.hive; -import static org.junit.Assert.*; +import static org.junit.Assert.assertEquals; import java.io.File; import java.io.IOException; +import java.io.PrintWriter; +import java.lang.reflect.Method; +import java.sql.Connection; import java.sql.Date; +import java.sql.DriverManager; +import java.sql.ResultSet; +import java.sql.ResultSetMetaData; +import java.sql.SQLException; +import java.sql.Statement; import java.text.SimpleDateFormat; import java.util.ArrayList; import java.util.HashMap; @@ -35,7 +43,10 @@ import org.apache.blur.thirdparty.thrift_0_9_0.TException; import org.apache.blur.thrift.BlurClient; import org.apache.blur.thrift.generated.Blur.Iface; import org.apache.blur.thrift.generated.BlurException; +import org.apache.blur.thrift.generated.BlurQuery; +import org.apache.blur.thrift.generated.BlurResults; import org.apache.blur.thrift.generated.ColumnDefinition; +import org.apache.blur.thrift.generated.Query; import org.apache.blur.thrift.generated.TableDescriptor; import org.apache.blur.utils.GCWatcher; import org.apache.hadoop.conf.Configuration; @@ -46,6 +57,9 @@ import org.apache.hadoop.fs.permission.FsAction; import org.apache.hadoop.fs.permission.FsPermission; import org.apache.hadoop.hive.serde2.SerDeException; import org.apache.hadoop.hive.serde2.objectinspector.ObjectInspector; +import org.apache.hadoop.mapred.MiniMRCluster; +import org.apache.hive.jdbc.HiveDriver; +import org.junit.After; import org.junit.AfterClass; import org.junit.Before; import org.junit.BeforeClass; @@ -53,14 +67,21 @@ import org.junit.Test; public class BlurSerDeTest { + private static final File METASTORE_DB_FILE = new File("metastore_db"); + private static final String FAM = "fam0"; private static final String YYYYMMDD = "yyyyMMdd"; + private static final String YYYY_MM_DD = "yyyy-MM-dd"; private static final String TEST = "test"; private static final File TMPDIR = new File(System.getProperty("blur.tmp.dir", "./target/tmp_BlurSerDeTest")); private static MiniCluster miniCluster; private static boolean externalProcesses = true; + private static final File WAREHOUSE = new File("./target/tmp/warehouse"); + private static final String COLUMN_SEP = new String(new char[] { 1 }); + private static final String ITEM_SEP = new String(new char[] { 2 }); @BeforeClass public static void startCluster() throws IOException { + System.setProperty("hadoop.log.dir", "./target/tmp_BlurSerDeTest_hadoop_log"); GCWatcher.init(0.60); LocalFileSystem localFS = FileSystem.getLocal(new Configuration()); File testDirectory = new File(TMPDIR, "blur-SerDe-test").getAbsoluteFile(); @@ -88,6 +109,8 @@ public class BlurSerDeTest { miniCluster.shutdownBlurCluster(); } + private Object mrMiniCluster; + @Before public void setup() throws BlurException, TException, IOException { String controllerConnectionStr = miniCluster.getControllerConnectionStr(); @@ -103,26 +126,45 @@ public class BlurSerDeTest { Map<String, String> props = new HashMap<String, String>(); props.put("dateFormat", YYYYMMDD); - client.addColumnDefinition(TEST, cd(false, "fam0", "string-col-single", "string")); - client.addColumnDefinition(TEST, cd(false, "fam0", "text-col-single", "text")); - client.addColumnDefinition(TEST, cd(false, "fam0", "stored-col-single", "stored")); - client.addColumnDefinition(TEST, cd(false, "fam0", "double-col-single", "double")); - client.addColumnDefinition(TEST, cd(false, "fam0", "float-col-single", "float")); - client.addColumnDefinition(TEST, cd(false, "fam0", "long-col-single", "long")); - client.addColumnDefinition(TEST, cd(false, "fam0", "int-col-single", "int")); - client.addColumnDefinition(TEST, cd(false, "fam0", "date-col-single", "date", props)); - - client.addColumnDefinition(TEST, cd(false, "fam0", "geo-col-single", "geo-pointvector")); - - client.addColumnDefinition(TEST, cd(true, "fam0", "string-col-multi", "string")); - client.addColumnDefinition(TEST, cd(true, "fam0", "text-col-multi", "text")); - client.addColumnDefinition(TEST, cd(true, "fam0", "stored-col-multi", "stored")); - client.addColumnDefinition(TEST, cd(true, "fam0", "double-col-multi", "double")); - client.addColumnDefinition(TEST, cd(true, "fam0", "float-col-multi", "float")); - client.addColumnDefinition(TEST, cd(true, "fam0", "long-col-multi", "long")); - client.addColumnDefinition(TEST, cd(true, "fam0", "int-col-multi", "int")); - client.addColumnDefinition(TEST, cd(true, "fam0", "date-col-multi", "date", props)); + client.addColumnDefinition(TEST, cd(false, FAM, "string-col-single", "string")); + client.addColumnDefinition(TEST, cd(false, FAM, "text-col-single", "text")); + client.addColumnDefinition(TEST, cd(false, FAM, "stored-col-single", "stored")); + client.addColumnDefinition(TEST, cd(false, FAM, "double-col-single", "double")); + client.addColumnDefinition(TEST, cd(false, FAM, "float-col-single", "float")); + client.addColumnDefinition(TEST, cd(false, FAM, "long-col-single", "long")); + client.addColumnDefinition(TEST, cd(false, FAM, "int-col-single", "int")); + client.addColumnDefinition(TEST, cd(false, FAM, "date-col-single", "date", props)); + + client.addColumnDefinition(TEST, cd(false, FAM, "geo-col-single", "geo-pointvector")); + + client.addColumnDefinition(TEST, cd(true, FAM, "string-col-multi", "string")); + client.addColumnDefinition(TEST, cd(true, FAM, "text-col-multi", "text")); + client.addColumnDefinition(TEST, cd(true, FAM, "stored-col-multi", "stored")); + client.addColumnDefinition(TEST, cd(true, FAM, "double-col-multi", "double")); + client.addColumnDefinition(TEST, cd(true, FAM, "float-col-multi", "float")); + client.addColumnDefinition(TEST, cd(true, FAM, "long-col-multi", "long")); + client.addColumnDefinition(TEST, cd(true, FAM, "int-col-multi", "int")); + client.addColumnDefinition(TEST, cd(true, FAM, "date-col-multi", "date", props)); + } + rmr(WAREHOUSE); + rmr(METASTORE_DB_FILE); + } + + @After + public void teardown() { + rmr(METASTORE_DB_FILE); + } + + private void rmr(File file) { + if (!file.exists()) { + return; + } + if (file.isDirectory()) { + for (File f : file.listFiles()) { + rmr(f); + } } + file.delete(); } private ColumnDefinition cd(boolean multiValue, String family, String columnName, String type) { @@ -146,7 +188,7 @@ public class BlurSerDeTest { Configuration conf = new Configuration(); Properties tbl = new Properties(); tbl.put(BlurSerDe.TABLE, TEST); - tbl.put(BlurSerDe.FAMILY, "fam0"); + tbl.put(BlurSerDe.FAMILY, FAM); tbl.put(BlurSerDe.ZK, miniCluster.getZkConnectionString()); blurSerDe.initialize(conf, tbl); @@ -207,6 +249,213 @@ public class BlurSerDeTest { assertEquals(list("1.0,2.0"), columns.get("geo-col-single")); } + @Test + public void test2() throws SQLException, ClassNotFoundException, IOException, BlurException, TException { + Class.forName(HiveDriver.class.getName()); + Connection connection = DriverManager.getConnection("jdbc:hive2://"); + + run(connection, "set hive.metastore.warehouse.dir=" + WAREHOUSE.toURI().toString()); + run(connection, "create database if not exists testdb"); + run(connection, "use testdb"); + + run(connection, "CREATE TABLE if not exists testtable ROW FORMAT SERDE 'org.apache.blur.hive.BlurSerDe' " + + "WITH SERDEPROPERTIES ( 'blur.zookeeper.connection'='" + miniCluster.getZkConnectionString() + "', " + + "'blur.table'='" + TEST + "', 'blur.family'='" + FAM + "' ) " + + "STORED BY 'org.apache.blur.hive.BlurHiveStorageHandler'"); + + run(connection, "desc testtable"); + + String createLoadTable = buildCreateLoadTable(connection); + run(connection, createLoadTable); + File dbDir = new File(WAREHOUSE, "testdb.db"); + File tableDir = new File(dbDir, "loadtable"); + int totalRecords = 100; + generateData(tableDir, totalRecords); + + run(connection, "select * from loadtable"); + + Configuration configuration = startMrMiniCluster(); + run(connection, "set mapred.job.tracker=" + configuration.get("mapred.job.tracker")); + run(connection, "insert into table testtable select * from loadtable"); + stopMrMiniCluster(); + connection.close(); + + Iface client = BlurClient.getClientFromZooKeeperConnectionStr(miniCluster.getZkConnectionString()); + BlurQuery blurQuery = new BlurQuery(); + Query query = new Query(); + query.setQuery("*"); + blurQuery.setQuery(query); + BlurResults results = client.query(TEST, blurQuery); + assertEquals(totalRecords, results.getTotalResults()); + } + + private void stopMrMiniCluster() { + callMethod(mrMiniCluster, "shutdown"); + } + + private Object callMethod(Object o, String methodName, Class<?>... classes) { + Class<? extends Object> clazz = o.getClass(); + try { + Method method = clazz.getDeclaredMethod(methodName, classes); + return method.invoke(o, new Object[] {}); + } catch (Exception e) { + throw new RuntimeException(e); + } + } + + private Configuration startMrMiniCluster() throws IOException { + mrMiniCluster = new MiniMRCluster(1, miniCluster.getFileSystemUri().toString(), 1); + return (Configuration) callMethod(mrMiniCluster, "createJobConf"); + } + + private void generateData(File file, int totalRecords) throws IOException { + SimpleDateFormat simpleDateFormat = new SimpleDateFormat(YYYY_MM_DD); + PrintWriter print = new PrintWriter(new File(file, "data")); + Date date = new Date(System.currentTimeMillis()); + for (int i = 0; i < totalRecords; i++) { + // rowid + print.print("rowid" + i); + print.print(COLUMN_SEP); + // recordid + print.print("recordid" + i); + print.print(COLUMN_SEP); + { + // date_col_multi + print.print(simpleDateFormat.format(date)); + print.print(ITEM_SEP); + print.print(simpleDateFormat.format(date)); + } + print.print(COLUMN_SEP); + // date_col_single + print.print(simpleDateFormat.format(date)); + print.print(COLUMN_SEP); + { + // double_col_multi + print.print("1.0"); + print.print(ITEM_SEP); + print.print("2.0"); + } + print.print(COLUMN_SEP); + // double_col_single + print.print("3.0"); + print.print(COLUMN_SEP); + + { + // float_col_multi + print.print("4.0"); + print.print(ITEM_SEP); + print.print("5.0"); + } + print.print(COLUMN_SEP); + // float_col_single + print.print("6.0"); + print.print(COLUMN_SEP); + + // geo_col_single + print.print("10.0"); + print.print(ITEM_SEP); + print.print("10.0"); + print.print(COLUMN_SEP); + + { + // int_col_multi + print.print("1"); + print.print(ITEM_SEP); + print.print("2"); + } + print.print(COLUMN_SEP); + // int_col_single + print.print("3"); + print.print(COLUMN_SEP); + + { + // long_col_multi + print.print("4"); + print.print(ITEM_SEP); + print.print("5"); + } + print.print(COLUMN_SEP); + // long_col_single + print.print("6"); + print.print(COLUMN_SEP); + + { + // stored_col_multi + print.print("stored_1"); + print.print(ITEM_SEP); + print.print("stored_2"); + } + print.print(COLUMN_SEP); + // stored_col_single + print.print("stored_3"); + print.print(COLUMN_SEP); + + { + // string_col_multi + print.print("string_1"); + print.print(ITEM_SEP); + print.print("string_2"); + } + print.print(COLUMN_SEP); + // string_col_single + print.print("string_3"); + print.print(COLUMN_SEP); + + { + // text_col_multi + print.print("text_1"); + print.print(ITEM_SEP); + print.print("text_2"); + } + print.print(COLUMN_SEP); + // text_col_single + print.print("text_3"); + print.println(); + } + print.close(); + + } + + private String buildCreateLoadTable(Connection connection) throws SQLException { + StringBuilder builder = new StringBuilder("create TABLE if not exists loadtable ("); + Statement statement = connection.createStatement(); + if (statement.execute("desc testtable")) { + ResultSet resultSet = statement.getResultSet(); + boolean first = true; + while (resultSet.next()) { + if (!first) { + builder.append(", "); + } + Object name = resultSet.getObject(1); + Object type = resultSet.getObject(2); + builder.append(name.toString()); + builder.append(' '); + builder.append(type.toString()); + first = false; + } + builder.append(")"); + return builder.toString(); + } + throw new RuntimeException("Can't build create table script."); + } + + private void run(Connection connection, String sql) throws SQLException { + System.out.println("Running:" + sql); + Statement statement = connection.createStatement(); + if (statement.execute(sql)) { + ResultSet resultSet = statement.getResultSet(); + while (resultSet.next()) { + ResultSetMetaData metaData = resultSet.getMetaData(); + int columnCount = metaData.getColumnCount(); + for (int i = 1; i <= columnCount; i++) { + System.out.print(resultSet.getObject(i) + "\t"); + } + System.out.println(); + } + } + statement.close(); + } + private List<String> list(String... sarray) { List<String> list = new ArrayList<String>(); for (String s : sarray) { http://git-wip-us.apache.org/repos/asf/incubator-blur/blob/e3854f06/blur-hive/src/test/java/test.hive ---------------------------------------------------------------------- diff --git a/blur-hive/src/test/java/test.hive b/blur-hive/src/test/java/test.hive index 5f950a8..5588904 100644 --- a/blur-hive/src/test/java/test.hive +++ b/blur-hive/src/test/java/test.hive @@ -49,7 +49,7 @@ create table if not exists input_data ( select * from input_data; -insert overwrite table test select * from input_data; +insert table test select * from input_data;
