Repository: carbondata
Updated Branches:
  refs/heads/master ad2c0e972 -> 7cd7623d1


[CARBONDATA-3160] Compaction support with MAP data type

This closes #2995


Project: http://git-wip-us.apache.org/repos/asf/carbondata/repo
Commit: http://git-wip-us.apache.org/repos/asf/carbondata/commit/7cd7623d
Tree: http://git-wip-us.apache.org/repos/asf/carbondata/tree/7cd7623d
Diff: http://git-wip-us.apache.org/repos/asf/carbondata/diff/7cd7623d

Branch: refs/heads/master
Commit: 7cd7623d19cd2eeba1831b8bdcf97327477e5c40
Parents: ad2c0e9
Author: dhatchayani <[email protected]>
Authored: Mon Dec 17 19:38:07 2018 +0530
Committer: ravipesala <[email protected]>
Committed: Tue Dec 18 12:15:43 2018 +0530

----------------------------------------------------------------------
 .../core/datastore/row/WriteStepRowUtil.java    |   4 +-
 .../complexType/TestCompactionComplexType.scala |  81 +++++++++++++-
 .../complexType/TestComplexDataType.scala       |   4 +-
 .../TestCreateDDLForComplexMapType.scala        |  69 ++++++++----
 .../spark/rdd/CarbonDataRDDFactory.scala        |  20 ++--
 .../CarbonAlterTableCompactionCommand.scala     |   5 -
 .../management/CarbonLoadDataCommand.scala      |  24 ++--
 .../store/CarbonFactDataHandlerModel.java       | 112 +++++--------------
 .../util/CarbonDataProcessorUtil.java           |   7 +-
 9 files changed, 179 insertions(+), 147 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/carbondata/blob/7cd7623d/core/src/main/java/org/apache/carbondata/core/datastore/row/WriteStepRowUtil.java
----------------------------------------------------------------------
diff --git 
a/core/src/main/java/org/apache/carbondata/core/datastore/row/WriteStepRowUtil.java
 
