Repository: carbondata Updated Branches: refs/heads/master 6e9ba6d19 -> 04084c73f
http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/core/src/main/java/org/apache/carbondata/core/stream/StreamPruner.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/stream/StreamPruner.java b/core/src/main/java/org/apache/carbondata/core/stream/StreamPruner.java index b40355b..178b9f1 100644 --- a/core/src/main/java/org/apache/carbondata/core/stream/StreamPruner.java +++ b/core/src/main/java/org/apache/carbondata/core/stream/StreamPruner.java @@ -100,7 +100,8 @@ public class StreamPruner { } byte[][] maxValue = streamFile.getMinMaxIndex().getMaxValues(); byte[][] minValue = streamFile.getMinMaxIndex().getMinValues(); - BitSet bitSet = filterExecuter.isScanRequired(maxValue, minValue); + BitSet bitSet = filterExecuter + .isScanRequired(maxValue, minValue, streamFile.getMinMaxIndex().getIsMinMaxSet()); if (!bitSet.isEmpty()) { return true; } else { http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/core/src/main/java/org/apache/carbondata/core/util/AbstractDataFileFooterConverter.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/util/AbstractDataFileFooterConverter.java b/core/src/main/java/org/apache/carbondata/core/util/AbstractDataFileFooterConverter.java index 1143ed5..b1dd580 100644 --- a/core/src/main/java/org/apache/carbondata/core/util/AbstractDataFileFooterConverter.java +++ b/core/src/main/java/org/apache/carbondata/core/util/AbstractDataFileFooterConverter.java @@ -19,6 +19,7 @@ package org.apache.carbondata.core.util; import java.io.IOException; import java.nio.ByteBuffer; import java.util.ArrayList; +import java.util.Arrays; import java.util.BitSet; import java.util.List; import java.util.Map; @@ -285,10 +286,21 @@ public abstract class AbstractDataFileFooterConverter { byte[][] currentMaxValue = blockletIndexList.get(0).getMinMaxIndex().getMaxValues().clone(); byte[][] minValue = null; byte[][] maxValue = null; + boolean[] blockletMinMaxFlag = null; + // flag at block level + boolean[] blockMinMaxFlag = blockletIndexList.get(0).getMinMaxIndex().getIsMinMaxSet(); for (int i = 1; i < blockletIndexList.size(); i++) { minValue = blockletIndexList.get(i).getMinMaxIndex().getMinValues(); maxValue = blockletIndexList.get(i).getMinMaxIndex().getMaxValues(); + blockletMinMaxFlag = blockletIndexList.get(i).getMinMaxIndex().getIsMinMaxSet(); for (int j = 0; j < maxValue.length; j++) { + // can be null for stores < 1.5.0 version + if (null != blockletMinMaxFlag && !blockletMinMaxFlag[i]) { + blockMinMaxFlag[i] = blockletMinMaxFlag[i]; + currentMaxValue[j] = new byte[0]; + currentMinValue[j] = new byte[0]; + continue; + } if (ByteUtil.UnsafeComparer.INSTANCE.compareTo(currentMinValue[j], minValue[j]) > 0) { currentMinValue[j] = minValue[j].clone(); } @@ -297,10 +309,14 @@ public abstract class AbstractDataFileFooterConverter { } } } - + if (null == blockMinMaxFlag) { + blockMinMaxFlag = new boolean[currentMaxValue.length]; + Arrays.fill(blockMinMaxFlag, true); + } BlockletMinMaxIndex minMax = new BlockletMinMaxIndex(); minMax.setMaxValues(currentMaxValue); minMax.setMinValues(currentMinValue); + minMax.setIsMinMaxSet(blockMinMaxFlag); blockletIndex.setMinMaxIndex(minMax); return blockletIndex; } @@ -418,9 +434,19 @@ public abstract class AbstractDataFileFooterConverter { blockletIndexThrift.getB_tree_index(); org.apache.carbondata.format.BlockletMinMaxIndex minMaxIndex = blockletIndexThrift.getMin_max_index(); + List<Boolean> isMinMaxSet = null; + // Below logic is added to handle backward compatibility + if (minMaxIndex.isSetMin_max_presence()) { + isMinMaxSet = minMaxIndex.getMin_max_presence(); + } else { + Boolean[] minMaxFlag = new Boolean[minMaxIndex.getMax_values().size()]; + Arrays.fill(minMaxFlag, true); + isMinMaxSet = Arrays.asList(minMaxFlag); + } return new BlockletIndex( new BlockletBTreeIndex(btreeIndex.getStart_key(), btreeIndex.getEnd_key()), - new BlockletMinMaxIndex(minMaxIndex.getMin_values(), minMaxIndex.getMax_values())); + new BlockletMinMaxIndex(minMaxIndex.getMin_values(), minMaxIndex.getMax_values(), + isMinMaxSet)); } /** http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/core/src/main/java/org/apache/carbondata/core/util/BlockletDataMapUtil.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/util/BlockletDataMapUtil.java b/core/src/main/java/org/apache/carbondata/core/util/BlockletDataMapUtil.java index 4a87e91..ac53b56 100644 --- a/core/src/main/java/org/apache/carbondata/core/util/BlockletDataMapUtil.java +++ b/core/src/main/java/org/apache/carbondata/core/util/BlockletDataMapUtil.java @@ -48,6 +48,7 @@ import org.apache.carbondata.core.indexstore.blockletindex.BlockletDataMapDistri import org.apache.carbondata.core.indexstore.blockletindex.BlockletDataMapFactory; import org.apache.carbondata.core.indexstore.blockletindex.SegmentIndexFileStore; import org.apache.carbondata.core.metadata.blocklet.DataFileFooter; +import org.apache.carbondata.core.metadata.blocklet.index.BlockletMinMaxIndex; import org.apache.carbondata.core.metadata.datatype.DataType; import org.apache.carbondata.core.metadata.datatype.DataTypes; import org.apache.carbondata.core.metadata.schema.table.CarbonTable; @@ -485,4 +486,23 @@ public class BlockletDataMapUtil { } return false; } + + /** + * Method to update the min max flag. For CACHE_LEVEL=BLOCK, for any column if min max is not + * written in any of the blocklet then for that column the flag will be false for the + * complete block + * + * @param minMaxIndex + * @param minMaxFlag + */ + public static void updateMinMaxFlag(BlockletMinMaxIndex minMaxIndex, boolean[] minMaxFlag) { + boolean[] isMinMaxSet = minMaxIndex.getIsMinMaxSet(); + if (null != isMinMaxSet) { + for (int i = 0; i < minMaxFlag.length; i++) { + if (!isMinMaxSet[i]) { + minMaxFlag[i] = isMinMaxSet[i]; + } + } + } + } } \ No newline at end of file http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/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 4be4f78..0167c9a 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 @@ -19,12 +19,16 @@ package org.apache.carbondata.core.util; import java.io.IOException; import java.nio.ByteBuffer; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; +import org.apache.carbondata.common.logging.LogService; +import org.apache.carbondata.common.logging.LogServiceFactory; import org.apache.carbondata.core.datastore.blocklet.BlockletEncodedColumnPage; import org.apache.carbondata.core.datastore.blocklet.EncodedBlocklet; import org.apache.carbondata.core.datastore.compression.CompressorFactory; import org.apache.carbondata.core.datastore.page.encoding.EncodedColumnPage; +import org.apache.carbondata.core.datastore.page.statistics.SimpleStatsResult; import org.apache.carbondata.core.datastore.page.statistics.TablePageStatistics; import org.apache.carbondata.core.metadata.ColumnarFormatVersion; import org.apache.carbondata.core.metadata.datatype.DataType; @@ -52,6 +56,9 @@ import org.apache.carbondata.format.SegmentInfo; */ public class CarbonMetadataUtil { + private static final LogService LOGGER = + LogServiceFactory.getLogService(CarbonMetadataUtil.class.getName()); + private CarbonMetadataUtil() { } @@ -105,9 +112,16 @@ public class CarbonMetadataUtil { if (minMaxIndex == null) { return null; } - + List<Boolean> isMinMaxSet = null; + if (minMaxIndex.isSetMin_max_presence()) { + isMinMaxSet = minMaxIndex.getMin_max_presence(); + } else { + Boolean[] minMaxFlag = new Boolean[minMaxIndex.getMax_values().size()]; + Arrays.fill(minMaxFlag, true); + isMinMaxSet = Arrays.asList(minMaxFlag); + } return new org.apache.carbondata.core.metadata.blocklet.index.BlockletMinMaxIndex( - minMaxIndex.getMin_values(), minMaxIndex.getMax_values()); + minMaxIndex.getMin_values(), minMaxIndex.getMax_values(), isMinMaxSet); } /** @@ -124,6 +138,7 @@ public class CarbonMetadataUtil { for (int i = 0; i < minMaxIndex.getMaxValues().length; i++) { blockletMinMaxIndex.addToMax_values(ByteBuffer.wrap(minMaxIndex.getMaxValues()[i])); blockletMinMaxIndex.addToMin_values(ByteBuffer.wrap(minMaxIndex.getMinValues()[i])); + blockletMinMaxIndex.addToMin_max_presence(minMaxIndex.getIsMinMaxSet()[i]); } return blockletMinMaxIndex; @@ -177,20 +192,28 @@ public class CarbonMetadataUtil { public static BlockletIndex getBlockletIndex(EncodedBlocklet encodedBlocklet, List<CarbonMeasure> carbonMeasureList) { BlockletMinMaxIndex blockletMinMaxIndex = new BlockletMinMaxIndex(); - + // merge writeMinMax flag for all the dimensions + List<Boolean> writeMinMaxFlag = + mergeWriteMinMaxFlagForAllPages(blockletMinMaxIndex, encodedBlocklet); // Calculating min/max for every each column. TablePageStatistics stats = new TablePageStatistics(getEncodedColumnPages(encodedBlocklet, true, 0), getEncodedColumnPages(encodedBlocklet, false, 0)); byte[][] minCol = stats.getDimensionMinValue().clone(); byte[][] maxCol = stats.getDimensionMaxValue().clone(); - 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++) { + // if writeMonMaxFlag is set to false for the dimension at index i, then update the page + // and blocklet min/max with empty byte array + if (!writeMinMaxFlag.get(i)) { + maxCol[i] = new byte[0]; + minCol[i] = new byte[0]; + continue; + } if (ByteUtil.UnsafeComparer.INSTANCE.compareTo(columnMaxData[i], maxCol[i]) > 0) { maxCol[i] = columnMaxData[i]; } @@ -250,6 +273,41 @@ public class CarbonMetadataUtil { } /** + * This method will combine the writeMinMax flag from all the pages. If any page for a given + * dimension has writeMinMax flag set to false then min max for that dimension will nto be + * written in any of the page and metadata + * + * @param blockletMinMaxIndex + * @param encodedBlocklet + */ + private static List<Boolean> mergeWriteMinMaxFlagForAllPages( + BlockletMinMaxIndex blockletMinMaxIndex, EncodedBlocklet encodedBlocklet) { + Boolean[] mergedWriteMinMaxFlag = + new Boolean[encodedBlocklet.getNumberOfDimension() + encodedBlocklet.getNumberOfMeasure()]; + // set writeMinMax flag to true for all the columns by default and then update if stats object + // has the this flag set to false + Arrays.fill(mergedWriteMinMaxFlag, true); + for (int i = 0; i < encodedBlocklet.getNumberOfDimension(); i++) { + for (int pageIndex = 0; pageIndex < encodedBlocklet.getNumberOfPages(); pageIndex++) { + EncodedColumnPage encodedColumnPage = + encodedBlocklet.getEncodedDimensionColumnPages().get(i).getEncodedColumnPageList() + .get(pageIndex); + SimpleStatsResult stats = encodedColumnPage.getStats(); + if (!stats.writeMinMax()) { + mergedWriteMinMaxFlag[i] = stats.writeMinMax(); + String columnName = encodedColumnPage.getActualPage().getColumnSpec().getFieldName(); + LOGGER.info("Min Max writing of blocklet ignored for column with name " + columnName); + break; + } + } + } + List<Boolean> min_max_presence = Arrays.asList(mergedWriteMinMaxFlag); + blockletMinMaxIndex.setMin_max_presence(min_max_presence); + return min_max_presence; + } + + /** + * Right now it is set to default values. We may use this in future * set the compressor. * before 1.5.0, we set a enum 'compression_codec'; * after 1.5.0, we use string 'compressor_name' instead http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/core/src/main/java/org/apache/carbondata/core/util/CarbonProperties.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/util/CarbonProperties.java b/core/src/main/java/org/apache/carbondata/core/util/CarbonProperties.java index 559320a..3438c4e 100644 --- a/core/src/main/java/org/apache/carbondata/core/util/CarbonProperties.java +++ b/core/src/main/java/org/apache/carbondata/core/util/CarbonProperties.java @@ -41,6 +41,7 @@ import static org.apache.carbondata.core.constants.CarbonCommonConstants.CARBON_ import static org.apache.carbondata.core.constants.CarbonCommonConstants.CARBON_DATA_FILE_VERSION; import static org.apache.carbondata.core.constants.CarbonCommonConstants.CARBON_DATE_FORMAT; import static org.apache.carbondata.core.constants.CarbonCommonConstants.CARBON_DYNAMIC_ALLOCATION_SCHEDULER_TIMEOUT; +import static org.apache.carbondata.core.constants.CarbonCommonConstants.CARBON_MINMAX_ALLOWED_BYTE_COUNT; import static org.apache.carbondata.core.constants.CarbonCommonConstants.CARBON_PREFETCH_BUFFERSIZE; import static org.apache.carbondata.core.constants.CarbonCommonConstants.CARBON_SCHEDULER_MIN_REGISTERED_RESOURCES_RATIO; import static org.apache.carbondata.core.constants.CarbonCommonConstants.CARBON_SCHEDULER_MIN_REGISTERED_RESOURCES_RATIO_DEFAULT; @@ -194,6 +195,9 @@ public final class CarbonProperties { case CARBON_LOAD_SORT_MEMORY_SPILL_PERCENTAGE: validateSortMemorySpillPercentage(); break; + case CARBON_MINMAX_ALLOWED_BYTE_COUNT: + validateStringCharacterLimit(); + break; // TODO : Validation for carbon.lock.type should be handled for addProperty flow default: // none @@ -258,6 +262,7 @@ public final class CarbonProperties { validateSortStorageMemory(); validateEnableQueryStatistics(); validateSortMemorySpillPercentage(); + validateStringCharacterLimit(); } /** @@ -1551,4 +1556,36 @@ public final class CarbonProperties { CarbonLoadOptionConstants.CARBON_LOAD_SORT_MEMORY_SPILL_PERCENTAGE_DEFAULT); } } + + /** + * This method validates the allowed character limit for storing min/max for string type + */ + private void validateStringCharacterLimit() { + int allowedCharactersLimit = 0; + try { + allowedCharactersLimit = Integer.parseInt(carbonProperties + .getProperty(CARBON_MINMAX_ALLOWED_BYTE_COUNT, + CarbonCommonConstants.CARBON_MINMAX_ALLOWED_BYTE_COUNT_DEFAULT)); + if (allowedCharactersLimit < CarbonCommonConstants.CARBON_MINMAX_ALLOWED_BYTE_COUNT_MIN + || allowedCharactersLimit + > CarbonCommonConstants.CARBON_MINMAX_ALLOWED_BYTE_COUNT_MAX) { + LOGGER.info("The min max byte limit for string type value \"" + allowedCharactersLimit + + "\" is invalid. Using the default value \"" + + CarbonCommonConstants.CARBON_MINMAX_ALLOWED_BYTE_COUNT_DEFAULT); + carbonProperties.setProperty(CARBON_MINMAX_ALLOWED_BYTE_COUNT, + CarbonCommonConstants.CARBON_MINMAX_ALLOWED_BYTE_COUNT_DEFAULT); + } else { + LOGGER.info( + "Considered value for min max byte limit for string is: " + allowedCharactersLimit); + carbonProperties + .setProperty(CARBON_MINMAX_ALLOWED_BYTE_COUNT, allowedCharactersLimit + ""); + } + } catch (NumberFormatException e) { + LOGGER.info("The min max byte limit for string type value \"" + allowedCharactersLimit + + "\" is invalid. Using the default value \"" + + CarbonCommonConstants.CARBON_MINMAX_ALLOWED_BYTE_COUNT_DEFAULT); + carbonProperties.setProperty(CARBON_MINMAX_ALLOWED_BYTE_COUNT, + CarbonCommonConstants.CARBON_MINMAX_ALLOWED_BYTE_COUNT_DEFAULT); + } + } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/core/src/test/java/org/apache/carbondata/core/indexstore/blockletindex/TestBlockletDataMap.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/carbondata/core/indexstore/blockletindex/TestBlockletDataMap.java b/core/src/test/java/org/apache/carbondata/core/indexstore/blockletindex/TestBlockletDataMap.java index 85de7c4..fee2e9d 100644 --- a/core/src/test/java/org/apache/carbondata/core/indexstore/blockletindex/TestBlockletDataMap.java +++ b/core/src/test/java/org/apache/carbondata/core/indexstore/blockletindex/TestBlockletDataMap.java @@ -35,7 +35,7 @@ public class TestBlockletDataMap extends AbstractDictionaryCacheTest { new MockUp<ImplicitIncludeFilterExecutorImpl>() { @Mock BitSet isFilterValuesPresentInBlockOrBlocklet(byte[][] maxValue, byte[][] minValue, - String uniqueBlockPath) { + String uniqueBlockPath, boolean[] isMinMaxSet) { BitSet bitSet = new BitSet(1); bitSet.set(8); return bitSet; @@ -45,13 +45,14 @@ public class TestBlockletDataMap extends AbstractDictionaryCacheTest { BlockDataMap blockletDataMap = new BlockletDataMap(); Method method = BlockDataMap.class .getDeclaredMethod("addBlockBasedOnMinMaxValue", FilterExecuter.class, byte[][].class, - byte[][].class, String.class, int.class); + byte[][].class, boolean[].class, String.class, int.class); method.setAccessible(true); byte[][] minValue = { ByteUtil.toBytes("sfds") }; byte[][] maxValue = { ByteUtil.toBytes("resa") }; + boolean[] minMaxFlag = new boolean[] {true}; Object result = method - .invoke(blockletDataMap, implicitIncludeFilterExecutor, minValue, maxValue, + .invoke(blockletDataMap, implicitIncludeFilterExecutor, minValue, maxValue, minMaxFlag, "/opt/store/default/carbon_table/Fact/Part0/Segment_0/part-0-0_batchno0-0-1514989110586.carbondata", 0); assert ((boolean) result); http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/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 14cd57a..7aa2236 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 @@ -172,10 +172,13 @@ public class CarbonMetadataUtilTest { List<ByteBuffer> maxList = new ArrayList<>(); maxList.add(ByteBuffer.wrap(byteArr1)); + List<Boolean> isMinMaxSet = new ArrayList<>(); + isMinMaxSet.add(true); + org.apache.carbondata.core.metadata.blocklet.index.BlockletMinMaxIndex blockletMinMaxIndex = new org.apache.carbondata.core.metadata.blocklet.index.BlockletMinMaxIndex(minList, - maxList); + maxList, isMinMaxSet); org.apache.carbondata.core.metadata.blocklet.index.BlockletBTreeIndex blockletBTreeIndex = new org.apache.carbondata.core.metadata.blocklet.index.BlockletBTreeIndex(startKey, http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/core/src/test/java/org/apache/carbondata/core/util/RangeFilterProcessorTest.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/carbondata/core/util/RangeFilterProcessorTest.java b/core/src/test/java/org/apache/carbondata/core/util/RangeFilterProcessorTest.java index 9b8be79..3fdce4e 100644 --- a/core/src/test/java/org/apache/carbondata/core/util/RangeFilterProcessorTest.java +++ b/core/src/test/java/org/apache/carbondata/core/util/RangeFilterProcessorTest.java @@ -320,7 +320,7 @@ public class RangeFilterProcessorTest { Deencapsulation.setField(range, "isDimensionPresentInCurrentBlock", true); Deencapsulation.setField(range, "lessThanExp", true); Deencapsulation.setField(range, "greaterThanExp", true); - result = range.isScanRequired(BlockMin, BlockMax, filterMinMax); + result = range.isScanRequired(BlockMin, BlockMax, filterMinMax, true); Assert.assertFalse(result); } @@ -336,7 +336,7 @@ public class RangeFilterProcessorTest { Deencapsulation.setField(range, "isDimensionPresentInCurrentBlock", true); Deencapsulation.setField(range, "lessThanExp", true); Deencapsulation.setField(range, "greaterThanExp", true); - result = range.isScanRequired(BlockMin, BlockMax, filterMinMax); + result = range.isScanRequired(BlockMin, BlockMax, filterMinMax, true); Assert.assertFalse(result); } @@ -352,7 +352,7 @@ public class RangeFilterProcessorTest { Deencapsulation.setField(range, "isDimensionPresentInCurrentBlock", true); Deencapsulation.setField(range, "lessThanExp", true); Deencapsulation.setField(range, "greaterThanExp", true); - result = range.isScanRequired(BlockMin, BlockMax, filterMinMax); + result = range.isScanRequired(BlockMin, BlockMax, filterMinMax, true); Assert.assertTrue(result); } @@ -370,7 +370,7 @@ public class RangeFilterProcessorTest { Deencapsulation.setField(range, "lessThanExp", true); Deencapsulation.setField(range, "greaterThanExp", true); - result = range.isScanRequired(BlockMin, BlockMax, filterMinMax); + result = range.isScanRequired(BlockMin, BlockMax, filterMinMax, true); rangeCovered = Deencapsulation.getField(range, "isRangeFullyCoverBlock"); Assert.assertTrue(result); Assert.assertTrue(rangeCovered); @@ -390,7 +390,7 @@ public class RangeFilterProcessorTest { Deencapsulation.setField(range, "lessThanExp", true); Deencapsulation.setField(range, "greaterThanExp", true); - result = range.isScanRequired(BlockMin, BlockMax, filterMinMax); + result = range.isScanRequired(BlockMin, BlockMax, filterMinMax, true); startBlockMinIsDefaultStart = Deencapsulation.getField(range, "startBlockMinIsDefaultStart"); Assert.assertTrue(result); Assert.assertTrue(startBlockMinIsDefaultStart); @@ -410,7 +410,7 @@ public class RangeFilterProcessorTest { Deencapsulation.setField(range, "lessThanExp", true); Deencapsulation.setField(range, "greaterThanExp", true); - result = range.isScanRequired(BlockMin, BlockMax, filterMinMax); + result = range.isScanRequired(BlockMin, BlockMax, filterMinMax, true); endBlockMaxisDefaultEnd = Deencapsulation.getField(range, "endBlockMaxisDefaultEnd"); Assert.assertTrue(result); Assert.assertTrue(endBlockMaxisDefaultEnd); http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/datamap/examples/src/minmaxdatamap/main/java/org/apache/carbondata/datamap/examples/MinMaxIndexDataMap.java ---------------------------------------------------------------------- diff --git a/datamap/examples/src/minmaxdatamap/main/java/org/apache/carbondata/datamap/examples/MinMaxIndexDataMap.java b/datamap/examples/src/minmaxdatamap/main/java/org/apache/carbondata/datamap/examples/MinMaxIndexDataMap.java index 5460247..40dc975 100644 --- a/datamap/examples/src/minmaxdatamap/main/java/org/apache/carbondata/datamap/examples/MinMaxIndexDataMap.java +++ b/datamap/examples/src/minmaxdatamap/main/java/org/apache/carbondata/datamap/examples/MinMaxIndexDataMap.java @@ -140,7 +140,7 @@ public class MinMaxIndexDataMap extends CoarseGrainDataMap { BitSet bitSet = filterExecuter.isScanRequired( readMinMaxDataMap[blkIdx][blkltIdx].getMaxValues(), - readMinMaxDataMap[blkIdx][blkltIdx].getMinValues()); + readMinMaxDataMap[blkIdx][blkltIdx].getMinValues(), null); if (!bitSet.isEmpty()) { String blockFileName = indexFilePath[blkIdx].substring( indexFilePath[blkIdx].lastIndexOf(File.separatorChar) + 1, http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/format/src/main/thrift/carbondata.thrift ---------------------------------------------------------------------- diff --git a/format/src/main/thrift/carbondata.thrift b/format/src/main/thrift/carbondata.thrift index 2423ffa..7130066 100644 --- a/format/src/main/thrift/carbondata.thrift +++ b/format/src/main/thrift/carbondata.thrift @@ -45,6 +45,7 @@ struct BlockletBTreeIndex{ struct BlockletMinMaxIndex{ 1: required list<binary> min_values; //Min value of all columns of one blocklet Bit-Packed 2: required list<binary> max_values; //Max value of all columns of one blocklet Bit-Packed + 3: optional list<bool> min_max_presence; // flag to specify whether min max is written for a column or not } /** http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/integration/spark-datasource/src/test/scala/org/apache/spark/sql/carbondata/datasource/TestCreateTableUsingSparkCarbonFileFormat.scala ---------------------------------------------------------------------- diff --git a/integration/spark-datasource/src/test/scala/org/apache/spark/sql/carbondata/datasource/TestCreateTableUsingSparkCarbonFileFormat.scala b/integration/spark-datasource/src/test/scala/org/apache/spark/sql/carbondata/datasource/TestCreateTableUsingSparkCarbonFileFormat.scala index 6a803fc..e6d4d48 100644 --- a/integration/spark-datasource/src/test/scala/org/apache/spark/sql/carbondata/datasource/TestCreateTableUsingSparkCarbonFileFormat.scala +++ b/integration/spark-datasource/src/test/scala/org/apache/spark/sql/carbondata/datasource/TestCreateTableUsingSparkCarbonFileFormat.scala @@ -19,24 +19,32 @@ package org.apache.spark.sql.carbondata.datasource import java.io.File import java.text.SimpleDateFormat +import java.util import java.util.{Date, Random} +import scala.collection.JavaConverters._ + import org.apache.commons.io.FileUtils import org.apache.commons.lang.RandomStringUtils import org.scalatest.{BeforeAndAfterAll, FunSuite} import org.apache.spark.util.SparkUtil import org.apache.spark.sql.carbondata.datasource.TestUtil.{spark, _} + import org.apache.carbondata.core.constants.CarbonCommonConstants import org.apache.carbondata.core.constants.{CarbonCommonConstants, CarbonV3DataFormatConstants} import org.apache.carbondata.core.datastore.filesystem.CarbonFile import org.apache.carbondata.core.datastore.impl.FileFactory import org.apache.carbondata.core.metadata.datatype.DataTypes -import org.apache.carbondata.core.util.{CarbonProperties, CarbonUtil} +import org.apache.carbondata.core.util.{CarbonProperties, CarbonUtil, DataFileFooterConverter} import org.apache.carbondata.sdk.file.{CarbonWriter, Field, Schema} import org.apache.hadoop.conf.Configuration import org.apache.spark.sql.Row import org.apache.spark.sql.carbondata.execution.datasources.CarbonFileIndexReplaceRule +import org.apache.carbondata.core.datamap.DataMapStoreManager +import org.apache.carbondata.core.metadata.AbsoluteTableIdentifier +import org.apache.carbondata.core.metadata.blocklet.DataFileFooter + class TestCreateTableUsingSparkCarbonFileFormat extends FunSuite with BeforeAndAfterAll { @@ -46,6 +54,9 @@ class TestCreateTableUsingSparkCarbonFileFormat extends FunSuite with BeforeAndA } override def afterAll(): Unit = { + CarbonProperties.getInstance() + .addProperty(CarbonCommonConstants.CARBON_MINMAX_ALLOWED_BYTE_COUNT, + CarbonCommonConstants.CARBON_MINMAX_ALLOWED_BYTE_COUNT_DEFAULT) spark.sql("DROP TABLE IF EXISTS sdkOutputTable") } @@ -328,8 +339,8 @@ class TestCreateTableUsingSparkCarbonFileFormat extends FunSuite with BeforeAndA assert(new File(filePath).exists()) cleanTestData() } - test("Read data having multi blocklet ") { - buildTestDataMuliBlockLet(700000) + test("Read data having multi blocklet and validate min max flag") { + buildTestDataMuliBlockLet(750000, 50000) assert(new File(writerPath).exists()) spark.sql("DROP TABLE IF EXISTS sdkOutputTable") @@ -342,41 +353,90 @@ class TestCreateTableUsingSparkCarbonFileFormat extends FunSuite with BeforeAndA s"""CREATE TABLE sdkOutputTable USING carbon LOCATION |'$writerPath' """.stripMargin) } - spark.sql("select count(*) from sdkOutputTable").show(false) - val result=checkAnswer(spark.sql("select count(*) from sdkOutputTable"),Seq(Row(700000))) + val result=checkAnswer(spark.sql("select count(*) from sdkOutputTable"),Seq(Row(800000))) if(result.isDefined){ assert(false,result.get) } + checkAnswer(spark + .sql( + "select count(*) from sdkOutputTable where from_email='Email for testing min max for " + + "allowed chars'"), + Seq(Row(50000))) + //expected answer for min max flag. FInally there should be 2 blocklets with one blocklet + // having min max flag as false for email column and other as true + val blocklet1MinMaxFlag = Array(true, true, true, true, true, false, true, true, true) + val blocklet2MinMaxFlag = Array(true, true, true, true, true, true, true, true, true) + val expectedMinMaxFlag = Array(blocklet1MinMaxFlag, blocklet2MinMaxFlag) + validateMinMaxFlag(expectedMinMaxFlag, 2) + spark.sql("DROP TABLE sdkOutputTable") // drop table should not delete the files assert(new File(writerPath).exists()) + clearDataMapCache cleanTestData() } - def buildTestDataMuliBlockLet(records :Int): Unit ={ + def buildTestDataMuliBlockLet(recordsInBlocklet1 :Int, recordsInBlocklet2 :Int): Unit ={ FileUtils.deleteDirectory(new File(writerPath)) val fields=new Array[Field](8) - fields(0)=new Field("myid",DataTypes.INT); - fields(1)=new Field("event_id",DataTypes.STRING); - fields(2)=new Field("eve_time",DataTypes.DATE); - fields(3)=new Field("ingestion_time",DataTypes.TIMESTAMP); - fields(4)=new Field("alldate",DataTypes.createArrayType(DataTypes.DATE)); - fields(5)=new Field("subject",DataTypes.STRING); - fields(6)=new Field("from_email",DataTypes.STRING); - fields(7)=new Field("sal",DataTypes.DOUBLE); + fields(0)=new Field("myid",DataTypes.INT) + fields(1)=new Field("event_id",DataTypes.STRING) + fields(2)=new Field("eve_time",DataTypes.DATE) + fields(3)=new Field("ingestion_time",DataTypes.TIMESTAMP) + fields(4)=new Field("alldate",DataTypes.createArrayType(DataTypes.DATE)) + fields(5)=new Field("subject",DataTypes.STRING) + fields(6)=new Field("from_email",DataTypes.STRING) + fields(7)=new Field("sal",DataTypes.DOUBLE) import scala.collection.JavaConverters._ + val emailDataBlocklet1 = "FromEmail" + val emailDataBlocklet2 = "Email for testing min max for allowed chars" try{ val options=Map("bad_records_action"->"FORCE","complex_delimiter_level_1"->"$").asJava val writer=CarbonWriter.builder().outputPath(writerPath).withBlockletSize(16).sortBy(Array("myid","ingestion_time","event_id")).withLoadOptions(options).buildWriterForCSVInput(new Schema(fields),spark.sessionState.newHadoopConf()) val timeF=new SimpleDateFormat("yyyy-MM-dd HH:mm:ss") val date_F=new SimpleDateFormat("yyyy-MM-dd") - for(i<-1 to records){ + for(i<- 1 to recordsInBlocklet1){ + val time=new Date(System.currentTimeMillis()) + writer.write(Array(""+i,"event_"+i,""+date_F.format(time),""+timeF.format(time),""+date_F.format(time)+"$"+date_F.format(time),"Subject_0",emailDataBlocklet1,""+new Random().nextDouble())) + } + for(i<- 1 to recordsInBlocklet2){ val time=new Date(System.currentTimeMillis()) - writer.write(Array(""+i,"event_"+i,""+date_F.format(time),""+timeF.format(time),""+date_F.format(time)+"$"+date_F.format(time),"Subject_0","FromEmail",""+new Random().nextDouble())) + writer.write(Array(""+i,"event_"+i,""+date_F.format(time),""+timeF.format(time),""+date_F.format(time)+"$"+date_F.format(time),"Subject_0",emailDataBlocklet2,""+new Random().nextDouble())) } writer.close() } } + /** + * read carbon index file and validate the min max flag written in each blocklet + * + * @param expectedMinMaxFlag + * @param numBlocklets + */ + private def validateMinMaxFlag(expectedMinMaxFlag: Array[Array[Boolean]], + numBlocklets: Int): Unit = { + val carbonFiles: Array[File] = new File(writerPath).listFiles() + val carbonIndexFile = carbonFiles.filter(file => file.getName.endsWith(".carbonindex"))(0) + val converter: DataFileFooterConverter = new DataFileFooterConverter(spark.sessionState + .newHadoopConf()) + val carbonIndexFilePath = FileFactory.getUpdatedFilePath(carbonIndexFile.getCanonicalPath) + val indexMetadata: List[DataFileFooter] = converter + .getIndexInfo(carbonIndexFilePath, null, false).asScala.toList + assert(indexMetadata.size == numBlocklets) + indexMetadata.zipWithIndex.foreach { filefooter => + val isMinMaxSet: Array[Boolean] = filefooter._1.getBlockletIndex.getMinMaxIndex.getIsMinMaxSet + assert(isMinMaxSet.sameElements(expectedMinMaxFlag(filefooter._2))) + } + } + + private def clearDataMapCache(): Unit = { + if (!spark.sparkContext.version.startsWith("2.1")) { + val mapSize = DataMapStoreManager.getInstance().getAllDataMaps.size() + DataMapStoreManager.getInstance() + .clearDataMaps(AbsoluteTableIdentifier.from(writerPath)) + assert(mapSize > DataMapStoreManager.getInstance().getAllDataMaps.size()) + } + } + test("Test with long string columns") { FileUtils.deleteDirectory(new File(writerPath)) // here we specify the long string column as varchar @@ -434,6 +494,7 @@ class TestCreateTableUsingSparkCarbonFileFormat extends FunSuite with BeforeAndA val op1 = spark.sql("select address from sdkOutputTableWithoutSchema limit 1").collectAsList() assert(op1.get(0).getString(0).length == 75000) spark.sql("DROP TABLE sdkOutputTableWithoutSchema") + clearDataMapCache cleanTestData() } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/integration/spark-datasource/src/test/scala/org/apache/spark/sql/carbondata/datasource/TestUtil.scala ---------------------------------------------------------------------- diff --git a/integration/spark-datasource/src/test/scala/org/apache/spark/sql/carbondata/datasource/TestUtil.scala b/integration/spark-datasource/src/test/scala/org/apache/spark/sql/carbondata/datasource/TestUtil.scala index 6727ca7..8b0eca8 100644 --- a/integration/spark-datasource/src/test/scala/org/apache/spark/sql/carbondata/datasource/TestUtil.scala +++ b/integration/spark-datasource/src/test/scala/org/apache/spark/sql/carbondata/datasource/TestUtil.scala @@ -25,6 +25,9 @@ import org.apache.spark.sql.{DataFrame, Row, SparkSession} import org.apache.spark.sql.catalyst.plans.logical import org.apache.spark.sql.catalyst.util.sideBySide +import org.apache.carbondata.core.constants.CarbonCommonConstants +import org.apache.carbondata.core.util.CarbonProperties + object TestUtil { val rootPath = new File(this.getClass.getResource("/").getPath @@ -44,6 +47,8 @@ object TestUtil { if (!spark.sparkContext.version.startsWith("2.1")) { spark.experimental.extraOptimizations = Seq(new CarbonFileIndexReplaceRule) } + CarbonProperties.getInstance() + .addProperty(CarbonCommonConstants.CARBON_MINMAX_ALLOWED_BYTE_COUNT, "40") def checkAnswer(df: DataFrame, expectedAnswer: java.util.List[Row]):Unit = { checkAnswer(df, expectedAnswer.asScala) match { http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/integration/spark2/src/main/scala/org/apache/carbondata/stream/CarbonStreamRecordReader.java ---------------------------------------------------------------------- diff --git a/integration/spark2/src/main/scala/org/apache/carbondata/stream/CarbonStreamRecordReader.java b/integration/spark2/src/main/scala/org/apache/carbondata/stream/CarbonStreamRecordReader.java index 6d69eb5..6c65285 100644 --- a/integration/spark2/src/main/scala/org/apache/carbondata/stream/CarbonStreamRecordReader.java +++ b/integration/spark2/src/main/scala/org/apache/carbondata/stream/CarbonStreamRecordReader.java @@ -421,8 +421,9 @@ public class CarbonStreamRecordReader extends RecordReader<Void, Object> { BlockletMinMaxIndex minMaxIndex = CarbonMetadataUtil.convertExternalMinMaxIndex( header.getBlocklet_index().getMin_max_index()); if (minMaxIndex != null) { - BitSet bitSet = - filter.isScanRequired(minMaxIndex.getMaxValues(), minMaxIndex.getMinValues()); + BitSet bitSet = filter + .isScanRequired(minMaxIndex.getMaxValues(), minMaxIndex.getMinValues(), + minMaxIndex.getIsMinMaxSet()); if (bitSet.isEmpty()) { return false; } else { http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/integration/spark2/src/test/scala/org/apache/spark/sql/CarbonGetTableDetailComandTestCase.scala ---------------------------------------------------------------------- diff --git a/integration/spark2/src/test/scala/org/apache/spark/sql/CarbonGetTableDetailComandTestCase.scala b/integration/spark2/src/test/scala/org/apache/spark/sql/CarbonGetTableDetailComandTestCase.scala index a49d5bb..ad6823d 100644 --- a/integration/spark2/src/test/scala/org/apache/spark/sql/CarbonGetTableDetailComandTestCase.scala +++ b/integration/spark2/src/test/scala/org/apache/spark/sql/CarbonGetTableDetailComandTestCase.scala @@ -43,9 +43,9 @@ class CarbonGetTableDetailCommandTestCase extends QueryTest with BeforeAndAfterA assertResult(2)(result.length) assertResult("table_info1")(result(0).getString(0)) // 2087 is the size of carbon table. Note that since 1.5.0, we add additional compressor name in metadata - assertResult(2187)(result(0).getLong(1)) + assertResult(2216)(result(0).getLong(1)) assertResult("table_info2")(result(1).getString(0)) - assertResult(2187)(result(1).getLong(1)) + assertResult(2216)(result(1).getLong(1)) } override def afterAll: Unit = { http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/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 8622fcd..25204cb 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 @@ -327,9 +327,10 @@ public class CarbonFactDataWriterImplV3 extends AbstractFactDataWriter { model.getSegmentProperties().getDimensions().size()); BlockletBTreeIndex bTreeIndex = new BlockletBTreeIndex(index.b_tree_index.getStart_key(), index.b_tree_index.getEnd_key()); - BlockletMinMaxIndex minMaxIndex = new BlockletMinMaxIndex(); - minMaxIndex.setMinValues(toByteArray(index.getMin_max_index().getMin_values())); - minMaxIndex.setMaxValues(toByteArray(index.getMin_max_index().getMax_values())); + BlockletMinMaxIndex minMaxIndex = + new BlockletMinMaxIndex(index.getMin_max_index().getMin_values(), + index.getMin_max_index().getMax_values(), + index.getMin_max_index().getMin_max_presence()); org.apache.carbondata.core.metadata.blocklet.index.BlockletIndex bIndex = new org.apache.carbondata.core.metadata.blocklet.index.BlockletIndex(bTreeIndex, minMaxIndex); http://git-wip-us.apache.org/repos/asf/carbondata/blob/04084c73/streaming/src/main/java/org/apache/carbondata/streaming/segment/StreamSegment.java ---------------------------------------------------------------------- diff --git a/streaming/src/main/java/org/apache/carbondata/streaming/segment/StreamSegment.java b/streaming/src/main/java/org/apache/carbondata/streaming/segment/StreamSegment.java index 744915d..51417c4 100644 --- a/streaming/src/main/java/org/apache/carbondata/streaming/segment/StreamSegment.java +++ b/streaming/src/main/java/org/apache/carbondata/streaming/segment/StreamSegment.java @@ -20,6 +20,7 @@ package org.apache.carbondata.streaming.segment; import java.io.File; import java.io.IOException; import java.util.ArrayList; +import java.util.Arrays; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -257,6 +258,10 @@ public class StreamSegment { CarbonUtil.getValueAsBytes(mrsStats[index].getDataType(), mrsStats[index].getMin()); } minMaxIndex.setMinValues(minIndexes); + // TODO: handle the min max writing for string type based on character limit for streaming + boolean[] isMinMaxSet = new boolean[dimStats.length + mrsStats.length]; + Arrays.fill(isMinMaxSet, true); + minMaxIndex.setIsMinMaxSet(isMinMaxSet); return minMaxIndex; }
