[CARBONDATA-2589][CARBONDATA-2590][CARBONDATA-2602]Local dictionary query Support
Supported Non filter query for local dictionary Supported Filter query on local dictionary Supported Query on complex column for primitive type local dictionary columns Local Dictionary support on Varchar columns Supported Vector reader on local dictionary This closes #2447 Project: http://git-wip-us.apache.org/repos/asf/carbondata/repo Commit: http://git-wip-us.apache.org/repos/asf/carbondata/commit/3a4b8813 Tree: http://git-wip-us.apache.org/repos/asf/carbondata/tree/3a4b8813 Diff: http://git-wip-us.apache.org/repos/asf/carbondata/diff/3a4b8813 Branch: refs/heads/master Commit: 3a4b88138296bc56db03dbe7fe771b60b874a573 Parents: 4935cb1 Author: kumarvishal09 <[email protected]> Authored: Thu Jul 5 11:57:47 2018 +0530 Committer: kunal642 <[email protected]> Committed: Tue Jul 10 11:05:53 2018 +0530 ---------------------------------------------------------------------- .../core/constants/CarbonCommonConstants.java | 10 ++ .../blocklet/BlockletEncodedColumnPage.java | 38 +++++-- .../chunk/impl/DimensionRawColumnChunk.java | 62 +++++++++++ .../impl/FixedLengthDimensionColumnPage.java | 2 +- .../impl/VariableLengthDimensionColumnPage.java | 28 ++++- ...mpressedDimensionChunkFileBasedReaderV1.java | 3 +- ...mpressedDimensionChunkFileBasedReaderV2.java | 3 +- ...mpressedDimensionChunkFileBasedReaderV3.java | 18 +-- .../chunk/store/ColumnPageWrapper.java | 11 +- .../chunk/store/DimensionChunkStoreFactory.java | 53 ++++++--- .../impl/LocalDictDimensionDataChunkStore.java | 109 +++++++++++++++++++ .../core/datastore/page/ColumnPage.java | 7 +- .../core/datastore/page/ComplexColumnPage.java | 2 +- .../datastore/page/LocalDictColumnPage.java | 24 +++- .../datastore/page/VarLengthColumnPageBase.java | 2 +- .../localdictionary/PageLevelDictionary.java | 28 ++++- .../ColumnLocalDictionaryGenerator.java | 11 +- .../carbondata/core/scan/filter/FilterUtil.java | 99 +++++++++++++---- .../executer/ExcludeFilterExecuterImpl.java | 8 +- .../executer/IncludeFilterExecuterImpl.java | 11 +- .../executer/RangeValueFilterExecuterImpl.java | 4 +- .../RowLevelRangeGrtThanFiterExecuterImpl.java | 2 +- ...elRangeGrtrThanEquaToFilterExecuterImpl.java | 2 +- ...velRangeLessThanEqualFilterExecuterImpl.java | 2 +- ...RowLevelRangeLessThanFilterExecuterImpl.java | 2 +- .../scan/result/vector/CarbonColumnVector.java | 11 ++ .../scan/result/vector/CarbonColumnarBatch.java | 1 - .../scan/result/vector/CarbonDictionary.java | 30 +++++ .../vector/impl/CarbonColumnVectorImpl.java | 33 +++++- .../vector/impl/CarbonDictionaryImpl.java | 54 +++++++++ .../examples/CarbonSessionExample.scala | 2 +- .../presto/CarbonColumnVectorWrapper.java | 22 +++- .../carbondata/spark/util/CarbonScalaUtil.scala | 8 +- .../vectorreader/CarbonDictionaryWrapper.java | 44 ++++++++ .../vectorreader/ColumnarVectorWrapper.java | 46 ++++++++ .../VectorizedCarbonRecordReader.java | 17 +++ .../carbondata/processing/store/TablePage.java | 3 +- 37 files changed, 713 insertions(+), 99 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/constants/CarbonCommonConstants.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/constants/CarbonCommonConstants.java b/core/src/main/java/org/apache/carbondata/core/constants/CarbonCommonConstants.java index c420853..7a64513 100644 --- a/core/src/main/java/org/apache/carbondata/core/constants/CarbonCommonConstants.java +++ b/core/src/main/java/org/apache/carbondata/core/constants/CarbonCommonConstants.java @@ -944,6 +944,16 @@ public final class CarbonCommonConstants { public static final String LOCAL_DICTIONARY_THRESHOLD_DEFAULT = "10000"; /** + * max dictionary threshold + */ + public static final int LOCAL_DICTIONARY_MAX = 100000; + + /** + * min dictionary threshold + */ + public static final int LOCAL_DICTIONARY_MIN = 1000; + + /** * Table property to specify the columns for which local dictionary needs to be generated. */ public static final String LOCAL_DICTIONARY_INCLUDE = "local_dictionary_include"; http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/blocklet/BlockletEncodedColumnPage.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/blocklet/BlockletEncodedColumnPage.java b/core/src/main/java/org/apache/carbondata/core/datastore/blocklet/BlockletEncodedColumnPage.java index 0fee09c..15adb2d 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/blocklet/BlockletEncodedColumnPage.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/blocklet/BlockletEncodedColumnPage.java @@ -73,8 +73,11 @@ public class BlockletEncodedColumnPage { */ private ArrayDeque<Future<FallbackEncodedColumnPage>> fallbackFutureQueue; + private String columnName; + BlockletEncodedColumnPage(ExecutorService fallbackExecutorService) { this.fallbackExecutorService = fallbackExecutorService; + this.fallbackFutureQueue = new ArrayDeque<>(); } /** @@ -92,29 +95,42 @@ public class BlockletEncodedColumnPage { // get first page dictionary this.pageLevelDictionary = encodedColumnPage.getPageDictionary(); } - encodedColumnPageList.add(encodedColumnPage); + this.encodedColumnPageList.add(encodedColumnPage); + this.columnName = encodedColumnPage.getActualPage().getColumnSpec().getFieldName(); return; } + // when first page was encoded without dictionary and next page encoded with dictionary + // in a blocklet + if (!isLocalDictEncoded && encodedColumnPage.isLocalDictGeneratedPage()) { + LOGGER.info( + "Local dictionary Fallback is initiated for column: " + this.columnName + " for page:" + + encodedColumnPageList.size()); + fallbackFutureQueue.add(fallbackExecutorService + .submit(new FallbackColumnPageEncoder(encodedColumnPage, encodedColumnPageList.size()))); + // fill null so once page is decoded again fill the re-encoded page again + this.encodedColumnPageList.add(null); + } // if local dictionary is false or column is encoded with local dictionary then // add a page - if (!isLocalDictEncoded || encodedColumnPage.isLocalDictGeneratedPage()) { - this.encodedColumnPageList.add(encodedColumnPage); + else if (!isLocalDictEncoded || encodedColumnPage.isLocalDictGeneratedPage()) { // merge page level dictionary values if (null != this.pageLevelDictionary) { pageLevelDictionary.mergerDictionaryValues(encodedColumnPage.getPageDictionary()); } - } else { - // if older pages were encoded with dictionary and new pages are without dictionary + this.encodedColumnPageList.add(encodedColumnPage); + } + // if all the older pages were encoded with dictionary and new pages are without dictionary + else { isLocalDictEncoded = false; pageLevelDictionary = null; - this.fallbackFutureQueue = new ArrayDeque<>(); - LOGGER.info( - "Local dictionary Fallback is initiated for column: " + encodedColumnPageList.get(0) - .getActualPage().getColumnSpec().getFieldName()); + LOGGER.warn("Local dictionary Fallback is initiated for column: " + this.columnName + + " for pages: 1 to " + encodedColumnPageList.size()); // submit all the older pages encoded with dictionary for fallback for (int pageIndex = 0; pageIndex < encodedColumnPageList.size(); pageIndex++) { - fallbackFutureQueue.add(fallbackExecutorService.submit( - new FallbackColumnPageEncoder(encodedColumnPageList.get(pageIndex), pageIndex))); + if (encodedColumnPageList.get(pageIndex).getActualPage().isLocalDictGeneratedPage()) { + fallbackFutureQueue.add(fallbackExecutorService.submit( + new FallbackColumnPageEncoder(encodedColumnPageList.get(pageIndex), pageIndex))); + } } //add to page list this.encodedColumnPageList.add(encodedColumnPage); http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/DimensionRawColumnChunk.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/DimensionRawColumnChunk.java b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/DimensionRawColumnChunk.java index d5320d6..0c0b6b7 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/DimensionRawColumnChunk.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/DimensionRawColumnChunk.java @@ -18,12 +18,23 @@ package org.apache.carbondata.core.datastore.chunk.impl; import java.io.IOException; import java.nio.ByteBuffer; +import java.util.BitSet; +import java.util.List; +import org.apache.carbondata.core.constants.CarbonCommonConstants; import org.apache.carbondata.core.datastore.FileReader; import org.apache.carbondata.core.datastore.chunk.AbstractRawColumnChunk; import org.apache.carbondata.core.datastore.chunk.DimensionColumnPage; import org.apache.carbondata.core.datastore.chunk.reader.DimensionColumnChunkReader; +import org.apache.carbondata.core.datastore.compression.CompressorFactory; +import org.apache.carbondata.core.datastore.page.ColumnPage; +import org.apache.carbondata.core.datastore.page.encoding.ColumnPageDecoder; +import org.apache.carbondata.core.datastore.page.encoding.DefaultEncodingFactory; import org.apache.carbondata.core.memory.MemoryException; +import org.apache.carbondata.core.scan.result.vector.CarbonDictionary; +import org.apache.carbondata.core.scan.result.vector.impl.CarbonDictionaryImpl; +import org.apache.carbondata.format.Encoding; +import org.apache.carbondata.format.LocalDictionaryChunk; /** * Contains raw dimension data, @@ -39,6 +50,8 @@ public class DimensionRawColumnChunk extends AbstractRawColumnChunk { private FileReader fileReader; + private CarbonDictionary localDictionary; + public DimensionRawColumnChunk(int columnIndex, ByteBuffer rawData, long offSet, int length, DimensionColumnChunkReader columnChunkReader) { super(columnIndex, rawData, offSet, length); @@ -126,4 +139,53 @@ public class DimensionRawColumnChunk extends AbstractRawColumnChunk { public FileReader getFileReader() { return fileReader; } + + public CarbonDictionary getLocalDictionary() { + if (null != getDataChunkV3().local_dictionary && null == localDictionary) { + try { + localDictionary = getDictionary(getDataChunkV3().local_dictionary); + } catch (IOException | MemoryException e) { + throw new RuntimeException(e); + } + } + return localDictionary; + } + + /** + * Below method will be used to get the local dictionary for a blocklet + * @param localDictionaryChunk + * local dictionary chunk thrift object + * @return local dictionary + * @throws IOException + * @throws MemoryException + */ + private CarbonDictionary getDictionary(LocalDictionaryChunk localDictionaryChunk) + throws IOException, MemoryException { + if (null != localDictionaryChunk) { + List<Encoding> encodings = localDictionaryChunk.getDictionary_meta().getEncoders(); + List<ByteBuffer> encoderMetas = localDictionaryChunk.getDictionary_meta().getEncoder_meta(); + ColumnPageDecoder decoder = + DefaultEncodingFactory.getInstance().createDecoder(encodings, encoderMetas); + ColumnPage decode = decoder.decode(localDictionaryChunk.getDictionary_data(), 0, + localDictionaryChunk.getDictionary_data().length); + BitSet usedDictionary = BitSet.valueOf(CompressorFactory.getInstance().getCompressor() + .unCompressByte(localDictionaryChunk.getDictionary_values())); + int length = usedDictionary.length(); + int index = 0; + byte[][] dictionary = new byte[length][]; + for (int i = 0; i < length; i++) { + if (usedDictionary.get(i)) { + dictionary[i] = decode.getBytes(index++); + } else { + dictionary[i] = null; + } + } + decode.freeMemory(); + // as dictionary values starts from 1 setting null default value + dictionary[1] = CarbonCommonConstants.MEMBER_DEFAULT_VAL_ARRAY; + return new CarbonDictionaryImpl(dictionary, usedDictionary.cardinality()); + } + return null; + } + } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/FixedLengthDimensionColumnPage.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/FixedLengthDimensionColumnPage.java b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/FixedLengthDimensionColumnPage.java index 570404a..e25fcb9 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/FixedLengthDimensionColumnPage.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/FixedLengthDimensionColumnPage.java @@ -47,7 +47,7 @@ public class FixedLengthDimensionColumnPage extends AbstractDimensionColumnPage dataChunk.length; dataChunkStore = DimensionChunkStoreFactory.INSTANCE .getDimensionChunkStore(columnValueSize, isExplicitSorted, numberOfRows, totalSize, - DimensionStoreType.FIXED_LENGTH); + DimensionStoreType.FIXED_LENGTH, null); dataChunkStore.putArray(invertedIndex, invertedIndexReverse, dataChunk); } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/VariableLengthDimensionColumnPage.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/VariableLengthDimensionColumnPage.java b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/VariableLengthDimensionColumnPage.java index 7394217..aea594c 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/VariableLengthDimensionColumnPage.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/impl/VariableLengthDimensionColumnPage.java @@ -21,6 +21,7 @@ import org.apache.carbondata.core.datastore.chunk.store.DimensionChunkStoreFacto import org.apache.carbondata.core.datastore.chunk.store.DimensionChunkStoreFactory.DimensionStoreType; import org.apache.carbondata.core.scan.executor.infos.KeyStructureInfo; import org.apache.carbondata.core.scan.result.vector.CarbonColumnVector; +import org.apache.carbondata.core.scan.result.vector.CarbonDictionary; import org.apache.carbondata.core.scan.result.vector.ColumnVectorInfo; /** @@ -32,14 +33,29 @@ public class VariableLengthDimensionColumnPage extends AbstractDimensionColumnPa * Constructor for this class */ public VariableLengthDimensionColumnPage(byte[] dataChunks, int[] invertedIndex, - int[] invertedIndexReverse, int numberOfRows, DimensionStoreType dimStoreType) { + int[] invertedIndexReverse, int numberOfRows, DimensionStoreType dimStoreType, + CarbonDictionary dictionary) { boolean isExplicitSorted = isExplicitSorted(invertedIndex); - long totalSize = null != invertedIndex ? - (dataChunks.length + (2 * numberOfRows * CarbonCommonConstants.INT_SIZE_IN_BYTE) + ( - numberOfRows * CarbonCommonConstants.INT_SIZE_IN_BYTE)) : - (dataChunks.length + (numberOfRows * CarbonCommonConstants.INT_SIZE_IN_BYTE)); + long totalSize = 0; + switch (dimStoreType) { + case LOCAL_DICT: + totalSize = null != invertedIndex ? + (dataChunks.length + (2 * numberOfRows * CarbonCommonConstants.INT_SIZE_IN_BYTE)) : + dataChunks.length; + break; + case VARIABLE_INT_LENGTH: + case VARIABLE_SHORT_LENGTH: + totalSize = null != invertedIndex ? + (dataChunks.length + (2 * numberOfRows * CarbonCommonConstants.INT_SIZE_IN_BYTE) + ( + numberOfRows * CarbonCommonConstants.INT_SIZE_IN_BYTE)) : + (dataChunks.length + (numberOfRows * CarbonCommonConstants.INT_SIZE_IN_BYTE)); + break; + default: + throw new UnsupportedOperationException("Invalidate dimension store type"); + } dataChunkStore = DimensionChunkStoreFactory.INSTANCE - .getDimensionChunkStore(0, isExplicitSorted, numberOfRows, totalSize, dimStoreType); + .getDimensionChunkStore(0, isExplicitSorted, numberOfRows, totalSize, dimStoreType, + dictionary); dataChunkStore.putArray(invertedIndex, invertedIndexReverse, dataChunks); } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v1/CompressedDimensionChunkFileBasedReaderV1.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v1/CompressedDimensionChunkFileBasedReaderV1.java b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v1/CompressedDimensionChunkFileBasedReaderV1.java index 92a9684..593d992 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v1/CompressedDimensionChunkFileBasedReaderV1.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v1/CompressedDimensionChunkFileBasedReaderV1.java @@ -152,7 +152,8 @@ public class CompressedDimensionChunkFileBasedReaderV1 extends AbstractChunkRead .hasEncoding(dataChunk.getEncodingList(), Encoding.DICTIONARY)) { columnDataChunk = new VariableLengthDimensionColumnPage(dataPage, invertedIndexes, invertedIndexesReverse, - numberOfRows, DimensionChunkStoreFactory.DimensionStoreType.VARIABLE_SHORT_LENGTH); + numberOfRows, DimensionChunkStoreFactory.DimensionStoreType.VARIABLE_SHORT_LENGTH, + null); } else { // to store fixed length column chunk values columnDataChunk = http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v2/CompressedDimensionChunkFileBasedReaderV2.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v2/CompressedDimensionChunkFileBasedReaderV2.java b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v2/CompressedDimensionChunkFileBasedReaderV2.java index 3cdbe1d..252a675 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v2/CompressedDimensionChunkFileBasedReaderV2.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v2/CompressedDimensionChunkFileBasedReaderV2.java @@ -176,7 +176,8 @@ public class CompressedDimensionChunkFileBasedReaderV2 extends AbstractChunkRead if (!hasEncoding(dimensionColumnChunk.encoders, Encoding.DICTIONARY)) { columnDataChunk = new VariableLengthDimensionColumnPage(dataPage, invertedIndexes, invertedIndexesReverse, - numberOfRows, DimensionChunkStoreFactory.DimensionStoreType.VARIABLE_SHORT_LENGTH); + numberOfRows, DimensionChunkStoreFactory.DimensionStoreType.VARIABLE_SHORT_LENGTH, + null); } else { // to store fixed length column chunk values columnDataChunk = http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v3/CompressedDimensionChunkFileBasedReaderV3.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v3/CompressedDimensionChunkFileBasedReaderV3.java b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v3/CompressedDimensionChunkFileBasedReaderV3.java index fee114d..eb5917b 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v3/CompressedDimensionChunkFileBasedReaderV3.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/reader/dimension/v3/CompressedDimensionChunkFileBasedReaderV3.java @@ -234,7 +234,7 @@ public class CompressedDimensionChunkFileBasedReaderV3 extends AbstractChunkRead throws IOException, MemoryException { if (isEncodedWithMeta(pageMetadata)) { ColumnPage decodedPage = decodeDimensionByMeta(pageMetadata, pageData, offset); - return new ColumnPageWrapper(decodedPage); + return new ColumnPageWrapper(decodedPage, rawColumnPage.getLocalDictionary()); } else { // following code is for backward compatibility return decodeDimensionLegacy(rawColumnPage, pageData, pageMetadata, offset); @@ -265,21 +265,25 @@ public class CompressedDimensionChunkFileBasedReaderV3 extends AbstractChunkRead CarbonUtil.getIntArray(pageData, offset, pageMetadata.rle_page_length); // uncompress the data with rle indexes dataPage = UnBlockIndexer.uncompressData(dataPage, rlePage, - eachColumnValueSize[rawColumnPage.getColumnIndex()]); + null == rawColumnPage.getLocalDictionary() ? + eachColumnValueSize[rawColumnPage.getColumnIndex()] : + 3); } DimensionColumnPage columnDataChunk = null; - // if no dictionary column then first create a no dictionary column chunk // and set to data chunk instance if (!hasEncoding(pageMetadata.encoders, Encoding.DICTIONARY)) { DimensionChunkStoreFactory.DimensionStoreType dimStoreType = - hasEncoding(pageMetadata.encoders, Encoding.DIRECT_COMPRESS_VARCHAR) ? - DimensionChunkStoreFactory.DimensionStoreType.VARIABLE_INT_LENGTH : - DimensionChunkStoreFactory.DimensionStoreType.VARIABLE_SHORT_LENGTH; + null != rawColumnPage.getLocalDictionary() ? + DimensionChunkStoreFactory.DimensionStoreType.LOCAL_DICT : + (hasEncoding(pageMetadata.encoders, Encoding.DIRECT_COMPRESS_VARCHAR) ? + DimensionChunkStoreFactory.DimensionStoreType.VARIABLE_INT_LENGTH : + DimensionChunkStoreFactory.DimensionStoreType.VARIABLE_SHORT_LENGTH); columnDataChunk = new VariableLengthDimensionColumnPage(dataPage, invertedIndexes, invertedIndexesReverse, - pageMetadata.getNumberOfRowsInpage(), dimStoreType); + pageMetadata.getNumberOfRowsInpage(), dimStoreType, + rawColumnPage.getLocalDictionary()); } else { // to store fixed length column chunk values columnDataChunk = http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/ColumnPageWrapper.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/ColumnPageWrapper.java b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/ColumnPageWrapper.java index c89ecc3..59e593a 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/ColumnPageWrapper.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/ColumnPageWrapper.java @@ -20,14 +20,19 @@ package org.apache.carbondata.core.datastore.chunk.store; import org.apache.carbondata.core.datastore.chunk.DimensionColumnPage; import org.apache.carbondata.core.datastore.page.ColumnPage; import org.apache.carbondata.core.scan.executor.infos.KeyStructureInfo; +import org.apache.carbondata.core.scan.result.vector.CarbonDictionary; import org.apache.carbondata.core.scan.result.vector.ColumnVectorInfo; +import org.apache.carbondata.core.util.CarbonUtil; public class ColumnPageWrapper implements DimensionColumnPage { private ColumnPage columnPage; - public ColumnPageWrapper(ColumnPage columnPage) { + private CarbonDictionary localDictionary; + + public ColumnPageWrapper(ColumnPage columnPage, CarbonDictionary localDictionary) { this.columnPage = columnPage; + this.localDictionary = localDictionary; } @Override @@ -55,6 +60,10 @@ public class ColumnPageWrapper implements DimensionColumnPage { @Override public byte[] getChunkData(int rowId) { + if (null != localDictionary) { + return localDictionary.getDictionaryValue(CarbonUtil + .getSurrogateInternal(columnPage.getBytes(rowId), 0, 3)); + } return columnPage.getBytes(rowId); } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/DimensionChunkStoreFactory.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/DimensionChunkStoreFactory.java b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/DimensionChunkStoreFactory.java index eccfd9c..c7bcef1 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/DimensionChunkStoreFactory.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/DimensionChunkStoreFactory.java @@ -18,12 +18,14 @@ package org.apache.carbondata.core.datastore.chunk.store; import org.apache.carbondata.core.constants.CarbonCommonConstants; +import org.apache.carbondata.core.datastore.chunk.store.impl.LocalDictDimensionDataChunkStore; import org.apache.carbondata.core.datastore.chunk.store.impl.safe.SafeFixedLengthDimensionDataChunkStore; import org.apache.carbondata.core.datastore.chunk.store.impl.safe.SafeVariableIntLengthDimensionDataChunkStore; import org.apache.carbondata.core.datastore.chunk.store.impl.safe.SafeVariableShortLengthDimensionDataChunkStore; import org.apache.carbondata.core.datastore.chunk.store.impl.unsafe.UnsafeFixedLengthDimensionDataChunkStore; import org.apache.carbondata.core.datastore.chunk.store.impl.unsafe.UnsafeVariableIntLengthDimensionDataChunkStore; import org.apache.carbondata.core.datastore.chunk.store.impl.unsafe.UnsafeVariableShortLengthDimensionDataChunkStore; +import org.apache.carbondata.core.scan.result.vector.CarbonDictionary; import org.apache.carbondata.core.util.CarbonProperties; /** @@ -62,26 +64,41 @@ public class DimensionChunkStoreFactory { * @return dimension store type */ public DimensionDataChunkStore getDimensionChunkStore(int columnValueSize, - boolean isInvertedIndex, int numberOfRows, long totalSize, DimensionStoreType storeType) { - + boolean isInvertedIndex, int numberOfRows, long totalSize, DimensionStoreType storeType, + CarbonDictionary dictionary) { if (isUnsafe) { - if (storeType == DimensionStoreType.FIXED_LENGTH) { - return new UnsafeFixedLengthDimensionDataChunkStore(totalSize, columnValueSize, - isInvertedIndex, numberOfRows); - } else if (storeType == DimensionStoreType.VARIABLE_SHORT_LENGTH) { - return new UnsafeVariableShortLengthDimensionDataChunkStore(totalSize, isInvertedIndex, - numberOfRows); - } else { - return new UnsafeVariableIntLengthDimensionDataChunkStore(totalSize, isInvertedIndex, - numberOfRows); + switch (storeType) { + case FIXED_LENGTH: + return new UnsafeFixedLengthDimensionDataChunkStore(totalSize, columnValueSize, + isInvertedIndex, numberOfRows); + case VARIABLE_SHORT_LENGTH: + return new UnsafeVariableShortLengthDimensionDataChunkStore(totalSize, isInvertedIndex, + numberOfRows); + case VARIABLE_INT_LENGTH: + return new UnsafeVariableIntLengthDimensionDataChunkStore(totalSize, isInvertedIndex, + numberOfRows); + case LOCAL_DICT: + return new LocalDictDimensionDataChunkStore( + new UnsafeFixedLengthDimensionDataChunkStore(totalSize, + 3, isInvertedIndex, numberOfRows), + dictionary); + default: + throw new UnsupportedOperationException("Invalid dimension store type"); } } else { - if (storeType == DimensionStoreType.FIXED_LENGTH) { - return new SafeFixedLengthDimensionDataChunkStore(isInvertedIndex, columnValueSize); - } else if (storeType == DimensionStoreType.VARIABLE_SHORT_LENGTH) { - return new SafeVariableShortLengthDimensionDataChunkStore(isInvertedIndex, numberOfRows); - } else { - return new SafeVariableIntLengthDimensionDataChunkStore(isInvertedIndex, numberOfRows); + switch (storeType) { + case FIXED_LENGTH: + return new SafeFixedLengthDimensionDataChunkStore(isInvertedIndex, columnValueSize); + case VARIABLE_SHORT_LENGTH: + return new SafeVariableShortLengthDimensionDataChunkStore(isInvertedIndex, numberOfRows); + case VARIABLE_INT_LENGTH: + return new SafeVariableIntLengthDimensionDataChunkStore(isInvertedIndex, numberOfRows); + case LOCAL_DICT: + return new LocalDictDimensionDataChunkStore( + new SafeFixedLengthDimensionDataChunkStore(isInvertedIndex, + 3), dictionary); + default: + throw new UnsupportedOperationException("Invalid dimension store type"); } } } @@ -90,6 +107,6 @@ public class DimensionChunkStoreFactory { * dimension store type enum */ public enum DimensionStoreType { - FIXED_LENGTH, VARIABLE_SHORT_LENGTH, VARIABLE_INT_LENGTH; + FIXED_LENGTH, VARIABLE_SHORT_LENGTH, VARIABLE_INT_LENGTH, LOCAL_DICT; } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/impl/LocalDictDimensionDataChunkStore.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/impl/LocalDictDimensionDataChunkStore.java b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/impl/LocalDictDimensionDataChunkStore.java new file mode 100644 index 0000000..a6bd23d --- /dev/null +++ b/core/src/main/java/org/apache/carbondata/core/datastore/chunk/store/impl/LocalDictDimensionDataChunkStore.java @@ -0,0 +1,109 @@ +/* + * 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.datastore.chunk.store.impl; + +import org.apache.carbondata.core.constants.CarbonCommonConstants; +import org.apache.carbondata.core.datastore.chunk.store.DimensionDataChunkStore; +import org.apache.carbondata.core.scan.result.vector.CarbonColumnVector; +import org.apache.carbondata.core.scan.result.vector.CarbonDictionary; + +/** + * Dimension chunk store for local dictionary encoded data. + * It's a decorator over dimension chunk store + */ +public class LocalDictDimensionDataChunkStore implements DimensionDataChunkStore { + + private DimensionDataChunkStore dimensionDataChunkStore; + + private CarbonDictionary dictionary; + + public LocalDictDimensionDataChunkStore(DimensionDataChunkStore dimensionDataChunkStore, + CarbonDictionary dictionary) { + this.dimensionDataChunkStore = dimensionDataChunkStore; + this.dictionary = dictionary; + } + + /** + * Below method will be used to put the rows and its metadata in offheap + * + * @param invertedIndex inverted index to be stored + * @param invertedIndexReverse inverted index reverse to be stored + * @param data data to be stored + */ + public void putArray(int[] invertedIndex, int[] invertedIndexReverse, byte[] data) { + this.dimensionDataChunkStore.putArray(invertedIndex, invertedIndexReverse, data); + } + + @Override public byte[] getRow(int rowId) { + return dictionary.getDictionaryValue(dimensionDataChunkStore.getSurrogate(rowId)); + } + + @Override public void fillRow(int rowId, CarbonColumnVector vector, int vectorRow) { + if (!dictionary.isDictionaryUsed()) { + vector.setDictionary(dictionary); + dictionary.setDictionaryUsed(); + } + int surrogate = dimensionDataChunkStore.getSurrogate(rowId); + if (surrogate == CarbonCommonConstants.MEMBER_DEFAULT_VAL_SURROGATE_KEY) { + vector.putNull(rowId); + vector.getDictionaryVector().putNull(rowId); + return; + } + vector.putNotNull(vectorRow); + vector.getDictionaryVector().putInt(vectorRow, dimensionDataChunkStore.getSurrogate(rowId)); + } + + @Override public void fillRow(int rowId, byte[] buffer, int offset) { + throw new UnsupportedOperationException("Operation not supported"); + } + + @Override public int getInvertedIndex(int rowId) { + return this.dimensionDataChunkStore.getInvertedIndex(rowId); + } + + @Override public int getInvertedReverseIndex(int rowId) { + return this.dimensionDataChunkStore.getInvertedReverseIndex(rowId); + } + + @Override public int getSurrogate(int rowId) { + throw new UnsupportedOperationException("Operation not supported"); + } + + @Override public int getColumnValueSize() { + throw new UnsupportedOperationException("Operation not supported"); + } + + @Override public boolean isExplicitSorted() { + return this.dimensionDataChunkStore.isExplicitSorted(); + } + + @Override public int compareTo(int rowId, byte[] compareValue) { + return dimensionDataChunkStore.compareTo(rowId, compareValue); + } + + /** + * Below method will be used to free the memory occupied by the column chunk + */ + @Override public void freeMemory() { + if (null != dimensionDataChunkStore) { + this.dimensionDataChunkStore.freeMemory(); + this.dictionary = null; + this.dimensionDataChunkStore = null; + } + } +} http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/page/ColumnPage.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/page/ColumnPage.java b/core/src/main/java/org/apache/carbondata/core/datastore/page/ColumnPage.java index 4ff1330..6077ddf 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/page/ColumnPage.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/page/ColumnPage.java @@ -163,15 +163,16 @@ public abstract class ColumnPage { } public static ColumnPage newLocalDictPage(TableSpec.ColumnSpec columnSpec, DataType dataType, - int pageSize, LocalDictionaryGenerator localDictionaryGenerator) throws MemoryException { + int pageSize, LocalDictionaryGenerator localDictionaryGenerator, + boolean isComplexTypePrimitive) throws MemoryException { if (unsafe) { return new LocalDictColumnPage(new UnsafeVarLengthColumnPage(columnSpec, dataType, pageSize), new UnsafeVarLengthColumnPage(columnSpec, DataTypes.BYTE_ARRAY, pageSize), - localDictionaryGenerator); + localDictionaryGenerator, isComplexTypePrimitive); } else { return new LocalDictColumnPage(new SafeVarLengthColumnPage(columnSpec, dataType, pageSize), new SafeVarLengthColumnPage(columnSpec, DataTypes.BYTE_ARRAY, pageSize), - localDictionaryGenerator); + localDictionaryGenerator, isComplexTypePrimitive); } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/page/ComplexColumnPage.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/page/ComplexColumnPage.java b/core/src/main/java/org/apache/carbondata/core/datastore/page/ComplexColumnPage.java index c6b650f..a170c8b 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/page/ComplexColumnPage.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/page/ComplexColumnPage.java @@ -83,7 +83,7 @@ public class ComplexColumnPage { TableSpec.ColumnSpec spec = TableSpec.ColumnSpec .newInstance(columnNames.get(i), DataTypes.BYTE_ARRAY, complexColumnType.get(i)); this.columnPages[i] = ColumnPage - .newLocalDictPage(spec, DataTypes.BYTE_ARRAY, pageSize, localDictionaryGenerator); + .newLocalDictPage(spec, DataTypes.BYTE_ARRAY, pageSize, localDictionaryGenerator, true); this.columnPages[i].setStatsCollector(new DummyStatsCollector()); } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/page/LocalDictColumnPage.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/page/LocalDictColumnPage.java b/core/src/main/java/org/apache/carbondata/core/datastore/page/LocalDictColumnPage.java index 2c7d3a7..a072852 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/page/LocalDictColumnPage.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/page/LocalDictColumnPage.java @@ -22,10 +22,13 @@ import java.math.BigDecimal; 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.keygenerator.KeyGenException; +import org.apache.carbondata.core.keygenerator.KeyGenerator; +import org.apache.carbondata.core.keygenerator.factory.KeyGeneratorFactory; import org.apache.carbondata.core.localdictionary.PageLevelDictionary; import org.apache.carbondata.core.localdictionary.exception.DictionaryThresholdReachedException; import org.apache.carbondata.core.localdictionary.generator.LocalDictionaryGenerator; -import org.apache.carbondata.core.util.ByteUtil; /** * Column page implementation for Local dictionary generated columns @@ -61,19 +64,26 @@ public class LocalDictColumnPage extends ColumnPage { */ private boolean isActualPageMemoryFreed; + private KeyGenerator keyGenerator; + + private int[] dummyKey; /** * Create a new column page with input data type and page size. */ protected LocalDictColumnPage(ColumnPage actualDataColumnPage, ColumnPage encodedColumnpage, - LocalDictionaryGenerator localDictionaryGenerator) { + LocalDictionaryGenerator localDictionaryGenerator, boolean isComplexTypePrimitive) { super(actualDataColumnPage.getColumnSpec(), actualDataColumnPage.getDataType(), actualDataColumnPage.getPageSize()); // if threshold is not reached then create page level dictionary // for encoding with local dictionary if (!localDictionaryGenerator.isThresholdReached()) { pageLevelDictionary = new PageLevelDictionary(localDictionaryGenerator, - actualDataColumnPage.getColumnSpec().getFieldName(), actualDataColumnPage.getDataType()); + actualDataColumnPage.getColumnSpec().getFieldName(), actualDataColumnPage.getDataType(), + isComplexTypePrimitive); this.encodedDataColumnPage = encodedColumnpage; + this.keyGenerator = KeyGeneratorFactory + .getKeyGenerator(new int[] { CarbonCommonConstants.LOCAL_DICTIONARY_MAX + 1 }); + this.dummyKey = new int[1]; } else { // else free the encoded column page memory as its of no use encodedColumnpage.freeMemory(); @@ -109,14 +119,18 @@ public class LocalDictColumnPage extends ColumnPage { if (null != pageLevelDictionary) { try { actualDataColumnPage.putBytes(rowId, bytes); - int dictionaryValue = pageLevelDictionary.getDictionaryValue(bytes); - encodedDataColumnPage.putBytes(rowId, ByteUtil.toBytes(dictionaryValue)); + dummyKey[0] = pageLevelDictionary.getDictionaryValue(bytes); + encodedDataColumnPage.putBytes(rowId, keyGenerator.generateKey(dummyKey)); } catch (DictionaryThresholdReachedException e) { LOGGER.error(e, "Local Dictionary threshold reached for the column: " + actualDataColumnPage .getColumnSpec().getFieldName()); pageLevelDictionary = null; encodedDataColumnPage.freeMemory(); encodedDataColumnPage = null; + } catch (KeyGenException e) { + LOGGER.error(e, "Unable to generate key for: " + actualDataColumnPage + .getColumnSpec().getFieldName()); + throw new RuntimeException(e); } } else { actualDataColumnPage.putBytes(rowId, bytes); http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/datastore/page/VarLengthColumnPageBase.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/datastore/page/VarLengthColumnPageBase.java b/core/src/main/java/org/apache/carbondata/core/datastore/page/VarLengthColumnPageBase.java index bd49b94..fa163bc 100644 --- a/core/src/main/java/org/apache/carbondata/core/datastore/page/VarLengthColumnPageBase.java +++ b/core/src/main/java/org/apache/carbondata/core/datastore/page/VarLengthColumnPageBase.java @@ -424,7 +424,7 @@ public abstract class VarLengthColumnPageBase extends ColumnPage { int offset = 0; byte[] data = new byte[totalLength]; for (int rowId = 0; rowId < rowOffset.size() - 1; rowId++) { - short length = (short) (rowOffset.get(rowId + 1) - rowOffset.get(rowId)); + int length = (rowOffset.get(rowId + 1) - rowOffset.get(rowId)); copyBytes(rowId, data, offset, length); offset += length; } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/localdictionary/PageLevelDictionary.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/localdictionary/PageLevelDictionary.java b/core/src/main/java/org/apache/carbondata/core/localdictionary/PageLevelDictionary.java index db95b03..10a4e6e 100644 --- a/core/src/main/java/org/apache/carbondata/core/localdictionary/PageLevelDictionary.java +++ b/core/src/main/java/org/apache/carbondata/core/localdictionary/PageLevelDictionary.java @@ -17,8 +17,10 @@ package org.apache.carbondata.core.localdictionary; import java.io.IOException; +import java.nio.ByteBuffer; import java.util.BitSet; +import org.apache.carbondata.core.constants.CarbonCommonConstants; import org.apache.carbondata.core.datastore.ColumnType; import org.apache.carbondata.core.datastore.TableSpec; import org.apache.carbondata.core.datastore.compression.CompressorFactory; @@ -54,12 +56,15 @@ public class PageLevelDictionary { private DataType dataType; + private boolean isComplexTypePrimitive; + public PageLevelDictionary(LocalDictionaryGenerator localDictionaryGenerator, String columnName, - DataType dataType) { + DataType dataType, boolean isComplexTypePrimitive) { this.localDictionaryGenerator = localDictionaryGenerator; this.usedDictionaryValues = new BitSet(); this.columnName = columnName; this.dataType = dataType; + this.isComplexTypePrimitive = isComplexTypePrimitive; } /** @@ -97,8 +102,12 @@ public class PageLevelDictionary { throws MemoryException, IOException { // TODO support for actual data type dictionary ColumnSPEC ColumnType columnType = ColumnType.PLAIN_VALUE; + boolean isVarcharType = false; + int lvSize = CarbonCommonConstants.SHORT_SIZE_IN_BYTE; if (DataTypes.VARCHAR == dataType) { columnType = ColumnType.PLAIN_LONG_VALUE; + lvSize = CarbonCommonConstants.INT_SIZE_IN_BYTE; + isVarcharType = true; } TableSpec.ColumnSpec spec = TableSpec.ColumnSpec.newInstance(columnName, DataTypes.BYTE_ARRAY, columnType); @@ -107,10 +116,23 @@ public class PageLevelDictionary { // TODO support data type specific stats collector for numeric data types dictionaryColumnPage.setStatsCollector(new DummyStatsCollector()); int rowId = 0; + ByteBuffer byteBuffer = null; for (int i = usedDictionaryValues.nextSetBit(0); i >= 0; i = usedDictionaryValues.nextSetBit(i + 1)) { - dictionaryColumnPage - .putData(rowId++, localDictionaryGenerator.getDictionaryKeyBasedOnValue(i)); + if (!isComplexTypePrimitive) { + dictionaryColumnPage + .putData(rowId++, localDictionaryGenerator.getDictionaryKeyBasedOnValue(i)); + } else { + byte[] dictionaryKeyBasedOnValue = localDictionaryGenerator.getDictionaryKeyBasedOnValue(i); + byteBuffer = ByteBuffer.allocate(lvSize + dictionaryKeyBasedOnValue.length); + if (!isVarcharType) { + byteBuffer.putShort((short) dictionaryKeyBasedOnValue.length); + } else { + byteBuffer.putInt(dictionaryKeyBasedOnValue.length); + } + byteBuffer.put(dictionaryKeyBasedOnValue); + dictionaryColumnPage.putData(rowId++, byteBuffer.array()); + } } // creating a encoder ColumnPageEncoder encoder = new DirectCompressCodec(DataTypes.BYTE_ARRAY).createEncoder(null); http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/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 index 8aa2e19..02ecb7f 100644 --- 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 @@ -33,6 +33,8 @@ public class ColumnLocalDictionaryGenerator implements LocalDictionaryGenerator */ private DictionaryStore dictionaryHolder; + private long currentSize; + public ColumnLocalDictionaryGenerator(int threshold, int lvLength) { // adding 1 to threshold for null value int newThreshold = threshold + 1; @@ -52,6 +54,7 @@ public class ColumnLocalDictionaryGenerator implements LocalDictionaryGenerator } catch (DictionaryThresholdReachedException e) { // do nothing } + currentSize += byteBuffer.array().length; } /** @@ -61,8 +64,12 @@ public class ColumnLocalDictionaryGenerator implements LocalDictionaryGenerator * @return dictionary value */ @Override public int generateDictionary(byte[] data) throws DictionaryThresholdReachedException { - int dictionaryValue = this.dictionaryHolder.putIfAbsent(data); - return dictionaryValue; + currentSize += data.length; + if (currentSize >= Integer.MAX_VALUE) { + throw new DictionaryThresholdReachedException( + "Unable to generate dictionary as Dictionary Size crossed 2GB limit"); + } + return this.dictionaryHolder.putIfAbsent(data); } /** http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/filter/FilterUtil.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/filter/FilterUtil.java b/core/src/main/java/org/apache/carbondata/core/scan/filter/FilterUtil.java index a657b91..c3e3b78 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/filter/FilterUtil.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/filter/FilterUtil.java @@ -1845,22 +1845,36 @@ public final class FilterUtil { } } + /** + * Below method will be called from include and exclude filter to convert filter values + * based on dictionary when local dictionary is present in blocklet. + * @param dictionary + * Dictionary + * @param actualFilterValues + * actual filter values + * @return encoded filter values + */ public static byte[][] getEncodedFilterValues(CarbonDictionary dictionary, byte[][] actualFilterValues) { if (null == dictionary) { return actualFilterValues; } - KeyGenerator keyGenerator = KeyGeneratorFactory.getKeyGenerator(new int[] { 100000 }); + KeyGenerator keyGenerator = KeyGeneratorFactory + .getKeyGenerator(new int[] { CarbonCommonConstants.LOCAL_DICTIONARY_MAX }); + int[] dummy = new int[1]; List<byte[]> encodedFilters = new ArrayList<>(); for (byte[] actualFilter : actualFilterValues) { - for (int i = 1; i < dictionary.getDictionaryValues().length; i++) { + for (int i = 1; i < dictionary.getDictionarySize(); i++) { + if (dictionary.getDictionaryValue(i) == null) { + continue; + } if (ByteUtil.UnsafeComparer.INSTANCE - .compareTo(actualFilter, dictionary.getDictionaryValues()[i]) - == 0) { + .compareTo(actualFilter, dictionary.getDictionaryValue(i)) == 0) { try { - encodedFilters.add(keyGenerator.generateKey(new int[] { i })); + dummy[0] = i; + encodedFilters.add(keyGenerator.generateKey(dummy)); } catch (KeyGenException e) { - //do nothing + LOGGER.error(e); } break; } @@ -1869,6 +1883,13 @@ public final class FilterUtil { return getSortedEncodedFilters(encodedFilters); } + /** + * Below method will be used to sort the filter values a filter are applied using incremental + * binary search + * @param encodedFilters + * encoded filter values + * @return sorted encoded filter values + */ private static byte[][] getSortedEncodedFilters(List<byte[]> encodedFilters) { java.util.Comparator<byte[]> filterNoDictValueComaparator = new java.util.Comparator<byte[]>() { @Override public int compare(byte[] filterMember1, byte[] filterMember2) { @@ -1879,15 +1900,28 @@ public final class FilterUtil { return encodedFilters.toArray(new byte[encodedFilters.size()][]); } - private static BitSet getIncludeDictionaryValues(Expression expression, + /** + * Below method will be used to get all the include filter values in case of range filters when + * blocklet is encoded with local dictionary + * @param expression + * filter expression + * @param dictionary + * dictionary + * @return include filter bitset + * @throws FilterUnsupportedException + */ + private static BitSet getIncludeDictFilterValuesForRange(Expression expression, CarbonDictionary dictionary) throws FilterUnsupportedException { ConditionalExpression conExp = (ConditionalExpression) expression; ColumnExpression columnExpression = conExp.getColumnList().get(0); BitSet includeFilterBitSet = new BitSet(); - for (int i = 2; i < dictionary.getDictionaryValues().length; i++) { + for (int i = 2; i < dictionary.getDictionarySize(); i++) { + if (null == dictionary.getDictionaryValue(i)) { + continue; + } try { RowIntf row = new RowImpl(); - String stringValue = new String(dictionary.getDictionaryValues()[i], + String stringValue = new String(dictionary.getDictionaryValue(i), Charset.forName(CarbonCommonConstants.DEFAULT_CHARSET)); row.setValues(new Object[] { DataTypeUtil.getDataBasedOnDataType(stringValue, columnExpression.getCarbonColumn().getDataType()) }); @@ -1904,8 +1938,21 @@ public final class FilterUtil { return includeFilterBitSet; } - public static byte[][] getEncodedFilterValues(BitSet includeDictValues, int dictSize, - boolean useExclude) { + /** + * Below method will used to get encoded filter values for range filter values + * when local dictionary is present in blocklet for columns + * If number of include filter is more than 60% of total dictionary size it will + * convert include to exclude + * @param includeDictValues + * include filter values + * @param carbonDictionary + * dictionary + * @param useExclude + * to check if using exclude will be more optimized + * @return encoded filter values + */ + private static byte[][] getEncodedFilterValuesForRange(BitSet includeDictValues, + CarbonDictionary carbonDictionary, boolean useExclude) { KeyGenerator keyGenerator = KeyGeneratorFactory .getKeyGenerator(new int[] { CarbonCommonConstants.LOCAL_DICTIONARY_MAX }); List<byte[]> encodedFilterValues = new ArrayList<>(); @@ -1918,38 +1965,50 @@ public final class FilterUtil { encodedFilterValues.add(keyGenerator.generateKey(dummy)); } } catch (KeyGenException e) { - // do nothing + LOGGER.error(e); } return encodedFilterValues.toArray(new byte[encodedFilterValues.size()][]); } else { try { - for (int i = 1; i < dictSize; i++) { - if (!includeDictValues.get(i)) { + for (int i = 1; i < carbonDictionary.getDictionarySize(); i++) { + if (!includeDictValues.get(i) && null != carbonDictionary.getDictionaryValue(i)) { dummy[0] = i; encodedFilterValues.add(keyGenerator.generateKey(dummy)); } } } catch (KeyGenException e) { - // do nothing + LOGGER.error(e); } } return getSortedEncodedFilters(encodedFilterValues); } - public static FilterExecuter getFilterExecutorForLocalDictionary( + /** + * Below method will be used to get filter executor instance for range filters + * when local dictonary is present for in blocklet + * @param rawColumnChunk + * raw column chunk + * @param exp + * filter expression + * @param isNaturalSorted + * is data was already sorted + * @return + */ + public static FilterExecuter getFilterExecutorForRangeFilters( DimensionRawColumnChunk rawColumnChunk, Expression exp, boolean isNaturalSorted) { BitSet includeDictionaryValues; try { includeDictionaryValues = - FilterUtil.getIncludeDictionaryValues(exp, rawColumnChunk.getLocalDictionary()); + FilterUtil.getIncludeDictFilterValuesForRange(exp, rawColumnChunk.getLocalDictionary()); } catch (FilterUnsupportedException e) { throw new RuntimeException(e); } boolean isExclude = includeDictionaryValues.cardinality() > 1 && FilterUtil - .isExcludeFilterNeedsToApply(rawColumnChunk.getLocalDictionary().getDictionarySize(), + .isExcludeFilterNeedsToApply(rawColumnChunk.getLocalDictionary().getDictionaryActualSize(), includeDictionaryValues.cardinality()); - byte[][] encodedFilterValues = FilterUtil.getEncodedFilterValues(includeDictionaryValues, - rawColumnChunk.getLocalDictionary().getDictionaryValues().length, isExclude); + byte[][] encodedFilterValues = FilterUtil + .getEncodedFilterValuesForRange(includeDictionaryValues, + rawColumnChunk.getLocalDictionary(), isExclude); FilterExecuter filterExecuter; if (!isExclude) { filterExecuter = new IncludeFilterExecuterImpl(encodedFilterValues, isNaturalSorted); http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/ExcludeFilterExecuterImpl.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/ExcludeFilterExecuterImpl.java b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/ExcludeFilterExecuterImpl.java index 7646550..71646c9 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/ExcludeFilterExecuterImpl.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/ExcludeFilterExecuterImpl.java @@ -297,8 +297,9 @@ public class ExcludeFilterExecuterImpl implements FilterExecuter { protected BitSet getFilteredIndexes(DimensionColumnPage dimensionColumnPage, int numberOfRows, boolean useBitsetPipeLine, BitSetGroup prvBitSetGroup, int pageNumber) { // check whether applying filtered based on previous bitset will be optimal - if (CarbonUtil.usePreviousFilterBitsetGroup(useBitsetPipeLine, prvBitSetGroup, pageNumber, - filterValues.length)) { + if (filterValues.length > 0 && CarbonUtil + .usePreviousFilterBitsetGroup(useBitsetPipeLine, prvBitSetGroup, pageNumber, + filterValues.length)) { return getFilteredIndexesUisngPrvBitset(dimensionColumnPage, prvBitSetGroup, pageNumber); } else { return getFilteredIndexes(dimensionColumnPage, numberOfRows); @@ -366,6 +367,9 @@ public class ExcludeFilterExecuterImpl implements FilterExecuter { DimensionColumnPage dimensionColumnPage, int numerOfRows) { BitSet bitSet = new BitSet(numerOfRows); bitSet.flip(0, numerOfRows); + if (filterValues.length == 0) { + return bitSet; + } int startIndex = 0; for (int i = 0; i < filterValues.length; i++) { if (startIndex >= numerOfRows) { http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/IncludeFilterExecuterImpl.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/IncludeFilterExecuterImpl.java b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/IncludeFilterExecuterImpl.java index 96ff7af..6af969e 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/IncludeFilterExecuterImpl.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/IncludeFilterExecuterImpl.java @@ -318,8 +318,9 @@ public class IncludeFilterExecuterImpl implements FilterExecuter { protected BitSet getFilteredIndexes(DimensionColumnPage dimensionColumnPage, int numberOfRows, boolean useBitsetPipeLine, BitSetGroup prvBitSetGroup, int pageNumber) { // check whether previous indexes can be optimal to apply filter on dimension column - if (CarbonUtil.usePreviousFilterBitsetGroup(useBitsetPipeLine, prvBitSetGroup, pageNumber, - filterValues.length)) { + if (filterValues.length > 0 && CarbonUtil + .usePreviousFilterBitsetGroup(useBitsetPipeLine, prvBitSetGroup, pageNumber, + filterValues.length)) { return getFilteredIndexesUisngPrvBitset(dimensionColumnPage, prvBitSetGroup, pageNumber, numberOfRows); } else { @@ -379,6 +380,9 @@ public class IncludeFilterExecuterImpl implements FilterExecuter { private BitSet setFilterdIndexToBitSetWithColumnIndex( DimensionColumnPage dimensionColumnPage, int numerOfRows) { BitSet bitSet = new BitSet(numerOfRows); + if (filterValues.length == 0) { + return bitSet; + } int startIndex = 0; for (int i = 0; i < filterValues.length; i++) { if (startIndex >= numerOfRows) { @@ -400,6 +404,9 @@ public class IncludeFilterExecuterImpl implements FilterExecuter { private BitSet setFilterdIndexToBitSet(DimensionColumnPage dimensionColumnPage, int numerOfRows) { BitSet bitSet = new BitSet(numerOfRows); + if (filterValues.length == 0) { + return bitSet; + } // binary search can only be applied if column is sorted and // inverted index exists for that column if (isNaturalSorted) { http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RangeValueFilterExecuterImpl.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RangeValueFilterExecuterImpl.java b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RangeValueFilterExecuterImpl.java index a4bbdd0..7919e35 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RangeValueFilterExecuterImpl.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RangeValueFilterExecuterImpl.java @@ -90,7 +90,7 @@ public class RangeValueFilterExecuterImpl extends ValueBasedFilterExecuterImpl { isRangeFullyCoverBlock = false; initDimensionChunkIndexes(); ifDefaultValueMatchesFilter(); - if (isDimensionPresentInCurrentBlock == true) { + if (isDimensionPresentInCurrentBlock) { isNaturalSorted = dimColEvaluatorInfo.getDimension().isUseInvertedIndex() && dimColEvaluatorInfo.getDimension().isSortColumn(); } @@ -373,7 +373,7 @@ public class RangeValueFilterExecuterImpl extends ValueBasedFilterExecuterImpl { if (null != rawColumnChunk.getLocalDictionary()) { if (null == filterExecuter) { filterExecuter = FilterUtil - .getFilterExecutorForLocalDictionary(rawColumnChunk, exp, isNaturalSorted); + .getFilterExecutorForRangeFilters(rawColumnChunk, exp, isNaturalSorted); if (filterExecuter instanceof ExcludeFilterExecuterImpl) { isExclude = true; } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeGrtThanFiterExecuterImpl.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeGrtThanFiterExecuterImpl.java b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeGrtThanFiterExecuterImpl.java index f2cb3dd..03e41eb 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeGrtThanFiterExecuterImpl.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeGrtThanFiterExecuterImpl.java @@ -211,7 +211,7 @@ public class RowLevelRangeGrtThanFiterExecuterImpl extends RowLevelFilterExecute if (null != rawColumnChunk.getLocalDictionary()) { if (null == filterExecuter) { filterExecuter = FilterUtil - .getFilterExecutorForLocalDictionary(rawColumnChunk, exp, isNaturalSorted); + .getFilterExecutorForRangeFilters(rawColumnChunk, exp, isNaturalSorted); if (filterExecuter instanceof ExcludeFilterExecuterImpl) { isExclude = true; } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeGrtrThanEquaToFilterExecuterImpl.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeGrtrThanEquaToFilterExecuterImpl.java b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeGrtrThanEquaToFilterExecuterImpl.java index b893f99..95cd356 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeGrtrThanEquaToFilterExecuterImpl.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeGrtrThanEquaToFilterExecuterImpl.java @@ -208,7 +208,7 @@ public class RowLevelRangeGrtrThanEquaToFilterExecuterImpl extends RowLevelFilte if (null != rawColumnChunk.getLocalDictionary()) { if (null == filterExecuter) { filterExecuter = FilterUtil - .getFilterExecutorForLocalDictionary(rawColumnChunk, exp, isNaturalSorted); + .getFilterExecutorForRangeFilters(rawColumnChunk, exp, isNaturalSorted); if (filterExecuter instanceof ExcludeFilterExecuterImpl) { isExclude = true; } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeLessThanEqualFilterExecuterImpl.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeLessThanEqualFilterExecuterImpl.java b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeLessThanEqualFilterExecuterImpl.java index 1b0a64a..dc5ed37 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeLessThanEqualFilterExecuterImpl.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeLessThanEqualFilterExecuterImpl.java @@ -203,7 +203,7 @@ public class RowLevelRangeLessThanEqualFilterExecuterImpl extends RowLevelFilter if (null != rawColumnChunk.getLocalDictionary()) { if (null == filterExecuter) { filterExecuter = FilterUtil - .getFilterExecutorForLocalDictionary(rawColumnChunk, exp, isNaturalSorted); + .getFilterExecutorForRangeFilters(rawColumnChunk, exp, isNaturalSorted); if (filterExecuter instanceof ExcludeFilterExecuterImpl) { isExclude = true; } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeLessThanFilterExecuterImpl.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeLessThanFilterExecuterImpl.java b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeLessThanFilterExecuterImpl.java index 5f68e34..c41a08d 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeLessThanFilterExecuterImpl.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/filter/executer/RowLevelRangeLessThanFilterExecuterImpl.java @@ -203,7 +203,7 @@ public class RowLevelRangeLessThanFilterExecuterImpl extends RowLevelFilterExecu if (null != rawColumnChunk.getLocalDictionary()) { if (null == filterExecuter) { filterExecuter = FilterUtil - .getFilterExecutorForLocalDictionary(rawColumnChunk, exp, isNaturalSorted); + .getFilterExecutorForRangeFilters(rawColumnChunk, exp, isNaturalSorted); if (filterExecuter instanceof ExcludeFilterExecuterImpl) { isExclude = true; } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonColumnVector.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonColumnVector.java b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonColumnVector.java index b606a50..3bee136 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonColumnVector.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonColumnVector.java @@ -57,6 +57,10 @@ public interface CarbonColumnVector { void putNulls(int rowId, int count); + void putNotNull(int rowId); + + void putNotNull(int rowId, int count); + boolean isNull(int rowId); void putObject(int rowId, Object obj); @@ -83,4 +87,11 @@ public interface CarbonColumnVector { void setFilteredRowsExist(boolean filteredRowsExist); + + void setDictionary(CarbonDictionary dictionary); + + boolean hasDictionary(); + + CarbonColumnVector getDictionaryVector(); + } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonColumnarBatch.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonColumnarBatch.java b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonColumnarBatch.java index 973ce0f..803715c 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonColumnarBatch.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonColumnarBatch.java @@ -86,5 +86,4 @@ public class CarbonColumnarBatch { } } } - } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonDictionary.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonDictionary.java b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonDictionary.java new file mode 100644 index 0000000..50d2ac5 --- /dev/null +++ b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/CarbonDictionary.java @@ -0,0 +1,30 @@ +/* + * 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.scan.result.vector; + +public interface CarbonDictionary { + + int getDictionaryActualSize(); + + int getDictionarySize(); + + boolean isDictionaryUsed(); + + void setDictionaryUsed(); + + byte[] getDictionaryValue(int index); +} http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/result/vector/impl/CarbonColumnVectorImpl.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/result/vector/impl/CarbonColumnVectorImpl.java b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/impl/CarbonColumnVectorImpl.java index e431aaf..cc6ddfc 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/result/vector/impl/CarbonColumnVectorImpl.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/impl/CarbonColumnVectorImpl.java @@ -24,8 +24,7 @@ import org.apache.carbondata.core.metadata.datatype.DataType; import org.apache.carbondata.core.metadata.datatype.DataTypes; import org.apache.carbondata.core.metadata.datatype.DecimalType; import org.apache.carbondata.core.scan.result.vector.CarbonColumnVector; - - +import org.apache.carbondata.core.scan.result.vector.CarbonDictionary; public class CarbonColumnVectorImpl implements CarbonColumnVector { @@ -59,6 +58,9 @@ public class CarbonColumnVectorImpl implements CarbonColumnVector { */ protected boolean anyNullsSet; + private CarbonDictionary carbonDictionary; + + private CarbonColumnVector dictionaryVector; public CarbonColumnVectorImpl(int batchSize, DataType dataType) { nullBytes = new BitSet(batchSize); @@ -78,6 +80,7 @@ public class CarbonColumnVectorImpl implements CarbonColumnVector { } else if (dataType instanceof DecimalType) { decimals = new BigDecimal[batchSize]; } else if (dataType == DataTypes.STRING || dataType == DataTypes.BYTE_ARRAY) { + dictionaryVector = new CarbonColumnVectorImpl(batchSize, DataTypes.INT); bytes = new byte[batchSize][]; } else { data = new Object[batchSize]; @@ -170,6 +173,13 @@ public class CarbonColumnVectorImpl implements CarbonColumnVector { anyNullsSet = true; } + @Override public void putNotNull(int rowId) { + + } + + @Override public void putNotNull(int rowId, int count) { + + } public boolean isNullAt(int rowId) { return nullBytes.get(rowId); @@ -203,7 +213,11 @@ public class CarbonColumnVectorImpl implements CarbonColumnVector { } else if (dataType instanceof DecimalType) { return decimals[rowId]; } else if (dataType == DataTypes.STRING || dataType == DataTypes.BYTE_ARRAY) { - return bytes[rowId]; + if (null != carbonDictionary) { + int dictKey = (Integer) dictionaryVector.getData(rowId); + return carbonDictionary.getDictionaryValue(dictKey); + } + return bytes[rowId]; } else { return data[rowId]; } @@ -227,6 +241,7 @@ public class CarbonColumnVectorImpl implements CarbonColumnVector { Arrays.fill(decimals, null); } else if (dataType == DataTypes.STRING || dataType == DataTypes.BYTE_ARRAY) { Arrays.fill(bytes, null); + this.dictionaryVector.reset(); } else { Arrays.fill(data, null); } @@ -251,6 +266,18 @@ public class CarbonColumnVectorImpl implements CarbonColumnVector { } + @Override public void setDictionary(CarbonDictionary dictionary) { + this.carbonDictionary = dictionary; + } + + @Override public boolean hasDictionary() { + return null != this.carbonDictionary; + } + + @Override public CarbonColumnVector getDictionaryVector() { + return dictionaryVector; + } + /** * Returns true if any of the nulls indicator are set for this column. This can be used * as an optimization to prevent setting nulls. http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/core/src/main/java/org/apache/carbondata/core/scan/result/vector/impl/CarbonDictionaryImpl.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/result/vector/impl/CarbonDictionaryImpl.java b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/impl/CarbonDictionaryImpl.java new file mode 100644 index 0000000..cc3a03c --- /dev/null +++ b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/impl/CarbonDictionaryImpl.java @@ -0,0 +1,54 @@ +/* + * 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.scan.result.vector.impl; + +import org.apache.carbondata.core.scan.result.vector.CarbonDictionary; + +public class CarbonDictionaryImpl implements CarbonDictionary { + + private byte[][] dictionary; + + private int actualSize; + + private boolean isDictUsed; + + public CarbonDictionaryImpl(byte[][] dictionary, int actualSize) { + this.dictionary = dictionary; + this.actualSize = actualSize; + } + + @Override public int getDictionaryActualSize() { + return actualSize; + } + + @Override public int getDictionarySize() { + return this.dictionary.length; + } + + @Override public boolean isDictionaryUsed() { + return this.isDictUsed; + } + + @Override public void setDictionaryUsed() { + this.isDictUsed = true; + } + + @Override public byte[] getDictionaryValue(int index) { + return dictionary[index]; + } + +} http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/examples/spark2/src/main/scala/org/apache/carbondata/examples/CarbonSessionExample.scala ---------------------------------------------------------------------- diff --git a/examples/spark2/src/main/scala/org/apache/carbondata/examples/CarbonSessionExample.scala b/examples/spark2/src/main/scala/org/apache/carbondata/examples/CarbonSessionExample.scala index 2883ded..e63e852 100644 --- a/examples/spark2/src/main/scala/org/apache/carbondata/examples/CarbonSessionExample.scala +++ b/examples/spark2/src/main/scala/org/apache/carbondata/examples/CarbonSessionExample.scala @@ -151,4 +151,4 @@ object CarbonSessionExample { spark.sql("DROP TABLE IF EXISTS carbonsession_table") spark.sql("DROP TABLE IF EXISTS stored_as_carbondata_table") } -} \ No newline at end of file +} http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/integration/presto/src/main/java/org/apache/carbondata/presto/CarbonColumnVectorWrapper.java ---------------------------------------------------------------------- diff --git a/integration/presto/src/main/java/org/apache/carbondata/presto/CarbonColumnVectorWrapper.java b/integration/presto/src/main/java/org/apache/carbondata/presto/CarbonColumnVectorWrapper.java index 4560241..b2c4c68 100644 --- a/integration/presto/src/main/java/org/apache/carbondata/presto/CarbonColumnVectorWrapper.java +++ b/integration/presto/src/main/java/org/apache/carbondata/presto/CarbonColumnVectorWrapper.java @@ -25,12 +25,12 @@ import org.apache.carbondata.core.metadata.datatype.DataType; import org.apache.carbondata.core.metadata.datatype.DataTypes; import org.apache.carbondata.core.metadata.datatype.StructField; import org.apache.carbondata.core.scan.result.vector.CarbonColumnVector; +import org.apache.carbondata.core.scan.result.vector.CarbonDictionary; import org.apache.carbondata.core.scan.result.vector.impl.CarbonColumnVectorImpl; import org.apache.spark.sql.types.ArrayType; import org.apache.spark.sql.types.BooleanType; import org.apache.spark.sql.types.DateType; -import org.apache.spark.sql.types.Decimal; import org.apache.spark.sql.types.DecimalType; import org.apache.spark.sql.types.DoubleType; import org.apache.spark.sql.types.FloatType; @@ -202,6 +202,14 @@ public class CarbonColumnVectorWrapper implements CarbonColumnVector { } } + @Override public void putNotNull(int rowId) { + + } + + @Override public void putNotNull(int rowId, int count) { + + } + @Override public boolean isNull(int rowId) { return columnVector.isNullAt(rowId); } @@ -238,6 +246,18 @@ public class CarbonColumnVectorWrapper implements CarbonColumnVector { this.filteredRowsExist = filteredRowsExist; } + @Override public void setDictionary(CarbonDictionary dictionary) { + this.columnVector.setDictionary(dictionary); + } + + @Override public boolean hasDictionary() { + return this.columnVector.hasDictionary(); + } + + @Override public CarbonColumnVector getDictionaryVector() { + return this.columnVector; + } + // TODO: this is copied from carbondata-spark-common module, use presto type instead of this private org.apache.carbondata.core.metadata.datatype.DataType convertSparkToCarbonDataType(org.apache.spark.sql.types.DataType dataType) { http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/integration/spark-common/src/main/scala/org/apache/carbondata/spark/util/CarbonScalaUtil.scala ---------------------------------------------------------------------- diff --git a/integration/spark-common/src/main/scala/org/apache/carbondata/spark/util/CarbonScalaUtil.scala b/integration/spark-common/src/main/scala/org/apache/carbondata/spark/util/CarbonScalaUtil.scala index b3f56a2..30c1874 100644 --- a/integration/spark-common/src/main/scala/org/apache/carbondata/spark/util/CarbonScalaUtil.scala +++ b/integration/spark-common/src/main/scala/org/apache/carbondata/spark/util/CarbonScalaUtil.scala @@ -686,7 +686,8 @@ object CarbonScalaUtil { // considered which is 1000 Try(localDictionaryThreshold.toInt) match { case scala.util.Success(value) => - if (value < 1000 || value > 100000) { + if (value < CarbonCommonConstants.LOCAL_DICTIONARY_MIN || + value > CarbonCommonConstants.LOCAL_DICTIONARY_MAX) { false } else { true @@ -727,4 +728,9 @@ object CarbonScalaUtil { } } } + + def isStringDataType(dataType: DataType): Boolean = { + dataType == StringType + } + } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/integration/spark2/src/main/java/org/apache/carbondata/spark/vectorreader/CarbonDictionaryWrapper.java ---------------------------------------------------------------------- diff --git a/integration/spark2/src/main/java/org/apache/carbondata/spark/vectorreader/CarbonDictionaryWrapper.java b/integration/spark2/src/main/java/org/apache/carbondata/spark/vectorreader/CarbonDictionaryWrapper.java new file mode 100644 index 0000000..044ed58 --- /dev/null +++ b/integration/spark2/src/main/java/org/apache/carbondata/spark/vectorreader/CarbonDictionaryWrapper.java @@ -0,0 +1,44 @@ +/* + * 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.spark.vectorreader; + +import org.apache.carbondata.core.scan.result.vector.CarbonDictionary; + +import org.apache.parquet.column.Dictionary; +import org.apache.parquet.column.Encoding; +import org.apache.parquet.io.api.Binary; + +public class CarbonDictionaryWrapper extends Dictionary { + + private Binary[] binaries; + + public CarbonDictionaryWrapper(Encoding encoding, CarbonDictionary dictionary) { + super(encoding); + binaries = new Binary[dictionary.getDictionarySize()]; + for (int i = 0; i < binaries.length; i++) { + binaries[i] = Binary.fromReusedByteArray(dictionary.getDictionaryValue(i)); + } + } + + @Override public int getMaxId() { + return binaries.length - 1; + } + + @Override public Binary decodeToBinary(int id) { + return binaries[id]; + } +} http://git-wip-us.apache.org/repos/asf/carbondata/blob/3a4b8813/integration/spark2/src/main/java/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapper.java ---------------------------------------------------------------------- diff --git a/integration/spark2/src/main/java/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapper.java b/integration/spark2/src/main/java/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapper.java index 9e0c102..4bba658 100644 --- a/integration/spark2/src/main/java/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapper.java +++ b/integration/spark2/src/main/java/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapper.java @@ -21,8 +21,10 @@ import java.math.BigDecimal; import org.apache.carbondata.core.metadata.datatype.DataType; import org.apache.carbondata.core.scan.result.vector.CarbonColumnVector; +import org.apache.carbondata.core.scan.result.vector.CarbonDictionary; import org.apache.carbondata.spark.util.CarbonScalaUtil; +import org.apache.parquet.column.Encoding; import org.apache.spark.sql.execution.vectorized.ColumnVector; import org.apache.spark.sql.types.Decimal; @@ -38,9 +40,15 @@ class ColumnarVectorWrapper implements CarbonColumnVector { private DataType blockDataType; + private CarbonColumnVector dictionaryVector; + ColumnarVectorWrapper(ColumnVector columnVector, boolean[] filteredRows) { this.columnVector = columnVector; this.filteredRows = filteredRows; + if (columnVector.getDictionaryIds() != null) { + this.dictionaryVector = + new ColumnarVectorWrapper(columnVector.getDictionaryIds(), filteredRows); + } } @Override public void putBoolean(int rowId, boolean value) { @@ -188,6 +196,25 @@ class ColumnarVectorWrapper implements CarbonColumnVector { } } + @Override public void putNotNull(int rowId) { + if (!filteredRows[rowId]) { + columnVector.putNotNull(counter++); + } + } + + @Override public void putNotNull(int rowId, int count) { + if (filteredRowsExist) { + for (int i = 0; i < count; i++) { + if (!filteredRows[rowId]) { + columnVector.putNotNull(counter++); + } + rowId++; + } + } else { + columnVector.putNotNulls(rowId, count); + } + } + @Override public boolean isNull(int rowId) { return columnVector.isNullAt(rowId); } @@ -204,6 +231,9 @@ class ColumnarVectorWrapper implements CarbonColumnVector { @Override public void reset() { counter = 0; filteredRowsExist = false; + if (null != dictionaryVector) { + dictionaryVector.reset(); + } } @Override public DataType getType() { @@ -223,4 +253,20 @@ class ColumnarVectorWrapper implements CarbonColumnVector { @Override public void setFilteredRowsExist(boolean filteredRowsExist) { this.filteredRowsExist = filteredRowsExist; } + + @Override public void setDictionary(CarbonDictionary dictionary) { + if (dictionary == null) { + columnVector.setDictionary(null); + } else { + columnVector.setDictionary(new CarbonDictionaryWrapper(Encoding.PLAIN, dictionary)); + } + } + + @Override public boolean hasDictionary() { + return columnVector.hasDictionary(); + } + + @Override public CarbonColumnVector getDictionaryVector() { + return dictionaryVector; + } }
