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++) {