b/core/src/main/java/org/apache/carbondata/core/datastore/row/WriteStepRowUtil.java
index 49716ac..48772a0 100644
--- 
a/core/src/main/java/org/apache/carbondata/core/datastore/row/WriteStepRowUtil.java
+++ 
b/core/src/main/java/org/apache/carbondata/core/datastore/row/WriteStepRowUtil.java
@@ -85,10 +85,10 @@ public class WriteStepRowUtil {
 
     // For Complex Type Columns
     byte[][] complexKeys = ((ByteArrayWrapper) row[0]).getComplexTypesKeys();
-    for (int i = segmentProperties.getNumberOfNoDictionaryDimension();
+    for (int i = segmentProperties.getNumberOfNoDictionaryDimension(), j = 0;
          i < segmentProperties.getNumberOfNoDictionaryDimension() + 
segmentProperties
              .getComplexDimensions().size(); i++) {
-      noDictAndComplexKeys[i] = complexKeys[i];
+      noDictAndComplexKeys[i] = complexKeys[j++];
     }
 
     // no dictionary and complex dimension

http://git-wip-us.apache.org/repos/asf/carbondata/blob/7cd7623d/integration/spark-common-test/src/test/scala/org/apache/carbondata/integration/spark/testsuite/complexType/TestCompactionComplexType.scala
----------------------------------------------------------------------
diff --git 
a/integration/spark-common-test/src/test/scala/org/apache/carbondata/integration/spark/testsuite/complexType/TestCompactionComplexType.scala
 
b/integration/spark-common-test/src/test/scala/org/apache/carbondata/integration/spark/testsuite/complexType/TestCompactionComplexType.scala
index a353ec0..e00d4b6 100644
--- 
a/integration/spark-common-test/src/test/scala/org/apache/carbondata/integration/spark/testsuite/complexType/TestCompactionComplexType.scala
+++ 
b/integration/spark-common-test/src/test/scala/org/apache/carbondata/integration/spark/testsuite/complexType/TestCompactionComplexType.scala
@@ -23,11 +23,26 @@ import scala.collection.mutable
 
 import org.apache.spark.sql.Row
 import org.apache.spark.sql.test.util.QueryTest
+import org.scalatest.BeforeAndAfterAll
 
 import org.apache.carbondata.core.constants.CarbonCommonConstants
 import org.apache.carbondata.core.util.CarbonProperties
 
-class TestCompactionComplexType extends QueryTest {
+class TestCompactionComplexType extends QueryTest with BeforeAndAfterAll {
+
+  private val compactionThreshold = CarbonProperties.getInstance()
+    .getProperty(CarbonCommonConstants.COMPACTION_SEGMENT_LEVEL_THRESHOLD,
+      CarbonCommonConstants.DEFAULT_SEGMENT_LEVEL_THRESHOLD)
+
+  override protected def beforeAll(): Unit = {
+    CarbonProperties.getInstance()
+      .addProperty(CarbonCommonConstants.COMPACTION_SEGMENT_LEVEL_THRESHOLD, 
"2,3")
+  }
+
+  override protected def afterAll(): Unit = {
+    CarbonProperties.getInstance()
+      .addProperty(CarbonCommonConstants.COMPACTION_SEGMENT_LEVEL_THRESHOLD, 
compactionThreshold)
+  }
 
   test("test INT with struct and array, Encoding INT-->BYTE") {
     sql("Drop table if exists adaptive")
@@ -989,4 +1004,68 @@ class TestCompactionComplexType extends QueryTest {
       ))
   }
 
+  test("complex type compaction") {
+    sql("drop table if exists complexcarbontable")
+    sql("create table complexcarbontable(deviceInformationId int, channelsId 
string," +
+        "ROMSize string, purchasedate string, mobile struct<imei:string, 
imsi:string>," +
+        "MAC array<string>, locationinfo array<struct<ActiveAreaId:int, 
ActiveCountry:string, " +
+        "ActiveProvince:string, Activecity:string, ActiveDistrict:string, 
ActiveStreet:string>>," +
+        "proddate 
struct<productionDate:string,activeDeactivedate:array<string>>, gamePointId " +
+        "double,contractNumber double) " +
+        "STORED BY 'org.apache.carbondata.format' " +
+        "TBLPROPERTIES ('DICTIONARY_INCLUDE'='deviceInformationId')"
+    )
+    sql(
+      s"LOAD DATA local inpath '$resourcesPath/complexdata.csv' INTO table " +
+      "complexcarbontable " +
+      "OPTIONS('DELIMITER'=',', 'QUOTECHAR'='\"', 
'FILEHEADER'='deviceInformationId,channelsId," +
+      
"ROMSize,purchasedate,mobile,MAC,locationinfo,proddate,gamePointId,contractNumber',"
 +
+      "'COMPLEX_DELIMITER_LEVEL_1'='$', 'COMPLEX_DELIMITER_LEVEL_2'=':')"
+    )
+    sql(
+      s"LOAD DATA local inpath '$resourcesPath/complexdata.csv' INTO table " +
+      "complexcarbontable " +
+      "OPTIONS('DELIMITER'=',', 'QUOTECHAR'='\"', 
'FILEHEADER'='deviceInformationId,channelsId," +
+      
"ROMSize,purchasedate,mobile,MAC,locationinfo,proddate,gamePointId,contractNumber',"
 +
+      "'COMPLEX_DELIMITER_LEVEL_1'='$', 'COMPLEX_DELIMITER_LEVEL_2'=':')"
+    )
+    sql("alter table complexcarbontable compact 'minor'")
+    sql(
+      "select locationinfo,proddate from complexcarbontable where 
deviceInformationId=1 limit 1")
+      .show(false)
+    checkAnswer(sql(
+      "select locationinfo,proddate from complexcarbontable where 
deviceInformationId=1 limit 1"),
+      Seq(Row(mutable
+        .WrappedArray
+        .make(Array(Row(7, "Chinese", "Hubei Province", "yichang", "yichang", 
"yichang"),
+          Row(7, "India", "New Delhi", "delhi", "delhi", "delhi"))),
+        Row("29-11-2015", mutable
+          .WrappedArray.make(Array("29-11-2015", "29-11-2015"))))))
+    sql("drop table if exists complexcarbontable")
+  }
+
+  test("test minor compaction with all complex types") {
+    sql("Drop table if exists adaptive")
+    sql(
+      "create table adaptive(roll int, student 
struct<id:SHORT,name:string,marks:array<SHORT>>, " +
+      "mapField map<int, string>) " +
+      "stored by 'carbondata'")
+    sql("insert into adaptive 
values(1,'11111\001abc\001200\002300\002400','1\002Nalla\0012" +
+        "\002Singh\0013\002Gupta\0014\002Kumar')")
+    sql("insert into adaptive 
values(1,'11111\001abc\001200\002300\002401','11\002Nalla\00112" +
+        "\002Singh\00113\002Gupta\00114\002Kumar')")
+    sql("insert into adaptive 
values(1,'11111\001abc\001200\002300\002402','21\002Nalla\00122" +
+        "\002Singh\00123\002Gupta\00124\002Kumar')")
+    sql("insert into adaptive 
values(1,'11111\001abc\001200\002300\002403','31\002Nalla\00132" +
+        "\002Singh\00133\002Gupta\00134\002Kumar')")
+    sql("alter table adaptive compact 'minor' ")
+    checkAnswer(sql("select * from adaptive"),
+      Seq(Row(1, Row(11111, "abc", mutable.WrappedArray.make(Array(200, 300, 
400))), Map(1 -> "Nalla", 2 -> "Singh", 3 -> "Gupta", 4 -> "Kumar")),
+        Row(1, Row(11111, "abc", mutable.WrappedArray.make(Array(200, 300, 
401))), Map(11 -> "Nalla", 12 -> "Singh", 13 -> "Gupta", 14 -> "Kumar")),
+        Row(1, Row(11111, "abc", mutable.WrappedArray.make(Array(200, 300, 
402))), Map(21 -> "Nalla", 22 -> "Singh", 23 -> "Gupta", 24 -> "Kumar")),
+        Row(1, Row(11111, "abc", mutable.WrappedArray.make(Array(200, 300, 
403))), Map(31 -> "Nalla", 32 -> "Singh", 33 -> "Gupta", 34 -> "Kumar"))
+      ))
+    sql("Drop table if exists adaptive")
+  }
+
 }

http://git-wip-us.apache.org/repos/asf/carbondata/blob/7cd7623d/integration/spark-common-test/src/test/scala/org/apache/carbondata/integration/spark/testsuite/complexType/TestComplexDataType.scala
----------------------------------------------------------------------
diff --git 
a/integration/spark-common-test/src/test/scala/org/apache/carbondata/integration/spark/testsuite/complexType/TestComplexDataType.scala
 
b/integration/spark-common-test/src/test/scala/org/apache/carbondata/integration/spark/testsuite/complexType/TestComplexDataType.scala
index 9cbd842..ee61716 100644
--- 
a/integration/spark-common-test/src/test/scala/org/apache/carbondata/integration/spark/testsuite/complexType/TestComplexDataType.scala
+++ 
b/integration/spark-common-test/src/test/scala/org/apache/carbondata/integration/spark/testsuite/complexType/TestComplexDataType.scala
@@ -902,7 +902,7 @@ class TestComplexDataType extends QueryTest with 
BeforeAndAfterAll {
     checkExistence(sql("select * from table1"),true,"2.9E9")
   }
 
-  test("test block compaction - auto merge") {
+  test("test compaction - auto merge") {
     sql("DROP TABLE IF EXISTS table1")
     CarbonProperties.getInstance()
       .addProperty(CarbonCommonConstants.ENABLE_AUTO_LOAD_MERGE, "true")
@@ -929,7 +929,7 @@ class TestComplexDataType extends QueryTest with 
BeforeAndAfterAll {
       "/Struct.csv' into table table1 options('delimiter'=','," +
       
"'quotechar'='\"','fileheader'='roll,person','complex_delimiter_level_1'='$'," +
       "'complex_delimiter_level_2'='&')")
-    checkExistence(sql("show segments for table table1"),false, "Compacted")
+    checkExistence(sql("show segments for table table1"),true, "Compacted")
     CarbonProperties.getInstance()
       .addProperty(CarbonCommonConstants.ENABLE_AUTO_LOAD_MERGE, "false")
   }

http://git-wip-us.apache.org/repos/asf/carbondata/blob/7cd7623d/integration/spark-common-test/src/test/scala/org/apache/carbondata/spark/testsuite/createTable/TestCreateDDLForComplexMapType.scala
----------------------------------------------------------------------
diff --git 
a/integration/spark-common-test/src/test/scala/org/apache/carbondata/spark/testsuite/createTable/TestCreateDDLForComplexMapType.scala
 
b/integration/spark-common-test/src/test/scala/org/apache/carbondata/spark/testsuite/createTable/TestCreateDDLForComplexMapType.scala
index 941364c..b8f7549 100644
--- 
a/integration/spark-common-test/src/test/scala/org/apache/carbondata/spark/testsuite/createTable/TestCreateDDLForComplexMapType.scala
+++ 
b/integration/spark-common-test/src/test/scala/org/apache/carbondata/spark/testsuite/createTable/TestCreateDDLForComplexMapType.scala
@@ -284,26 +284,6 @@ class TestCreateDDLForComplexMapType extends QueryTest 
with BeforeAndAfterAll {
 
   }
 
-  test("Test Compaction blocking") {
-    sql("DROP TABLE IF EXISTS carbon")
-
-    sql(
-      s"""
-         | CREATE TABLE carbon(
-         | a INT,
-         | mapField map<INT,STRING>
-         | )
-         | STORED BY 'carbondata'
-         | """
-        .stripMargin)
-
-    val exception = intercept[UnsupportedOperationException](
-      sql("ALTER table carbon compact 'minor'")
-    )
-    assertResult("Compaction is unsupported for Table containing Map 
Columns")(exception
-      .getMessage())
-  }
-
   test("Test Load duplicate keys data in map") {
     sql("DROP TABLE IF EXISTS carbon")
     sql(
@@ -424,6 +404,55 @@ class TestCreateDDLForComplexMapType extends QueryTest 
with BeforeAndAfterAll {
     ))
   }
 
+  test("test compaction with map data type") {
+    sql("DROP TABLE IF EXISTS carbon")
+    sql(
+      s"""
+         | CREATE TABLE carbon(
+         | mapField map<INT,STRING>
+         | )
+         | STORED BY 'carbondata'
+         | """
+        .stripMargin)
+    sql(
+      s"""
+         | LOAD DATA LOCAL INPATH '$path'
+         | INTO TABLE carbon OPTIONS(
+         | 'header' = 'false')
+       """.stripMargin)
+    sql(
+      s"""
+         | LOAD DATA LOCAL INPATH '$path'
+         | INTO TABLE carbon OPTIONS(
+         | 'header' = 'false')
+       """.stripMargin)
+    sql(
+      s"""
+         | LOAD DATA LOCAL INPATH '$path'
+         | INTO TABLE carbon OPTIONS(
+         | 'header' = 'false')
+       """.stripMargin)
+    sql(
+      s"""
+         | LOAD DATA LOCAL INPATH '$path'
+         | INTO TABLE carbon OPTIONS(
+         | 'header' = 'false')
+       """.stripMargin)
+    sql("alter table carbon compact 'minor'")
+    sql("show segments for table carbon").show(false)
+    checkAnswer(sql("select * from carbon"), Seq(
+      Row(Map(1 -> "Nalla", 2 -> "Singh", 4 -> "Kumar")),
+      Row(Map(10 -> "Nallaa", 20 -> "Sissngh", 100 -> "Gusspta", 40 -> 
"Kumar")),
+      Row(Map(1 -> "Nalla", 2 -> "Singh", 4 -> "Kumar")),
+      Row(Map(10 -> "Nallaa", 20 -> "Sissngh", 100 -> "Gusspta", 40 -> 
"Kumar")),
+      Row(Map(1 -> "Nalla", 2 -> "Singh", 4 -> "Kumar")),
+      Row(Map(10 -> "Nallaa", 20 -> "Sissngh", 100 -> "Gusspta", 40 -> 
"Kumar")),
+      Row(Map(1 -> "Nalla", 2 -> "Singh", 4 -> "Kumar")),
+      Row(Map(10 -> "Nallaa", 20 -> "Sissngh", 100 -> "Gusspta", 40 -> 
"Kumar"))
+    ))
+    sql("DROP TABLE IF EXISTS carbon")
+  }
+
   test("Sort Column table property blocking for Map type") {
     sql("DROP TABLE IF EXISTS carbon")
     val exception1 = intercept[Exception] {

http://git-wip-us.apache.org/repos/asf/carbondata/blob/7cd7623d/integration/spark2/src/main/scala/org/apache/carbondata/spark/rdd/CarbonDataRDDFactory.scala
----------------------------------------------------------------------
diff --git 
a/integration/spark2/src/main/scala/org/apache/carbondata/spark/rdd/CarbonDataRDDFactory.scala
 
b/integration/spark2/src/main/scala/org/apache/carbondata/spark/rdd/CarbonDataRDDFactory.scala
index 34c6592..3b69f9e 100644
--- 
a/integration/spark2/src/main/scala/org/apache/carbondata/spark/rdd/CarbonDataRDDFactory.scala
+++ 
b/integration/spark2/src/main/scala/org/apache/carbondata/spark/rdd/CarbonDataRDDFactory.scala
@@ -571,19 +571,13 @@ object CarbonDataRDDFactory {
         if (carbonTable.isHivePartitionTable) {
           carbonLoadModel.setFactTimeStamp(System.currentTimeMillis())
         }
-        // Block compaction for table containing complex datatype
-        if (carbonTable.getTableInfo.getFactTable.getListOfColumns.asScala
-          .exists(m => m.getDataType.isComplexType)) {
-          LOGGER.warn("Compaction is skipped as table contains complex 
columns")
-        } else {
-          val compactedSegments = new util.ArrayList[String]()
-          handleSegmentMerging(sqlContext,
-            carbonLoadModel,
-            carbonTable,
-            compactedSegments,
-            operationContext)
-          carbonLoadModel.setMergedSegmentIds(compactedSegments)
-        }
+        val compactedSegments = new util.ArrayList[String]()
+        handleSegmentMerging(sqlContext,
+          carbonLoadModel,
+          carbonTable,
+          compactedSegments,
+          operationContext)
+        carbonLoadModel.setMergedSegmentIds(compactedSegments)
         writtenSegment
       } catch {
         case e: Exception =>

http://git-wip-us.apache.org/repos/asf/carbondata/blob/7cd7623d/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonAlterTableCompactionCommand.scala
----------------------------------------------------------------------
diff --git 
a/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonAlterTableCompactionCommand.scala
 
b/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonAlterTableCompactionCommand.scala
index a908a84..6defb0a 100644
--- 
a/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonAlterTableCompactionCommand.scala
+++ 
b/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonAlterTableCompactionCommand.scala
@@ -88,11 +88,6 @@ case class CarbonAlterTableCompactionCommand(
     if (!table.getTableInfo.isTransactionalTable) {
       throw new MalformedCarbonCommandException("Unsupported operation on non 
transactional table")
     }
-    if (table.getTableInfo.getFactTable.getListOfColumns.asScala
-      .exists((m => DataTypes.isMapType(m.getDataType)))) {
-      throw new UnsupportedOperationException(
-        "Compaction is unsupported for Table containing Map Columns")
-    }
     if (CarbonUtil.hasAggregationDataMap(table) ||
         (table.isChildDataMap && null == 
operationContext.getProperty(table.getTableName))) {
       // If the compaction request is of 'streaming' type then we need to 
generate loadCommands

http://git-wip-us.apache.org/repos/asf/carbondata/blob/7cd7623d/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonLoadDataCommand.scala
----------------------------------------------------------------------
diff --git 
a/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonLoadDataCommand.scala
 
b/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonLoadDataCommand.scala
index 3d2924c..cd4b0ae 100644
--- 
a/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonLoadDataCommand.scala
+++ 
b/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/management/CarbonLoadDataCommand.scala
@@ -831,21 +831,15 @@ case class CarbonLoadDataCommand(
     }
     try {
       carbonLoadModel.setFactTimeStamp(System.currentTimeMillis())
-      // Block compaction for table containing complex datatype
-      if (table.getTableInfo.getFactTable.getListOfColumns.asScala
-        .exists(m => m.getDataType.isComplexType)) {
-        LOGGER.warn("Compaction is skipped as table contains complex columns")
-      } else {
-        val compactedSegments = new util.ArrayList[String]()
-        // Trigger auto compaction
-        CarbonDataRDDFactory.handleSegmentMerging(
-          sparkSession.sqlContext,
-          carbonLoadModel,
-          table,
-          compactedSegments,
-          operationContext)
-        carbonLoadModel.setMergedSegmentIds(compactedSegments)
-      }
+      val compactedSegments = new util.ArrayList[String]()
+      // Trigger auto compaction
+      CarbonDataRDDFactory.handleSegmentMerging(
+        sparkSession.sqlContext,
+        carbonLoadModel,
+        table,
+        compactedSegments,
+        operationContext)
+      carbonLoadModel.setMergedSegmentIds(compactedSegments)
     } catch {
       case e: Exception =>
         throw new Exception(

http://git-wip-us.apache.org/repos/asf/carbondata/blob/7cd7623d/processing/src/main/java/org/apache/carbondata/processing/store/CarbonFactDataHandlerModel.java
----------------------------------------------------------------------
diff --git 
a/processing/src/main/java/org/apache/carbondata/processing/store/CarbonFactDataHandlerModel.java
 
b/processing/src/main/java/org/apache/carbondata/processing/store/CarbonFactDataHandlerModel.java
index b502da2..e759c02 100644
--- 
a/processing/src/main/java/org/apache/carbondata/processing/store/CarbonFactDataHandlerModel.java
+++ 
b/processing/src/main/java/org/apache/carbondata/processing/store/CarbonFactDataHandlerModel.java
@@ -44,10 +44,7 @@ import org.apache.carbondata.core.util.CarbonProperties;
 import org.apache.carbondata.core.util.CarbonUtil;
 import org.apache.carbondata.core.util.path.CarbonTablePath;
 import org.apache.carbondata.processing.datamap.DataMapWriterListener;
-import org.apache.carbondata.processing.datatypes.ArrayDataType;
 import org.apache.carbondata.processing.datatypes.GenericDataType;
-import org.apache.carbondata.processing.datatypes.PrimitiveDataType;
-import org.apache.carbondata.processing.datatypes.StructDataType;
 import org.apache.carbondata.processing.loading.CarbonDataLoadConfiguration;
 import org.apache.carbondata.processing.loading.DataField;
 import 
org.apache.carbondata.processing.loading.constants.DataLoadProcessorConstants;
@@ -246,24 +243,9 @@ public class CarbonFactDataHandlerModel {
     }
 
     //To Set MDKey Index of each primitive type in complex type
-    int surrIndex = simpleDimsCount;
-    Iterator<Map.Entry<String, GenericDataType>> complexMap =
-        
CarbonDataProcessorUtil.getComplexTypesMap(configuration.getDataFields(), 
configuration)
-            .entrySet().iterator();
-    Map<Integer, GenericDataType> complexIndexMap = new 
HashMap<>(complexDimensionCount);
-    while (complexMap.hasNext()) {
-      Map.Entry<String, GenericDataType> complexDataType = complexMap.next();
-      complexDataType.getValue().setOutputArrayIndex(0);
-      complexIndexMap.put(simpleDimsCount, complexDataType.getValue());
-      simpleDimsCount++;
-      List<GenericDataType> primitiveTypes = new ArrayList<GenericDataType>();
-      complexDataType.getValue().getAllPrimitiveChildren(primitiveTypes);
-      for (GenericDataType eachPrimitive : primitiveTypes) {
-        if (eachPrimitive.getIsColumnDictionary()) {
-          eachPrimitive.setSurrogateIndex(surrIndex++);
-        }
-      }
-    }
+    Map<Integer, GenericDataType> complexIndexMap = getComplexMap(
+        
configuration.getDataLoadProperty(DataLoadProcessorConstants.SERIALIZATION_NULL_FORMAT)
+            .toString(), simpleDimsCount, configuration.getDataFields());
 
     CarbonDataFileAttributes carbonDataFileAttributes =
         new CarbonDataFileAttributes(Long.parseLong(configuration.getTaskNo()),
@@ -375,7 +357,7 @@ public class CarbonFactDataHandlerModel {
     carbonFactDataHandlerModel.setColCardinality(formattedCardinality);
 
     carbonFactDataHandlerModel.setComplexIndexMap(
-        convertComplexDimensionToGenericDataType(segmentProperties,
+        convertComplexDimensionToComplexIndexMap(segmentProperties,
             loadModel.getSerializationNullFormat()));
     DataType[] measureDataTypes = new 
DataType[segmentProperties.getMeasures().size()];
     int i = 0;
@@ -417,75 +399,39 @@ public class CarbonFactDataHandlerModel {
    * @param isNullFormat
    * @return
    */
-  private static Map<Integer, GenericDataType> 
convertComplexDimensionToGenericDataType(
+  private static Map<Integer, GenericDataType> 
convertComplexDimensionToComplexIndexMap(
       SegmentProperties segmentProperties, String isNullFormat) {
     List<CarbonDimension> complexDimensions = 
segmentProperties.getComplexDimensions();
-    Map<Integer, GenericDataType> complexIndexMap = new 
HashMap<>(complexDimensions.size());
-    int dimensionCount = -1;
-    if (segmentProperties.getDimensions().size() == 0) {
-      dimensionCount = 0;
-    } else {
-      dimensionCount = segmentProperties.getDimensions().size() - 
segmentProperties
-          .getNumberOfNoDictionaryDimension() - 
segmentProperties.getComplexDimensions().size();
-    }
-    for (CarbonDimension carbonDimension : complexDimensions) {
-      if (carbonDimension.isComplex()) {
-        GenericDataType genericDataType;
-        DataType dataType = carbonDimension.getDataType();
-        if (DataTypes.isArrayType(dataType)) {
-          genericDataType =
-              new ArrayDataType(carbonDimension.getColName(), "", 
carbonDimension.getColumnId());
-        } else if (DataTypes.isStructType(dataType)) {
-          genericDataType =
-              new StructDataType(carbonDimension.getColName(), "", 
carbonDimension.getColumnId());
-        } else {
-          // Add Primitive type.
-          throw new RuntimeException("Primitive Type should not be coming in 
first loop");
-        }
-        if (carbonDimension.getNumberOfChild() > 0) {
-          addChildrenForComplex(carbonDimension.getListOfChildDimensions(), 
genericDataType,
-              isNullFormat);
-        }
-        genericDataType.setOutputArrayIndex(0);
-        complexIndexMap.put(dimensionCount++, genericDataType);
-      }
-
+    int simpleDimsCount = segmentProperties.getDimensions().size() - 
segmentProperties
+        .getNumberOfNoDictionaryDimension();
+    DataField[] dataFields = new DataField[complexDimensions.size()];
+    int i = 0;
+    for (CarbonColumn complexDimension : complexDimensions) {
+      dataFields[i++] = new DataField(complexDimension);
     }
-    return complexIndexMap;
+    return getComplexMap(isNullFormat, simpleDimsCount, dataFields);
   }
 
-  private static void addChildrenForComplex(List<CarbonDimension> 
listOfChildDimensions,
-      GenericDataType genericDataType, String isNullFormat) {
-    for (CarbonDimension carbonDimension : listOfChildDimensions) {
-      String parentColName =
-          carbonDimension.getColName().substring(0, 
carbonDimension.getColName().lastIndexOf("."));
-      DataType dataType = carbonDimension.getDataType();
-      if (DataTypes.isArrayType(dataType)) {
-        GenericDataType arrayGeneric =
-            new ArrayDataType(carbonDimension.getColName(), parentColName,
-                carbonDimension.getColumnId());
-        if (carbonDimension.getNumberOfChild() > 0) {
-          addChildrenForComplex(carbonDimension.getListOfChildDimensions(), 
arrayGeneric,
-              isNullFormat);
-        }
-        genericDataType.addChildren(arrayGeneric);
-      } else if (DataTypes.isStructType(dataType)) {
-        GenericDataType structGeneric =
-            new StructDataType(carbonDimension.getColName(), parentColName,
-                carbonDimension.getColumnId());
-        if (carbonDimension.getNumberOfChild() > 0) {
-          addChildrenForComplex(carbonDimension.getListOfChildDimensions(), 
structGeneric,
-              isNullFormat);
+  private static Map<Integer, GenericDataType> getComplexMap(String 
isNullFormat,
+      int simpleDimsCount, DataField[] dataFields) {
+    int surrIndex = simpleDimsCount;
+    Iterator<Map.Entry<String, GenericDataType>> complexMap =
+        CarbonDataProcessorUtil.getComplexTypesMap(dataFields, 
isNullFormat).entrySet().iterator();
+    Map<Integer, GenericDataType> complexIndexMap = new 
HashMap<>(dataFields.length);
+    while (complexMap.hasNext()) {
+      Map.Entry<String, GenericDataType> complexDataType = complexMap.next();
+      complexDataType.getValue().setOutputArrayIndex(0);
+      complexIndexMap.put(simpleDimsCount, complexDataType.getValue());
+      simpleDimsCount++;
+      List<GenericDataType> primitiveTypes = new ArrayList<GenericDataType>();
+      complexDataType.getValue().getAllPrimitiveChildren(primitiveTypes);
+      for (GenericDataType eachPrimitive : primitiveTypes) {
+        if (eachPrimitive.getIsColumnDictionary()) {
+          eachPrimitive.setSurrogateIndex(surrIndex++);
         }
-        genericDataType.addChildren(structGeneric);
-      } else {
-        // Primitive Data Type
-        genericDataType.addChildren(
-            new 
PrimitiveDataType(carbonDimension.getColumnSchema().getColumnName(),
-                dataType, parentColName, carbonDimension.getColumnId(),
-                
carbonDimension.getColumnSchema().hasEncoding(Encoding.DICTIONARY), 
isNullFormat));
       }
     }
+    return complexIndexMap;
   }
 
   /**

http://git-wip-us.apache.org/repos/asf/carbondata/blob/7cd7623d/processing/src/main/java/org/apache/carbondata/processing/util/CarbonDataProcessorUtil.java
----------------------------------------------------------------------
diff --git 
a/processing/src/main/java/org/apache/carbondata/processing/util/CarbonDataProcessorUtil.java
 
b/processing/src/main/java/org/apache/carbondata/processing/util/CarbonDataProcessorUtil.java
index 98b2543..1cff96d 100644
--- 
a/processing/src/main/java/org/apache/carbondata/processing/util/CarbonDataProcessorUtil.java
+++ 
b/processing/src/main/java/org/apache/carbondata/processing/util/CarbonDataProcessorUtil.java
@@ -311,17 +311,12 @@ public final class CarbonDataProcessorUtil {
 
   // TODO: need to simplify it. Not required create string first.
   public static Map<String, GenericDataType> getComplexTypesMap(DataField[] 
dataFields,
-      CarbonDataLoadConfiguration configuration) {
+      String nullFormat) {
     String complexTypeString = getComplexTypeString(dataFields);
 
     if (null == complexTypeString || complexTypeString.equals("")) {
       return new LinkedHashMap<>();
     }
-
-    String nullFormat =
-        
configuration.getDataLoadProperty(DataLoadProcessorConstants.SERIALIZATION_NULL_FORMAT)
-            .toString();
-
     Map<String, GenericDataType> complexTypesMap = new LinkedHashMap<String, 
GenericDataType>();
     String[] hierarchies = 
complexTypeString.split(CarbonCommonConstants.SEMICOLON_SPC_CHARACTER);
     for (int i = 0; i < hierarchies.length; i++) {

Reply via email to