Repository: carbondata Updated Branches: refs/heads/master 3a4b88138 -> faf74209a
[CARBONDATA-2585][CARBONDATA-2586]Fix local dictionary support for preagg and set localdict info in column schema This PR fixes local dictionary support for preaggregate and set the column dict info of each column in column schema read and write for backward compatibility. This closes #2451 Project: http://git-wip-us.apache.org/repos/asf/carbondata/repo Commit: http://git-wip-us.apache.org/repos/asf/carbondata/commit/faf74209 Tree: http://git-wip-us.apache.org/repos/asf/carbondata/tree/faf74209 Diff: http://git-wip-us.apache.org/repos/asf/carbondata/diff/faf74209 Branch: refs/heads/master Commit: faf74209a54aee2331a1cc967776d75aa49fc4c3 Parents: 3a4b881 Author: akashrn5 <[email protected]> Authored: Thu Jul 5 11:23:40 2018 +0530 Committer: kunal642 <[email protected]> Committed: Tue Jul 10 12:23:10 2018 +0530 ---------------------------------------------------------------------- .../core/metadata/schema/table/CarbonTable.java | 4 ++- .../schema/table/column/ColumnSchema.java | 2 ++ .../preaaggregate/PreAggregateTableHelper.scala | 38 +++++++++++++++----- 3 files changed, 35 insertions(+), 9 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/carbondata/blob/faf74209/core/src/main/java/org/apache/carbondata/core/metadata/schema/table/CarbonTable.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/metadata/schema/table/CarbonTable.java b/core/src/main/java/org/apache/carbondata/core/metadata/schema/table/CarbonTable.java index 80853f9..4cb42e6 100644 --- a/core/src/main/java/org/apache/carbondata/core/metadata/schema/table/CarbonTable.java +++ b/core/src/main/java/org/apache/carbondata/core/metadata/schema/table/CarbonTable.java @@ -378,7 +378,9 @@ public class CarbonTable implements Serializable { fillVisibleDimensions(tableSchema.getTableName()); fillVisibleMeasures(tableSchema.getTableName()); addImplicitDimension(dimensionOrdinal, implicitDimensions); - + CarbonUtil.setLocalDictColumnsToWrapperSchema(tableSchema.getListOfColumns(), + tableSchema.getTableProperties(), tableSchema.getTableProperties() + .get(CarbonCommonConstants.LOCAL_DICTIONARY_ENABLE)); dimensionOrdinalMax = dimensionOrdinal; } http://git-wip-us.apache.org/repos/asf/carbondata/blob/faf74209/core/src/main/java/org/apache/carbondata/core/metadata/schema/table/column/ColumnSchema.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/metadata/schema/table/column/ColumnSchema.java b/core/src/main/java/org/apache/carbondata/core/metadata/schema/table/column/ColumnSchema.java index 786e873..9bb9e85 100644 --- a/core/src/main/java/org/apache/carbondata/core/metadata/schema/table/column/ColumnSchema.java +++ b/core/src/main/java/org/apache/carbondata/core/metadata/schema/table/column/ColumnSchema.java @@ -553,6 +553,7 @@ public class ColumnSchema implements Serializable, Writable { parentColumnTableRelations.get(i).write(out); } } + out.writeBoolean(isLocalDictColumn); } @Override @@ -601,5 +602,6 @@ public class ColumnSchema implements Serializable, Writable { parentColumnTableRelations.add(parentColumnTableRelation); } } + this.isLocalDictColumn = in.readBoolean(); } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/faf74209/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/preaaggregate/PreAggregateTableHelper.scala ---------------------------------------------------------------------- diff --git a/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/preaaggregate/PreAggregateTableHelper.scala b/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/preaaggregate/PreAggregateTableHelper.scala index c37ba89..bc03be2 100644 --- a/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/preaaggregate/PreAggregateTableHelper.scala +++ b/integration/spark2/src/main/scala/org/apache/spark/sql/execution/command/preaaggregate/PreAggregateTableHelper.scala @@ -136,14 +136,36 @@ case class PreAggregateTableHelper( parentTable.getTableInfo.getFactTable.getTableProperties.asScala .getOrElse(CarbonCommonConstants.LOCAL_DICTIONARY_THRESHOLD, CarbonCommonConstants.LOCAL_DICTIONARY_THRESHOLD_DEFAULT)) - tableProperties - .put(CarbonCommonConstants.LOCAL_DICTIONARY_INCLUDE, - parentTable.getTableInfo.getFactTable.getTableProperties.asScala - .getOrElse(CarbonCommonConstants.LOCAL_DICTIONARY_INCLUDE, "")) - tableProperties - .put(CarbonCommonConstants.LOCAL_DICTIONARY_EXCLUDE, - parentTable.getTableInfo.getFactTable.getTableProperties.asScala - .getOrElse(CarbonCommonConstants.LOCAL_DICTIONARY_EXCLUDE, "")) + val parentDictInclude = parentTable.getTableInfo.getFactTable.getTableProperties.asScala + .getOrElse(CarbonCommonConstants.LOCAL_DICTIONARY_INCLUDE, "").split(",") + + val parentDictExclude = parentTable.getTableInfo.getFactTable.getTableProperties.asScala + .getOrElse(CarbonCommonConstants.LOCAL_DICTIONARY_EXCLUDE, "").split(",") + + val newDictInclude = + parentDictInclude.flatMap(parentcol => + fields.collect { + case col if fieldRelationMap(col).aggregateFunction.isEmpty && + parentcol.equals(fieldRelationMap(col). + columnTableRelationList.get.head.parentColumnName) => + col.column + }) + + val newDictExclude = parentDictExclude.flatMap(parentcol => + fields.collect { + case col if fieldRelationMap(col).aggregateFunction.isEmpty && + parentcol.equals(fieldRelationMap(col). + columnTableRelationList.get.head.parentColumnName) => + col.column + }) + if (newDictInclude.nonEmpty) { + tableProperties + .put(CarbonCommonConstants.LOCAL_DICTIONARY_INCLUDE, newDictInclude.mkString(",")) + } + if (newDictExclude.nonEmpty) { + tableProperties + .put(CarbonCommonConstants.LOCAL_DICTIONARY_EXCLUDE, newDictExclude.mkString(",")) + } val tableIdentifier = TableIdentifier(parentTable.getTableName + "_" + dataMapName, Some(parentTable.getDatabaseName))
