Repository: carbondata Updated Branches: refs/heads/master 5804d7570 -> e7103397d
http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/core/src/main/java/org/apache/carbondata/core/localdictionary/dictionaryholder/DictionaryStore.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/localdictionary/dictionaryholder/DictionaryStore.java b/core/src/main/java/org/apache/carbondata/core/localdictionary/dictionaryholder/DictionaryStore.java new file mode 100644 index 0000000..226104b --- /dev/null +++ b/core/src/main/java/org/apache/carbondata/core/localdictionary/dictionaryholder/DictionaryStore.java @@ -0,0 +1,50 @@ +/* + * 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.carbondata.core.localdictionary.dictionaryholder; + +import org.apache.carbondata.core.localdictionary.exception.DictionaryThresholdReachedException; + +/** + * Interface for storing the dictionary key and value. + * Concrete implementation can be of map based or trie based. + */ +public interface DictionaryStore { + + /** + * Below method will be used to add dictionary value to dictionary holder + * if it is already present in the holder then it will return exiting dictionary value. + * @param key + * dictionary key + * @return dictionary value + */ + int putIfAbsent(byte[] key) throws DictionaryThresholdReachedException; + + /** + * Below method to get the current size of dictionary + * @return true if threshold of store reached + */ + boolean isThresholdReached(); + + /** + * Below method will be used to get the dictionary key based on value + * @param value + * dictionary value + * @return dictionary key based on value + */ + byte[] getDictionaryKeyBasedOnValue(int value); + +} http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/core/src/main/java/org/apache/carbondata/core/localdictionary/dictionaryholder/MapBasedDictionaryStore.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/localdictionary/dictionaryholder/MapBasedDictionaryStore.java b/core/src/main/java/org/apache/carbondata/core/localdictionary/dictionaryholder/MapBasedDictionaryStore.java new file mode 100644 index 0000000..05ca002 --- /dev/null +++ b/core/src/main/java/org/apache/carbondata/core/localdictionary/dictionaryholder/MapBasedDictionaryStore.java @@ -0,0 +1,137 @@ +/* + * 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.carbondata.core.localdictionary.dictionaryholder; + +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +import org.apache.carbondata.core.cache.dictionary.DictionaryByteArrayWrapper; +import org.apache.carbondata.core.localdictionary.exception.DictionaryThresholdReachedException; + +/** + * Map based dictionary holder class, it will use map to hold + * the dictionary key and its value + */ +public class MapBasedDictionaryStore implements DictionaryStore { + + /** + * use to assign dictionary value to new key + */ + private int lastAssignValue; + + /** + * to maintain dictionary key value + */ + private final Map<DictionaryByteArrayWrapper, Integer> dictionary; + + /** + * maintaining array for reverse lookup + * otherwise iterating everytime in map for reverse lookup will be slowdown the performance + * It will only maintain the reference + */ + private DictionaryByteArrayWrapper[] referenceDictionaryArray; + + /** + * dictionary threshold to check if threshold is reached + */ + private int dictionaryThreshold; + + /** + * for checking threshold is reached or not + */ + private boolean isThresholdReached; + + public MapBasedDictionaryStore(int dictionaryThreshold) { + this.dictionaryThreshold = dictionaryThreshold; + this.dictionary = new ConcurrentHashMap<>(); + this.referenceDictionaryArray = new DictionaryByteArrayWrapper[dictionaryThreshold]; + } + + /** + * Below method will be used to add dictionary value to dictionary holder + * if it is already present in the holder then it will return exiting dictionary value. + * + * @param data dictionary key + * @return dictionary value + */ + @Override public int putIfAbsent(byte[] data) throws DictionaryThresholdReachedException { + // check if threshold has already reached + checkIfThresholdReached(); + DictionaryByteArrayWrapper key = new DictionaryByteArrayWrapper(data); + // get the dictionary value + Integer value = dictionary.get(key); + // if value is null then dictionary is not present in store + if (null == value) { + // aquire the lock + synchronized (dictionary) { + // check threshold + checkIfThresholdReached(); + // get the value again as other thread might have added + value = dictionary.get(key); + // double chekcing + if (null == value) { + // increment the value + value = ++lastAssignValue; + // if new value is greater than threshold + if (value > dictionaryThreshold) { + // clear the dictionary + dictionary.clear(); + referenceDictionaryArray = null; + // set the threshold boolean to true + isThresholdReached = true; + // throw exception + checkIfThresholdReached(); + } + // add to reference array + // position is -1 as dictionary value starts from 1 + this.referenceDictionaryArray[value - 1] = key; + dictionary.put(key, value); + } + } + } + return value; + } + + private void checkIfThresholdReached() throws DictionaryThresholdReachedException { + if (isThresholdReached) { + throw new DictionaryThresholdReachedException( + "Unable to generate dictionary value. Dictionary threshold reached"); + } + } + + /** + * Below method to get the current size of dictionary + * + * @return + */ + @Override public boolean isThresholdReached() { + return isThresholdReached; + } + + /** + * Below method will be used to get the dictionary key based on value + * + * @param value dictionary value + * Caller will take of passing proper value + * @return dictionary key based on value + */ + @Override public byte[] getDictionaryKeyBasedOnValue(int value) { + assert referenceDictionaryArray != null; + // reference array index will be -1 of the value as dictionary value starts from 1 + return referenceDictionaryArray[value - 1].getData(); + } +} http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/core/src/main/java/org/apache/carbondata/core/localdictionary/exception/DictionaryThresholdReachedException.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/localdictionary/exception/DictionaryThresholdReachedException.java b/core/src/main/java/org/apache/carbondata/core/localdictionary/exception/DictionaryThresholdReachedException.java new file mode 100644 index 0000000..7d648e0 --- /dev/null +++ b/core/src/main/java/org/apache/carbondata/core/localdictionary/exception/DictionaryThresholdReachedException.java @@ -0,0 +1,87 @@ +/* + * 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.carbondata.core.localdictionary.exception; + +import java.util.Locale; + +public class DictionaryThresholdReachedException extends Exception { + /** + * default serial version ID. + */ + private static final long serialVersionUID = 1L; + + /** + * The Error message. + */ + private String msg = ""; + + /** + * Constructor + * + * @param msg The error message for this exception. + */ + public DictionaryThresholdReachedException(String msg) { + super(msg); + this.msg = msg; + } + + /** + * Constructor + * + * @param msg exception message + * @param throwable detail exception + */ + public DictionaryThresholdReachedException(String msg, Throwable throwable) { + super(msg, throwable); + this.msg = msg; + } + + /** + * Constructor + * + * @param throwable exception + */ + public DictionaryThresholdReachedException(Throwable throwable) { + super(throwable); + } + + /** + * This method is used to get the localized message. + * + * @param locale - A Locale object represents a specific geographical, + * political, or cultural region. + * @return - Localized error message. + */ + public String getLocalizedMessage(Locale locale) { + return ""; + } + + /** + * getLocalizedMessage + */ + @Override public String getLocalizedMessage() { + return super.getLocalizedMessage(); + } + + /** + * getMessage + */ + public String getMessage() { + return this.msg; + } +} + http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/core/src/main/java/org/apache/carbondata/core/localdictionary/generator/ColumnLocalDictionaryGenerator.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/localdictionary/generator/ColumnLocalDictionaryGenerator.java b/core/src/main/java/org/apache/carbondata/core/localdictionary/generator/ColumnLocalDictionaryGenerator.java new file mode 100644 index 0000000..5ae9e27 --- /dev/null +++ b/core/src/main/java/org/apache/carbondata/core/localdictionary/generator/ColumnLocalDictionaryGenerator.java @@ -0,0 +1,75 @@ +/* + * 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.carbondata.core.localdictionary.generator; + +import org.apache.carbondata.core.constants.CarbonCommonConstants; +import org.apache.carbondata.core.localdictionary.dictionaryholder.DictionaryStore; +import org.apache.carbondata.core.localdictionary.dictionaryholder.MapBasedDictionaryStore; +import org.apache.carbondata.core.localdictionary.exception.DictionaryThresholdReachedException; + +/** + * Class to generate local dictionary for column + */ +public class ColumnLocalDictionaryGenerator implements LocalDictionaryGenerator { + + /** + * dictionary holder to hold dictionary values + */ + private DictionaryStore dictionaryHolder; + + public ColumnLocalDictionaryGenerator(int threshold) { + // adding 1 to threshold for null value + int newThreshold = threshold + 1; + this.dictionaryHolder = new MapBasedDictionaryStore(newThreshold); + // for handling null values + try { + dictionaryHolder.putIfAbsent(CarbonCommonConstants.MEMBER_DEFAULT_VAL_ARRAY); + } catch (DictionaryThresholdReachedException e) { + // do nothing + } + } + + /** + * Below method will be used to generate dictionary + * @param data + * data for which dictionary needs to be generated + * @return dictionary value + */ + @Override public int generateDictionary(byte[] data) throws DictionaryThresholdReachedException { + int dictionaryValue = this.dictionaryHolder.putIfAbsent(data); + return dictionaryValue; + } + + /** + * Below method will be used to check if threshold is reached + * for dictionary for particular column + * @return true if dictionary threshold reached for column + */ + @Override public boolean isThresholdReached() { + return this.dictionaryHolder.isThresholdReached(); + } + + /** + * Below method will be used to get the dictionary key based on value + * @param value + * dictionary value + * @return dictionary key based on value + */ + @Override public byte[] getDictionaryKeyBasedOnValue(int value) { + return this.dictionaryHolder.getDictionaryKeyBasedOnValue(value); + } +} http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/core/src/main/java/org/apache/carbondata/core/localdictionary/generator/LocalDictionaryGenerator.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/localdictionary/generator/LocalDictionaryGenerator.java b/core/src/main/java/org/apache/carbondata/core/localdictionary/generator/LocalDictionaryGenerator.java new file mode 100644 index 0000000..553c65b --- /dev/null +++ b/core/src/main/java/org/apache/carbondata/core/localdictionary/generator/LocalDictionaryGenerator.java @@ -0,0 +1,48 @@ +/* + * 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.carbondata.core.localdictionary.generator; + +import org.apache.carbondata.core.localdictionary.exception.DictionaryThresholdReachedException; + +/** + * Interface for generating dictionary for column + */ +public interface LocalDictionaryGenerator { + + /** + * Below method will be used to generate dictionary + * @param data + * data for which dictionary needs to be generated + * @return dictionary value + */ + int generateDictionary(byte[] data) throws DictionaryThresholdReachedException; + + /** + * Below method will be used to check if threshold is reached + * for dictionary for particular column + * @return true if dictionary threshold reached for column + */ + boolean isThresholdReached(); + + /** + * Below method will be used to get the dictionary key based on value + * @param value + * dictionary value + * @return dictionary key based on value + */ + byte[] getDictionaryKeyBasedOnValue(int value); +} http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/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 2cb19ea..68bd749 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 @@ -482,7 +482,7 @@ public class CarbonTable implements Serializable { * @return */ public boolean isLocalDictionaryEnabled() { - return isLocalDictionaryEnabled; + return false; } /** http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/core/src/main/java/org/apache/carbondata/core/util/CarbonMetadataUtil.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/util/CarbonMetadataUtil.java b/core/src/main/java/org/apache/carbondata/core/util/CarbonMetadataUtil.java index af5121c..58de030 100644 --- a/core/src/main/java/org/apache/carbondata/core/util/CarbonMetadataUtil.java +++ b/core/src/main/java/org/apache/carbondata/core/util/CarbonMetadataUtil.java @@ -23,7 +23,9 @@ import java.util.List; import java.util.Set; import org.apache.carbondata.core.datastore.block.SegmentProperties; -import org.apache.carbondata.core.datastore.page.EncodedTablePage; +import org.apache.carbondata.core.datastore.blocklet.BlockletEncodedColumnPage; +import org.apache.carbondata.core.datastore.blocklet.EncodedBlocklet; +import org.apache.carbondata.core.datastore.page.encoding.EncodedColumnPage; import org.apache.carbondata.core.datastore.page.statistics.TablePageStatistics; import org.apache.carbondata.core.metadata.ColumnarFormatVersion; import org.apache.carbondata.core.metadata.datatype.DataType; @@ -44,6 +46,7 @@ import org.apache.carbondata.format.Encoding; import org.apache.carbondata.format.FileFooter3; import org.apache.carbondata.format.FileHeader; import org.apache.carbondata.format.IndexHeader; +import org.apache.carbondata.format.LocalDictionaryChunk; import org.apache.carbondata.format.SegmentInfo; /** @@ -124,18 +127,39 @@ public class CarbonMetadataUtil { return numberOfRows; } - public static BlockletIndex getBlockletIndex(List<EncodedTablePage> encodedTablePageList, + private static EncodedColumnPage[] getEncodedColumnPages(EncodedBlocklet encodedBlocklet, + boolean isDimension, int pageIndex) { + int size = + isDimension ? encodedBlocklet.getNumberOfDimension() : encodedBlocklet.getNumberOfMeasure(); + EncodedColumnPage [] encodedPages = new EncodedColumnPage[size]; + + for (int i = 0; i < size; i++) { + if (isDimension) { + encodedPages[i] = + encodedBlocklet.getEncodedDimensionColumnPages().get(i).getEncodedColumnPageList() + .get(pageIndex); + } else { + encodedPages[i] = + encodedBlocklet.getEncodedMeasureColumnPages().get(i).getEncodedColumnPageList() + .get(pageIndex); + } + } + return encodedPages; + } + public static BlockletIndex getBlockletIndex(EncodedBlocklet encodedBlocklet, List<CarbonMeasure> carbonMeasureList) { BlockletMinMaxIndex blockletMinMaxIndex = new BlockletMinMaxIndex(); + // Calculating min/max for every each column. - TablePageStatistics stats = new TablePageStatistics(encodedTablePageList.get(0).getDimensions(), - encodedTablePageList.get(0).getMeasures()); + TablePageStatistics stats = + new TablePageStatistics(getEncodedColumnPages(encodedBlocklet, true, 0), + getEncodedColumnPages(encodedBlocklet, false, 0)); byte[][] minCol = stats.getDimensionMinValue().clone(); byte[][] maxCol = stats.getDimensionMaxValue().clone(); - for (EncodedTablePage encodedTablePage : encodedTablePageList) { - stats = new TablePageStatistics(encodedTablePage.getDimensions(), - encodedTablePage.getMeasures()); + for (int pageIndex = 0; pageIndex < encodedBlocklet.getNumberOfPages(); pageIndex++) { + stats = new TablePageStatistics(getEncodedColumnPages(encodedBlocklet, true, pageIndex), + getEncodedColumnPages(encodedBlocklet, false, pageIndex)); byte[][] columnMaxData = stats.getDimensionMaxValue(); byte[][] columnMinData = stats.getDimensionMinValue(); for (int i = 0; i < maxCol.length; i++) { @@ -155,16 +179,16 @@ public class CarbonMetadataUtil { blockletMinMaxIndex.addToMin_values(ByteBuffer.wrap(min)); } - stats = new TablePageStatistics(encodedTablePageList.get(0).getDimensions(), - encodedTablePageList.get(0).getMeasures()); + stats = new TablePageStatistics(getEncodedColumnPages(encodedBlocklet, true, 0), + getEncodedColumnPages(encodedBlocklet, false, 0)); byte[][] measureMaxValue = stats.getMeasureMaxValue().clone(); byte[][] measureMinValue = stats.getMeasureMinValue().clone(); byte[] minVal = null; byte[] maxVal = null; - for (int i = 1; i < encodedTablePageList.size(); i++) { + for (int i = 1; i < encodedBlocklet.getNumberOfPages(); i++) { for (int j = 0; j < measureMinValue.length; j++) { - stats = new TablePageStatistics( - encodedTablePageList.get(i).getDimensions(), encodedTablePageList.get(i).getMeasures()); + stats = new TablePageStatistics(getEncodedColumnPages(encodedBlocklet, true, i), + getEncodedColumnPages(encodedBlocklet, false, i)); minVal = stats.getMeasureMinValue()[j]; maxVal = stats.getMeasureMaxValue()[j]; if (compareMeasureData(measureMaxValue[j], maxVal, carbonMeasureList.get(j).getDataType()) @@ -185,10 +209,11 @@ public class CarbonMetadataUtil { blockletMinMaxIndex.addToMin_values(ByteBuffer.wrap(min)); } BlockletBTreeIndex blockletBTreeIndex = new BlockletBTreeIndex(); - byte[] startKey = encodedTablePageList.get(0).getPageKey().serializeStartKey(); + byte[] startKey = encodedBlocklet.getPageMetadataList().get(0).serializeStartKey(); blockletBTreeIndex.setStart_key(startKey); - byte[] endKey = encodedTablePageList.get( - encodedTablePageList.size() - 1).getPageKey().serializeEndKey(); + byte[] endKey = + encodedBlocklet.getPageMetadataList().get(encodedBlocklet.getPageMetadataList().size() - 1) + .serializeEndKey(); blockletBTreeIndex.setEnd_key(endKey); BlockletIndex blockletIndex = new BlockletIndex(); blockletIndex.setMin_max_index(blockletMinMaxIndex); @@ -300,7 +325,8 @@ public class CarbonMetadataUtil { /** * return DataChunk3 that contains the input DataChunk2 list */ - public static DataChunk3 getDataChunk3(List<DataChunk2> dataChunksList) { + public static DataChunk3 getDataChunk3(List<DataChunk2> dataChunksList, + LocalDictionaryChunk encodedDictionary) { int offset = 0; DataChunk3 dataChunk = new DataChunk3(); List<Integer> pageOffsets = new ArrayList<>(); @@ -313,6 +339,7 @@ public class CarbonMetadataUtil { pageLengths.add(length); offset += length; } + dataChunk.setLocal_dictionary(encodedDictionary); dataChunk.setData_chunk_list(dataChunksList); dataChunk.setPage_length(pageLengths); dataChunk.setPage_offset(pageOffsets); @@ -323,26 +350,32 @@ public class CarbonMetadataUtil { * return DataChunk3 for the dimension column (specifed by `columnIndex`) * in `encodedTablePageList` */ - public static DataChunk3 getDimensionDataChunk3(List<EncodedTablePage> encodedTablePageList, - int columnIndex) throws IOException { - List<DataChunk2> dataChunksList = new ArrayList<>(encodedTablePageList.size()); - for (EncodedTablePage encodedTablePage : encodedTablePageList) { - dataChunksList.add(encodedTablePage.getDimension(columnIndex).getPageMetadata()); + public static DataChunk3 getDimensionDataChunk3(EncodedBlocklet encodedBlocklet, + int columnIndex) { + List<DataChunk2> dataChunksList = new ArrayList<>(); + BlockletEncodedColumnPage blockletEncodedColumnPage = + encodedBlocklet.getEncodedDimensionColumnPages().get(columnIndex); + for (EncodedColumnPage encodedColumnPage : blockletEncodedColumnPage + .getEncodedColumnPageList()) { + dataChunksList.add(encodedColumnPage.getPageMetadata()); } - return CarbonMetadataUtil.getDataChunk3(dataChunksList); + return CarbonMetadataUtil + .getDataChunk3(dataChunksList, blockletEncodedColumnPage.getEncodedDictionary()); } /** * return DataChunk3 for the measure column (specifed by `columnIndex`) * in `encodedTablePageList` */ - public static DataChunk3 getMeasureDataChunk3(List<EncodedTablePage> encodedTablePageList, - int columnIndex) throws IOException { - List<DataChunk2> dataChunksList = new ArrayList<>(encodedTablePageList.size()); - for (EncodedTablePage encodedTablePage : encodedTablePageList) { - dataChunksList.add(encodedTablePage.getMeasure(columnIndex).getPageMetadata()); + public static DataChunk3 getMeasureDataChunk3(EncodedBlocklet encodedBlocklet, int columnIndex) { + List<DataChunk2> dataChunksList = new ArrayList<>(); + BlockletEncodedColumnPage blockletEncodedColumnPage = + encodedBlocklet.getEncodedMeasureColumnPages().get(columnIndex); + for (EncodedColumnPage encodedColumnPage : blockletEncodedColumnPage + .getEncodedColumnPageList()) { + dataChunksList.add(encodedColumnPage.getPageMetadata()); } - return CarbonMetadataUtil.getDataChunk3(dataChunksList); + return CarbonMetadataUtil.getDataChunk3(dataChunksList, null); } private static int compareMeasureData(byte[] first, byte[] second, DataType dataType) { http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/core/src/test/java/org/apache/carbondata/core/util/CarbonMetadataUtilTest.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/carbondata/core/util/CarbonMetadataUtilTest.java b/core/src/test/java/org/apache/carbondata/core/util/CarbonMetadataUtilTest.java index da31ea3..2909dc4 100644 --- a/core/src/test/java/org/apache/carbondata/core/util/CarbonMetadataUtilTest.java +++ b/core/src/test/java/org/apache/carbondata/core/util/CarbonMetadataUtilTest.java @@ -173,71 +173,6 @@ public class CarbonMetadataUtilTest { assertEquals(indexHeader, indexheaderResult); } - @Test public void testConvertFileFooter() throws Exception { - int[] cardinality = { 1, 2, 3, 4, 5 }; - - org.apache.carbondata.core.metadata.schema.table.column.ColumnSchema colSchema = - new org.apache.carbondata.core.metadata.schema.table.column.ColumnSchema(); - org.apache.carbondata.core.metadata.schema.table.column.ColumnSchema colSchema1 = - new org.apache.carbondata.core.metadata.schema.table.column.ColumnSchema(); - List<org.apache.carbondata.core.metadata.schema.table.column.ColumnSchema> - columnSchemaList = new ArrayList<>(); - columnSchemaList.add(colSchema); - columnSchemaList.add(colSchema1); - - SegmentProperties segmentProperties = new SegmentProperties(columnSchemaList, cardinality); - - final EncodedColumnPage measure = new EncodedColumnPage(new DataChunk2(), new byte[]{0,1}, - PrimitivePageStatsCollector.newInstance( - org.apache.carbondata.core.metadata.datatype.DataTypes.BYTE)); - new MockUp<EncodedTablePage>() { - @SuppressWarnings("unused") @Mock - public EncodedColumnPage getMeasure(int measureIndex) { - return measure; - } - }; - - new MockUp<TablePageKey>() { - @SuppressWarnings("unused") @Mock - public byte[] serializeStartKey() { - return new byte[]{1, 2}; - } - - @SuppressWarnings("unused") @Mock - public byte[] serializeEndKey() { - return new byte[]{1, 2}; - } - }; - - TablePageKey key = new TablePageKey(3, segmentProperties, false); - EncodedTablePage encodedTablePage = EncodedTablePage.newInstance(3, new EncodedColumnPage[0], new EncodedColumnPage[0], - key); - - List<EncodedTablePage> encodedTablePageList = new ArrayList<>(); - encodedTablePageList.add(encodedTablePage); - - BlockletInfo3 blockletInfoColumnar1 = new BlockletInfo3(); - - List<BlockletInfo3> blockletInfoColumnarList = new ArrayList<>(); - blockletInfoColumnarList.add(blockletInfoColumnar1); - - byte[] byteMaxArr = "1".getBytes(); - byte[] byteMinArr = "2".getBytes(); - - BlockletIndex index = getBlockletIndex(encodedTablePageList, segmentProperties.getMeasures()); - List<BlockletIndex> indexList = new ArrayList<>(); - indexList.add(index); - - BlockletMinMaxIndex blockletMinMaxIndex = new BlockletMinMaxIndex(); - blockletMinMaxIndex.addToMax_values(ByteBuffer.wrap(byteMaxArr)); - blockletMinMaxIndex.addToMin_values(ByteBuffer.wrap(byteMinArr)); - FileFooter3 footer = convertFileFooterVersion3(blockletInfoColumnarList, - indexList, - cardinality, 2); - assertEquals(footer.getBlocklet_index_list(), indexList); - - } - @Test public void testGetBlockIndexInfo() throws Exception { byte[] startKey = { 1, 2, 3, 4, 5 }; byte[] endKey = { 9, 3, 5, 5, 5 }; http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/format/src/main/thrift/carbondata.thrift ---------------------------------------------------------------------- diff --git a/format/src/main/thrift/carbondata.thrift b/format/src/main/thrift/carbondata.thrift index 1c15f3d..a495b6d 100644 --- a/format/src/main/thrift/carbondata.thrift +++ b/format/src/main/thrift/carbondata.thrift @@ -145,6 +145,7 @@ struct DataChunk3{ 1: required list<DataChunk2> data_chunk_list; // List of data chunk 2: optional list<i32> page_offset; // Offset of each chunk 3: optional list<i32> page_length; // Length of each chunk + 4: optional LocalDictionaryChunk local_dictionary; // to store blocklet local dictionary values } /** @@ -230,4 +231,15 @@ struct BlockletHeader{ 3: optional BlockletIndex blocklet_index; // Index for the following blocklet 4: required BlockletInfo blocklet_info; // Info for the following blocklet 5: optional dictionary.ColumnDictionaryChunk dictionary; // Blocklet local dictionary +} + +struct LocalDictionaryChunk { + 1: required LocalDictionaryChunkMeta dictionary_meta + 2: required binary dictionary_data; // the values in dictionary order, each value is represented in binary format + 3: required binary dictionary_values; // surrogate keys used in the blocklet +} + +struct LocalDictionaryChunkMeta { + 1: required list<schema.Encoding> encoders; // The List of encoders overriden at node level + 2: required list<binary> encoder_meta; // Extra information required by encoders } \ No newline at end of file http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/processing/src/main/java/org/apache/carbondata/processing/datatypes/ArrayDataType.java ---------------------------------------------------------------------- diff --git a/processing/src/main/java/org/apache/carbondata/processing/datatypes/ArrayDataType.java b/processing/src/main/java/org/apache/carbondata/processing/datatypes/ArrayDataType.java index 4ce80a6..da34746 100644 --- a/processing/src/main/java/org/apache/carbondata/processing/datatypes/ArrayDataType.java +++ b/processing/src/main/java/org/apache/carbondata/processing/datatypes/ArrayDataType.java @@ -65,10 +65,12 @@ public class ArrayDataType implements GenericDataType<ArrayObject> { */ private int dataCounter; - private ArrayDataType(int outputArrayIndex, int dataCounter, GenericDataType children) { + private ArrayDataType(int outputArrayIndex, int dataCounter, GenericDataType children, + String name) { this.outputArrayIndex = outputArrayIndex; this.dataCounter = dataCounter; this.children = children; + this.name = name; } @@ -108,7 +110,7 @@ public class ArrayDataType implements GenericDataType<ArrayObject> { * return column unique id */ @Override - public String getColumnId() { + public String getColumnNames() { return columnId; } @@ -285,7 +287,8 @@ public class ArrayDataType implements GenericDataType<ArrayObject> { @Override public GenericDataType<ArrayObject> deepCopy() { - return new ArrayDataType(this.outputArrayIndex, this.dataCounter, this.children.deepCopy()); + return new ArrayDataType(this.outputArrayIndex, this.dataCounter, this.children.deepCopy(), + this.name); } @Override @@ -293,4 +296,10 @@ public class ArrayDataType implements GenericDataType<ArrayObject> { type.add(ColumnType.COMPLEX_ARRAY); children.getChildrenType(type); } + + @Override public void getColumnNames(List<String> columnNameList) { + columnNameList.add(name); + children.getColumnNames(columnNameList); + } + } \ No newline at end of file http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/processing/src/main/java/org/apache/carbondata/processing/datatypes/GenericDataType.java ---------------------------------------------------------------------- diff --git a/processing/src/main/java/org/apache/carbondata/processing/datatypes/GenericDataType.java b/processing/src/main/java/org/apache/carbondata/processing/datatypes/GenericDataType.java index 8b1ccf2..049bf57 100644 --- a/processing/src/main/java/org/apache/carbondata/processing/datatypes/GenericDataType.java +++ b/processing/src/main/java/org/apache/carbondata/processing/datatypes/GenericDataType.java @@ -100,7 +100,7 @@ public interface GenericDataType<T> { /** * @return column uuid string */ - String getColumnId(); + String getColumnNames(); /** * set array index to be referred while creating metadata column @@ -159,4 +159,6 @@ public interface GenericDataType<T> { void getChildrenType(List<ColumnType> type); + void getColumnNames(List<String> columnNameList); + } http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/processing/src/main/java/org/apache/carbondata/processing/datatypes/PrimitiveDataType.java ---------------------------------------------------------------------- diff --git a/processing/src/main/java/org/apache/carbondata/processing/datatypes/PrimitiveDataType.java b/processing/src/main/java/org/apache/carbondata/processing/datatypes/PrimitiveDataType.java index 3a477ce..5d22e55 100644 --- a/processing/src/main/java/org/apache/carbondata/processing/datatypes/PrimitiveDataType.java +++ b/processing/src/main/java/org/apache/carbondata/processing/datatypes/PrimitiveDataType.java @@ -235,7 +235,7 @@ public class PrimitiveDataType implements GenericDataType<Object> { * get column unique id */ @Override - public String getColumnId() { + public String getColumnNames() { return columnId; } @@ -536,11 +536,15 @@ public class PrimitiveDataType implements GenericDataType<Object> { dataType.nullformat = this.nullformat; dataType.setKeySize(this.keySize); dataType.setSurrogateIndex(this.index); - + dataType.name = this.name; return dataType; } public void getChildrenType(List<ColumnType> type) { type.add(ColumnType.COMPLEX_PRIMITIVE); } + + @Override public void getColumnNames(List<String> columnNameList) { + columnNameList.add(name); + } } \ No newline at end of file http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/processing/src/main/java/org/apache/carbondata/processing/datatypes/StructDataType.java ---------------------------------------------------------------------- diff --git a/processing/src/main/java/org/apache/carbondata/processing/datatypes/StructDataType.java b/processing/src/main/java/org/apache/carbondata/processing/datatypes/StructDataType.java index b66eef7..4d3ba87 100644 --- a/processing/src/main/java/org/apache/carbondata/processing/datatypes/StructDataType.java +++ b/processing/src/main/java/org/apache/carbondata/processing/datatypes/StructDataType.java @@ -60,10 +60,12 @@ public class StructDataType implements GenericDataType<StructObject> { */ private int dataCounter; - private StructDataType(List<GenericDataType> children, int outputArrayIndex, int dataCounter) { + private StructDataType(List<GenericDataType> children, int outputArrayIndex, int dataCounter, + String name) { this.children = children; this.outputArrayIndex = outputArrayIndex; this.dataCounter = dataCounter; + this.name = name; } /** @@ -113,7 +115,7 @@ public class StructDataType implements GenericDataType<StructObject> { * get column unique id */ @Override - public String getColumnId() { + public String getColumnNames() { return columnId; } @@ -318,7 +320,7 @@ public class StructDataType implements GenericDataType<StructObject> { for (GenericDataType child : children) { childrenClone.add(child.deepCopy()); } - return new StructDataType(childrenClone, this.outputArrayIndex, this.dataCounter); + return new StructDataType(childrenClone, this.outputArrayIndex, this.dataCounter, this.name); } public void getChildrenType(List<ColumnType> type) { @@ -327,4 +329,11 @@ public class StructDataType implements GenericDataType<StructObject> { children.get(i).getChildrenType(type); } } + + @Override public void getColumnNames(List<String> columnNameList) { + columnNameList.add(name); + for (int i = 0; i < children.size(); i++) { + children.get(i).getColumnNames(columnNameList); + } + } } \ No newline at end of file http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/processing/src/main/java/org/apache/carbondata/processing/store/CarbonFactDataHandlerColumnar.java ---------------------------------------------------------------------- diff --git a/processing/src/main/java/org/apache/carbondata/processing/store/CarbonFactDataHandlerColumnar.java b/processing/src/main/java/org/apache/carbondata/processing/store/CarbonFactDataHandlerColumnar.java index 5fe3261..f3cb9c3 100644 --- a/processing/src/main/java/org/apache/carbondata/processing/store/CarbonFactDataHandlerColumnar.java +++ b/processing/src/main/java/org/apache/carbondata/processing/store/CarbonFactDataHandlerColumnar.java @@ -49,7 +49,6 @@ import org.apache.carbondata.core.util.CarbonProperties; import org.apache.carbondata.core.util.CarbonThreadFactory; import org.apache.carbondata.core.util.CarbonUtil; import org.apache.carbondata.processing.datatypes.GenericDataType; -import org.apache.carbondata.processing.loading.sort.SortScopeOptions; import org.apache.carbondata.processing.store.writer.CarbonFactDataWriter; /** @@ -137,44 +136,19 @@ public class CarbonFactDataHandlerColumnar implements CarbonFactHandler { } private void initParameters(CarbonFactDataHandlerModel model) { - SortScopeOptions.SortScope sortScope = model.getSortScope(); this.colGrpModel = model.getSegmentProperties().getColumnGroupModel(); - - // in compaction flow the measure with decimal type will come as spark decimal. - // need to convert it to byte array. - if (model.isCompactionFlow()) { - try { - numberOfCores = Integer.parseInt(CarbonProperties.getInstance() - .getProperty(CarbonCommonConstants.NUM_CORES_COMPACTING, - CarbonCommonConstants.NUM_CORES_DEFAULT_VAL)); - } catch (NumberFormatException exc) { - LOGGER.error("Configured value for property " + CarbonCommonConstants.NUM_CORES_COMPACTING - + "is wrong.Falling back to the default value " - + CarbonCommonConstants.NUM_CORES_DEFAULT_VAL); - numberOfCores = Integer.parseInt(CarbonCommonConstants.NUM_CORES_DEFAULT_VAL); - } - } else { - numberOfCores = CarbonProperties.getInstance().getNumberOfCores(); - } - - if (sortScope != null && sortScope.equals(SortScopeOptions.SortScope.GLOBAL_SORT)) { - numberOfCores = 1; - } - // Overriding it to the task specified cores. - if (model.getWritingCoresCount() > 0) { - numberOfCores = model.getWritingCoresCount(); - } - + this.numberOfCores = model.getNumberOfCores(); blockletProcessingCount = new AtomicInteger(0); - producerExecutorService = Executors.newFixedThreadPool(numberOfCores, - new CarbonThreadFactory("ProducerPool:" + model.getTableName() - + ", range: " + model.getBucketId())); + producerExecutorService = Executors.newFixedThreadPool(model.getNumberOfCores(), + new CarbonThreadFactory( + "ProducerPool_" + System.nanoTime() + ":" + model.getTableName() + ", range: " + model + .getBucketId())); producerExecutorServiceTaskList = new ArrayList<>(CarbonCommonConstants.DEFAULT_COLLECTION_SIZE); LOGGER.info("Initializing writer executors"); - consumerExecutorService = Executors - .newFixedThreadPool(1, new CarbonThreadFactory("ConsumerPool:" + model.getTableName() - + ", range: " + model.getBucketId())); + consumerExecutorService = Executors.newFixedThreadPool(1, new CarbonThreadFactory( + "ConsumerPool_" + System.nanoTime() + ":" + model.getTableName() + ", range: " + model + .getBucketId())); consumerExecutorServiceTaskList = new ArrayList<>(1); semaphore = new Semaphore(numberOfCores); tablePageList = new TablePageList(); http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/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 27249ab..5b12229 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 @@ -23,10 +23,14 @@ import java.util.Iterator; import java.util.List; import java.util.Map; +import org.apache.carbondata.common.logging.LogService; +import org.apache.carbondata.common.logging.LogServiceFactory; import org.apache.carbondata.core.constants.CarbonCommonConstants; import org.apache.carbondata.core.datastore.TableSpec; import org.apache.carbondata.core.datastore.block.SegmentProperties; import org.apache.carbondata.core.keygenerator.KeyGenerator; +import org.apache.carbondata.core.localdictionary.generator.ColumnLocalDictionaryGenerator; +import org.apache.carbondata.core.localdictionary.generator.LocalDictionaryGenerator; import org.apache.carbondata.core.metadata.AbsoluteTableIdentifier; import org.apache.carbondata.core.metadata.CarbonMetadata; import org.apache.carbondata.core.metadata.CarbonTableIdentifier; @@ -35,6 +39,7 @@ import org.apache.carbondata.core.metadata.schema.table.CarbonTable; import org.apache.carbondata.core.metadata.schema.table.column.CarbonDimension; import org.apache.carbondata.core.metadata.schema.table.column.CarbonMeasure; import org.apache.carbondata.core.metadata.schema.table.column.ColumnSchema; +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; @@ -50,6 +55,12 @@ import org.apache.carbondata.processing.util.CarbonDataProcessorUtil; public class CarbonFactDataHandlerModel { /** + * LOGGER + */ + private static final LogService LOGGER = + LogServiceFactory.getLogService(CarbonFactDataHandlerModel.class.getName()); + + /** * dbName */ private String databaseName; @@ -163,6 +174,10 @@ public class CarbonFactDataHandlerModel { private short writingCoresCount; + private Map<String, LocalDictionaryGenerator> columnLocalDictGenMap; + + private int numberOfCores; + /** * Create the model using @{@link CarbonDataLoadConfiguration} */ @@ -272,7 +287,8 @@ public class CarbonFactDataHandlerModel { } carbonFactDataHandlerModel.dataMapWriterlistener = listener; carbonFactDataHandlerModel.writingCoresCount = configuration.getWritingCoresCount(); - + setLocalDictToModel(carbonTable, wrapperColumnSchema, carbonFactDataHandlerModel); + setNumberOfCores(carbonFactDataHandlerModel); return carbonFactDataHandlerModel; } @@ -340,8 +356,9 @@ public class CarbonFactDataHandlerModel { carbonFactDataHandlerModel.getTaskExtension(), String.valueOf(loadModel.getFactTimeStamp()), loadModel.getSegmentId())); - + setLocalDictToModel(carbonTable, wrapperColumnSchema, carbonFactDataHandlerModel); carbonFactDataHandlerModel.dataMapWriterlistener = listener; + setNumberOfCores(carbonFactDataHandlerModel); return carbonFactDataHandlerModel; } @@ -623,5 +640,86 @@ public class CarbonFactDataHandlerModel { return dataMapWriterlistener; } + public Map<String, LocalDictionaryGenerator> getColumnLocalDictGenMap() { + return columnLocalDictGenMap; + } + + /** + * This method prepares a map which will have column and local dictionary generator mapping for + * all the local dictionary columns. + * @param carbonTable + * @param wrapperColumnSchema + * @param carbonFactDataHandlerModel + */ + private static void setLocalDictToModel(CarbonTable carbonTable, + List<ColumnSchema> wrapperColumnSchema, + CarbonFactDataHandlerModel carbonFactDataHandlerModel) { + boolean islocalDictEnabled = carbonTable.isLocalDictionaryEnabled(); + // creates a map only if local dictionary is enabled, else map will be null + Map<String, LocalDictionaryGenerator> columnLocalDictGenMap = new HashMap<>(); + if (islocalDictEnabled) { + int localDictionaryThreshold = carbonTable.getLocalDictionaryThreshold(); + for (ColumnSchema columnSchema : wrapperColumnSchema) { + // check whether the column is local dictionary column or not + if (columnSchema.isLocalDictColumn()) { + columnLocalDictGenMap.put(columnSchema.getColumnName(), + new ColumnLocalDictionaryGenerator(localDictionaryThreshold)); + } + } + } + if (islocalDictEnabled) { + LOGGER.info("Local dictionary is enabled for table: " + carbonTable.getTableUniqueName()); + LOGGER.info( + "Local dictionary threshold for table: " + carbonTable.getTableUniqueName() + " is: " + + carbonTable.getLocalDictionaryThreshold()); + Iterator<Map.Entry<String, LocalDictionaryGenerator>> iterator = + columnLocalDictGenMap.entrySet().iterator(); + StringBuilder stringBuilder = new StringBuilder(); + while (iterator.hasNext()) { + Map.Entry<String, LocalDictionaryGenerator> next = iterator.next(); + stringBuilder.append(next.getKey()); + stringBuilder.append(','); + } + LOGGER.info("Local dictionary will be generated for the columns:" + stringBuilder.toString() + + " for table: " + carbonTable.getTableUniqueName()); + } + carbonFactDataHandlerModel.setColumnLocalDictGenMap(columnLocalDictGenMap); + } + + public void setColumnLocalDictGenMap( + Map<String, LocalDictionaryGenerator> columnLocalDictGenMap) { + this.columnLocalDictGenMap = columnLocalDictGenMap; + } + + private static void setNumberOfCores(CarbonFactDataHandlerModel model) { + // in compaction flow the measure with decimal type will come as spark decimal. + // need to convert it to byte array. + if (model.isCompactionFlow()) { + try { + model.numberOfCores = Integer.parseInt(CarbonProperties.getInstance() + .getProperty(CarbonCommonConstants.NUM_CORES_COMPACTING, + CarbonCommonConstants.NUM_CORES_DEFAULT_VAL)); + } catch (NumberFormatException exc) { + LOGGER.error("Configured value for property " + CarbonCommonConstants.NUM_CORES_COMPACTING + + "is wrong.Falling back to the default value " + + CarbonCommonConstants.NUM_CORES_DEFAULT_VAL); + model.numberOfCores = Integer.parseInt(CarbonCommonConstants.NUM_CORES_DEFAULT_VAL); + } + } else { + model.numberOfCores = CarbonProperties.getInstance().getNumberOfCores(); + } + + if (model.sortScope != null && model.sortScope.equals(SortScopeOptions.SortScope.GLOBAL_SORT)) { + model.numberOfCores = 1; + } + // Overriding it to the task specified cores. + if (model.getWritingCoresCount() > 0) { + model.numberOfCores = model.getWritingCoresCount(); + } + } + + public int getNumberOfCores() { + return numberOfCores; + } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/processing/src/main/java/org/apache/carbondata/processing/store/TablePage.java ---------------------------------------------------------------------- diff --git a/processing/src/main/java/org/apache/carbondata/processing/store/TablePage.java b/processing/src/main/java/org/apache/carbondata/processing/store/TablePage.java index b1b966b..c634a6d 100644 --- a/processing/src/main/java/org/apache/carbondata/processing/store/TablePage.java +++ b/processing/src/main/java/org/apache/carbondata/processing/store/TablePage.java @@ -45,6 +45,7 @@ import org.apache.carbondata.core.datastore.page.statistics.PrimitivePageStatsCo import org.apache.carbondata.core.datastore.row.CarbonRow; import org.apache.carbondata.core.datastore.row.WriteStepRowUtil; import org.apache.carbondata.core.keygenerator.KeyGenException; +import org.apache.carbondata.core.localdictionary.generator.LocalDictionaryGenerator; import org.apache.carbondata.core.memory.MemoryException; import org.apache.carbondata.core.metadata.datatype.DataType; import org.apache.carbondata.core.metadata.datatype.DataTypes; @@ -64,8 +65,8 @@ public class TablePage { // one vector to make it efficient for sorting private ColumnPage[] dictDimensionPages; private ColumnPage[] noDictDimensionPages; - private ComplexColumnPage[] complexDimensionPages; private ColumnPage[] measurePages; + private ComplexColumnPage[] complexDimensionPages; // the num of rows in this page, it must be less than short value (65536) private int pageSize; @@ -104,19 +105,26 @@ public class TablePage { page.setStatsCollector(KeyPageStatsCollector.newInstance(DataTypes.BYTE_ARRAY)); dictDimensionPages[tmpNumDictDimIdx++] = page; } else { + // will be encoded using string page + LocalDictionaryGenerator localDictionaryGenerator = + model.getColumnLocalDictGenMap().get(spec.getFieldName()); + DataType dataType = DataTypes.STRING; if (DataTypes.VARCHAR == spec.getSchemaDataType()) { - page = ColumnPage.newPage(spec, DataTypes.VARCHAR, pageSize); + dataType = DataTypes.VARCHAR; + } + if (null != localDictionaryGenerator) { + page = ColumnPage.newLocalDictPage(spec, dataType, pageSize, localDictionaryGenerator); + } else { + page = ColumnPage.newPage(spec, dataType, pageSize); + } + if (DataTypes.VARCHAR == dataType) { page.setStatsCollector(LVLongStringStatsCollector.newInstance()); } else { - // In previous implementation, other data types such as string, date and timestamp - // will be encoded using string page - page = ColumnPage.newPage(spec, DataTypes.STRING, pageSize); page.setStatsCollector(LVShortStringStatsCollector.newInstance()); } noDictDimensionPages[tmpNumNoDictDimIdx++] = page; } } - complexDimensionPages = new ComplexColumnPage[model.getComplexColumnCount()]; for (int i = 0; i < complexDimensionPages.length; i++) { // here we still do not the depth of the complex column, it will be initialized when @@ -137,6 +145,7 @@ public class TablePage { PrimitivePageStatsCollector.newInstance(dataTypes[i])); measurePages[i] = page; } + boolean hasNoDictionary = noDictDimensionPages.length > 0; this.key = new TablePageKey(pageSize, model.getSegmentProperties(), hasNoDictionary); @@ -225,8 +234,16 @@ public class TablePage { // initialize the page if first row if (rowId == 0) { List<ColumnType> complexColumnType = new ArrayList<>(); + List<String> columnNames = new ArrayList<>(); complexDataType.getChildrenType(complexColumnType); - complexDimensionPages[index] = new ComplexColumnPage(pageSize, complexColumnType); + complexDataType.getColumnNames(columnNames); + complexDimensionPages[index] = new ComplexColumnPage(complexColumnType); + try { + complexDimensionPages[index] + .initialize(model.getColumnLocalDictGenMap(), columnNames, pageSize); + } catch (MemoryException e) { + throw new RuntimeException(e); + } } int depthInComplexColumn = complexDimensionPages[index].getDepth(); @@ -253,7 +270,7 @@ public class TablePage { } for (int depth = 0; depth < depthInComplexColumn; depth++) { - complexDimensionPages[index].putComplexData(rowId, depth, encodedComplexColumnar.get(depth)); + complexDimensionPages[index].putComplexData(depth, encodedComplexColumnar.get(depth)); } } @@ -267,6 +284,11 @@ public class TablePage { for (ColumnPage page : measurePages) { page.freeMemory(); } + for (ComplexColumnPage page : complexDimensionPages) { + if (null != page) { + page.freeMemory(); + } + } } // Adds length as a short element (first 2 bytes) to the head of the input byte array http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/processing/src/main/java/org/apache/carbondata/processing/store/writer/AbstractFactDataWriter.java ---------------------------------------------------------------------- diff --git a/processing/src/main/java/org/apache/carbondata/processing/store/writer/AbstractFactDataWriter.java b/processing/src/main/java/org/apache/carbondata/processing/store/writer/AbstractFactDataWriter.java index b76722b..3082b91 100644 --- a/processing/src/main/java/org/apache/carbondata/processing/store/writer/AbstractFactDataWriter.java +++ b/processing/src/main/java/org/apache/carbondata/processing/store/writer/AbstractFactDataWriter.java @@ -149,6 +149,8 @@ public abstract class AbstractFactDataWriter implements CarbonFactDataWriter { */ private boolean enableDirectlyWriteData2Hdfs = false; + protected ExecutorService fallbackExecutorService; + public AbstractFactDataWriter(CarbonFactDataHandlerModel model) { this.model = model; blockIndexInfoList = new ArrayList<>(); @@ -197,6 +199,14 @@ public abstract class AbstractFactDataWriter implements CarbonFactDataWriter { blockletMetadata = new ArrayList<BlockletInfo3>(); blockletIndex = new ArrayList<>(); listener = this.model.getDataMapWriterlistener(); + if (model.getColumnLocalDictGenMap().size() > 0) { + int numberOfCores = 1; + if (model.getNumberOfCores() > 1) { + numberOfCores = model.getNumberOfCores() / 2; + } + fallbackExecutorService = Executors.newFixedThreadPool(numberOfCores, new CarbonThreadFactory( + "FallbackPool:" + model.getTableName() + ", range: " + model.getBucketId())); + } } /** @@ -415,6 +425,9 @@ public abstract class AbstractFactDataWriter implements CarbonFactDataWriter { } catch (InterruptedException | ExecutionException | IOException e) { throw new CarbonDataWriterException(e); } + if (null != fallbackExecutorService) { + fallbackExecutorService.shutdownNow(); + } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/processing/src/main/java/org/apache/carbondata/processing/store/writer/v3/BlockletDataHolder.java ---------------------------------------------------------------------- diff --git a/processing/src/main/java/org/apache/carbondata/processing/store/writer/v3/BlockletDataHolder.java b/processing/src/main/java/org/apache/carbondata/processing/store/writer/v3/BlockletDataHolder.java index 36fda3c..7607cf0 100644 --- a/processing/src/main/java/org/apache/carbondata/processing/store/writer/v3/BlockletDataHolder.java +++ b/processing/src/main/java/org/apache/carbondata/processing/store/writer/v3/BlockletDataHolder.java @@ -16,29 +16,34 @@ */ package org.apache.carbondata.processing.store.writer.v3; -import java.util.ArrayList; -import java.util.List; +import java.util.concurrent.ExecutorService; +import org.apache.carbondata.core.datastore.blocklet.EncodedBlocklet; import org.apache.carbondata.core.datastore.page.EncodedTablePage; import org.apache.carbondata.processing.store.TablePage; public class BlockletDataHolder { - private List<EncodedTablePage> encodedTablePage; + + /** + * current data size + */ private long currentSize; - public BlockletDataHolder() { - this.encodedTablePage = new ArrayList<>(); + private EncodedBlocklet encodedBlocklet; + + public BlockletDataHolder(ExecutorService fallbackpool) { + encodedBlocklet = new EncodedBlocklet(fallbackpool); } public void clear() { - encodedTablePage.clear(); currentSize = 0; + encodedBlocklet.clear(); } public void addPage(TablePage rawTablePage) { EncodedTablePage encodedTablePage = rawTablePage.getEncodedTablePage(); - this.encodedTablePage.add(encodedTablePage); currentSize += encodedTablePage.getEncodedSize(); + encodedBlocklet.addEncodedTablePage(encodedTablePage); } public long getSize() { @@ -47,19 +52,14 @@ public class BlockletDataHolder { } public int getNumberOfPagesAdded() { - return encodedTablePage.size(); + return encodedBlocklet.getNumberOfPages(); } public int getTotalRows() { - int rows = 0; - for (EncodedTablePage nh : encodedTablePage) { - rows += nh.getPageSize(); - } - return rows; + return encodedBlocklet.getBlockletSize(); } - public List<EncodedTablePage> getEncodedTablePages() { - return encodedTablePage; + public EncodedBlocklet getEncodedBlocklet() { + return encodedBlocklet; } - } http://git-wip-us.apache.org/repos/asf/carbondata/blob/e7103397/processing/src/main/java/org/apache/carbondata/processing/store/writer/v3/CarbonFactDataWriterImplV3.java ---------------------------------------------------------------------- diff --git a/processing/src/main/java/org/apache/carbondata/processing/store/writer/v3/CarbonFactDataWriterImplV3.java b/processing/src/main/java/org/apache/carbondata/processing/store/writer/v3/CarbonFactDataWriterImplV3.java index d1deef1..e562f26 100644 --- a/processing/src/main/java/org/apache/carbondata/processing/store/writer/v3/CarbonFactDataWriterImplV3.java +++ b/processing/src/main/java/org/apache/carbondata/processing/store/writer/v3/CarbonFactDataWriterImplV3.java @@ -25,8 +25,9 @@ import org.apache.carbondata.common.logging.LogService; import org.apache.carbondata.common.logging.LogServiceFactory; import org.apache.carbondata.core.constants.CarbonCommonConstants; import org.apache.carbondata.core.constants.CarbonV3DataFormatConstants; +import org.apache.carbondata.core.datastore.blocklet.BlockletEncodedColumnPage; +import org.apache.carbondata.core.datastore.blocklet.EncodedBlocklet; import org.apache.carbondata.core.datastore.exception.CarbonDataWriterException; -import org.apache.carbondata.core.datastore.page.EncodedTablePage; import org.apache.carbondata.core.datastore.page.encoding.EncodedColumnPage; import org.apache.carbondata.core.metadata.blocklet.BlockletInfo; import org.apache.carbondata.core.metadata.blocklet.index.BlockletBTreeIndex; @@ -76,7 +77,7 @@ public class CarbonFactDataWriterImplV3 extends AbstractFactDataWriter { blockletSizeThreshold = fileSizeInBytes; LOGGER.info("Blocklet size configure for table is: " + blockletSizeThreshold); } - blockletDataHolder = new BlockletDataHolder(); + blockletDataHolder = new BlockletDataHolder(fallbackExecutorService); } @Override protected void writeBlockletInfoToFile() @@ -110,14 +111,15 @@ public class CarbonFactDataWriterImplV3 extends AbstractFactDataWriter { */ @Override public void writeTablePage(TablePage tablePage) throws CarbonDataWriterException,IOException { + // condition for writting all the pages if (!tablePage.isLastPage()) { boolean isAdded = false; // check if size more than blocklet size then write the page to file - if (blockletDataHolder.getSize() + tablePage.getEncodedTablePage().getEncodedSize() >= - blockletSizeThreshold) { + if (blockletDataHolder.getSize() + tablePage.getEncodedTablePage().getEncodedSize() + >= blockletSizeThreshold) { // if blocklet size exceeds threshold, write blocklet data - if (blockletDataHolder.getEncodedTablePages().size() == 0) { + if (blockletDataHolder.getNumberOfPagesAdded() == 0) { isAdded = true; addPageData(tablePage); } @@ -164,12 +166,13 @@ public class CarbonFactDataWriterImplV3 extends AbstractFactDataWriter { */ private void writeBlockletToFile() { // get the list of all encoded table page - List<EncodedTablePage> encodedTablePageList = blockletDataHolder.getEncodedTablePages(); - int numDimensions = encodedTablePageList.get(0).getNumDimensions(); - int numMeasures = encodedTablePageList.get(0).getNumMeasures(); + EncodedBlocklet encodedBlocklet = blockletDataHolder.getEncodedBlocklet(); + int numDimensions = encodedBlocklet.getNumberOfDimension(); + int numMeasures = encodedBlocklet.getNumberOfMeasure(); + // get data chunks for all the column byte[][] dataChunkBytes = new byte[numDimensions + numMeasures][]; - long metadataSize = fillDataChunk(encodedTablePageList, dataChunkBytes); + long metadataSize = fillDataChunk(encodedBlocklet, dataChunkBytes); // calculate the total size of data to be written long blockletSize = blockletDataHolder.getSize() + metadataSize; // to check if data size will exceed the block size then create a new file @@ -199,27 +202,22 @@ public class CarbonFactDataWriterImplV3 extends AbstractFactDataWriter { /** * Fill dataChunkBytes and return total size of page metadata */ - private long fillDataChunk(List<EncodedTablePage> encodedTablePageList, byte[][] dataChunkBytes) { + private long fillDataChunk(EncodedBlocklet encodedBlocklet, byte[][] dataChunkBytes) { int size = 0; - int numDimensions = encodedTablePageList.get(0).getNumDimensions(); - int numMeasures = encodedTablePageList.get(0).getNumMeasures(); + int numDimensions = encodedBlocklet.getNumberOfDimension(); + int numMeasures = encodedBlocklet.getNumberOfMeasure(); int measureStartIndex = numDimensions; // calculate the size of data chunks - try { - for (int i = 0; i < numDimensions; i++) { - dataChunkBytes[i] = CarbonUtil.getByteArray( - CarbonMetadataUtil.getDimensionDataChunk3(encodedTablePageList, i)); - size += dataChunkBytes[i].length; - } - for (int i = 0; i < numMeasures; i++) { - dataChunkBytes[measureStartIndex] = CarbonUtil.getByteArray( - CarbonMetadataUtil.getMeasureDataChunk3(encodedTablePageList, i)); - size += dataChunkBytes[measureStartIndex].length; - measureStartIndex++; - } - } catch (IOException e) { - LOGGER.error(e, "Problem while getting the data chunks"); - throw new CarbonDataWriterException("Problem while getting the data chunks", e); + for (int i = 0; i < numDimensions; i++) { + dataChunkBytes[i] = + CarbonUtil.getByteArray(CarbonMetadataUtil.getDimensionDataChunk3(encodedBlocklet, i)); + size += dataChunkBytes[i].length; + } + for (int i = 0; i < numMeasures; i++) { + dataChunkBytes[measureStartIndex] = + CarbonUtil.getByteArray(CarbonMetadataUtil.getMeasureDataChunk3(encodedBlocklet, i)); + size += dataChunkBytes[measureStartIndex].length; + measureStartIndex++; } return size; } @@ -250,33 +248,30 @@ public class CarbonFactDataWriterImplV3 extends AbstractFactDataWriter { List<Long> currentDataChunksOffset = new ArrayList<>(); // to maintain the length of each data chunk in blocklet List<Integer> currentDataChunksLength = new ArrayList<>(); - List<EncodedTablePage> encodedTablePages = blockletDataHolder.getEncodedTablePages(); - int numberOfDimension = encodedTablePages.get(0).getNumDimensions(); - int numberOfMeasures = encodedTablePages.get(0).getNumMeasures(); + EncodedBlocklet encodedBlocklet = blockletDataHolder.getEncodedBlocklet(); + int numberOfDimension = encodedBlocklet.getNumberOfDimension(); + int numberOfMeasures = encodedBlocklet.getNumberOfMeasure(); ByteBuffer buffer = null; long dimensionOffset = 0; long measureOffset = 0; - int numberOfRows = 0; - // calculate the number of rows in each blocklet - for (EncodedTablePage encodedTablePage : encodedTablePages) { - numberOfRows += encodedTablePage.getPageSize(); - } for (int i = 0; i < numberOfDimension; i++) { currentDataChunksOffset.add(offset); currentDataChunksLength.add(dataChunkBytes[i].length); buffer = ByteBuffer.wrap(dataChunkBytes[i]); currentOffsetInFile += fileChannel.write(buffer); offset += dataChunkBytes[i].length; - for (EncodedTablePage encodedTablePage : encodedTablePages) { - EncodedColumnPage dimension = encodedTablePage.getDimension(i); - buffer = dimension.getEncodedData(); + BlockletEncodedColumnPage blockletEncodedColumnPage = + encodedBlocklet.getEncodedDimensionColumnPages().get(i); + for (EncodedColumnPage dimensionPage : blockletEncodedColumnPage + .getEncodedColumnPageList()) { + buffer = dimensionPage.getEncodedData(); int bufferSize = buffer.limit(); currentOffsetInFile += fileChannel.write(buffer); offset += bufferSize; } } dimensionOffset = offset; - int dataChunkStartIndex = encodedTablePages.get(0).getNumDimensions(); + int dataChunkStartIndex = encodedBlocklet.getNumberOfDimension(); for (int i = 0; i < numberOfMeasures; i++) { currentDataChunksOffset.add(offset); currentDataChunksLength.add(dataChunkBytes[dataChunkStartIndex].length); @@ -284,9 +279,11 @@ public class CarbonFactDataWriterImplV3 extends AbstractFactDataWriter { currentOffsetInFile += fileChannel.write(buffer); offset += dataChunkBytes[dataChunkStartIndex].length; dataChunkStartIndex++; - for (EncodedTablePage encodedTablePage : encodedTablePages) { - EncodedColumnPage measure = encodedTablePage.getMeasure(i); - buffer = measure.getEncodedData(); + BlockletEncodedColumnPage blockletEncodedColumnPage = + encodedBlocklet.getEncodedMeasureColumnPages().get(i); + for (EncodedColumnPage measurePage : blockletEncodedColumnPage + .getEncodedColumnPageList()) { + buffer = measurePage.getEncodedData(); int bufferSize = buffer.limit(); currentOffsetInFile += fileChannel.write(buffer); offset += bufferSize; @@ -295,10 +292,11 @@ public class CarbonFactDataWriterImplV3 extends AbstractFactDataWriter { measureOffset = offset; blockletIndex.add( CarbonMetadataUtil.getBlockletIndex( - encodedTablePages, model.getSegmentProperties().getMeasures())); + encodedBlocklet, model.getSegmentProperties().getMeasures())); BlockletInfo3 blockletInfo3 = - new BlockletInfo3(numberOfRows, currentDataChunksOffset, currentDataChunksLength, - dimensionOffset, measureOffset, blockletDataHolder.getEncodedTablePages().size()); + new BlockletInfo3(encodedBlocklet.getBlockletSize(), currentDataChunksOffset, + currentDataChunksLength, dimensionOffset, measureOffset, + encodedBlocklet.getNumberOfPages()); blockletMetadata.add(blockletInfo3); }
