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;
 
 
 

Reply via email to