Repository: carbondata Updated Branches: refs/heads/master e0baa9b9f -> 3d3b6ff16
http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/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 803715c..471f9b2 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 @@ -56,7 +56,9 @@ public class CarbonColumnarBatch { actualSize = 0; rowCounter = 0; rowsFiltered = 0; - Arrays.fill(filteredRows, false); + if (filteredRows != null) { + Arrays.fill(filteredRows, false); + } for (int i = 0; i < columnVectors.length; i++) { columnVectors[i].reset(); } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/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 index 50d2ac5..2147c43 100644 --- 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 @@ -27,4 +27,6 @@ public interface CarbonDictionary { void setDictionaryUsed(); byte[] getDictionaryValue(int index); + + byte[][] getAllDictionaryValues(); } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/core/src/main/java/org/apache/carbondata/core/scan/result/vector/ColumnVectorInfo.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/result/vector/ColumnVectorInfo.java b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/ColumnVectorInfo.java index 59117dd..d127728 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/result/vector/ColumnVectorInfo.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/result/vector/ColumnVectorInfo.java @@ -16,7 +16,10 @@ */ package org.apache.carbondata.core.scan.result.vector; +import java.util.BitSet; + import org.apache.carbondata.core.keygenerator.directdictionary.DirectDictionaryGenerator; +import org.apache.carbondata.core.metadata.datatype.DecimalConverterFactory; import org.apache.carbondata.core.scan.filter.GenericQueryType; import org.apache.carbondata.core.scan.model.ProjectionDimension; import org.apache.carbondata.core.scan.model.ProjectionMeasure; @@ -32,6 +35,8 @@ public class ColumnVectorInfo implements Comparable<ColumnVectorInfo> { public DirectDictionaryGenerator directDictionaryGenerator; public MeasureDataVectorProcessor.MeasureVectorFiller measureVectorFiller; public GenericQueryType genericQueryType; + public BitSet deletedRows; + public DecimalConverterFactory.DecimalConverter decimalConverter; @Override public int compareTo(ColumnVectorInfo o) { return ordinal - o.ordinal; http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/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 f8f663f..5dfd6ca 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 @@ -146,7 +146,7 @@ public class CarbonColumnVectorImpl implements CarbonColumnVector { } } - @Override public void putBytes(int rowId, byte[] value) { + @Override public void putByteArray(int rowId, byte[] value) { bytes[rowId] = value; } @@ -160,7 +160,7 @@ public class CarbonColumnVectorImpl implements CarbonColumnVector { } } - @Override public void putBytes(int rowId, int offset, int length, byte[] value) { + @Override public void putByteArray(int rowId, int offset, int length, byte[] value) { bytes[rowId] = new byte[length]; System.arraycopy(value, offset, bytes[rowId], 0, length); } @@ -227,6 +227,31 @@ public class CarbonColumnVectorImpl implements CarbonColumnVector { } } + public Object getDataArray() { + if (dataType == DataTypes.BOOLEAN || dataType == DataTypes.BYTE) { + return byteArr; + } else if (dataType == DataTypes.SHORT) { + return shorts; + } else if (dataType == DataTypes.INT) { + return ints; + } else if (dataType == DataTypes.LONG || dataType == DataTypes.TIMESTAMP) { + return longs; + } else if (dataType == DataTypes.FLOAT) { + return floats; + } else if (dataType == DataTypes.DOUBLE) { + return doubles; + } else if (dataType instanceof DecimalType) { + return decimals; + } else if (dataType == DataTypes.STRING || dataType == DataTypes.BYTE_ARRAY) { + if (null != carbonDictionary) { + return ints; + } + return bytes; + } else { + return data; + } + } + @Override public void reset() { nullBytes.clear(); if (dataType == DataTypes.BOOLEAN || dataType == DataTypes.BYTE) { @@ -287,4 +312,42 @@ public class CarbonColumnVectorImpl implements CarbonColumnVector { * as an optimization to prevent setting nulls. */ public final boolean anyNullsSet() { return anyNullsSet; } + + @Override public void putFloats(int rowId, int count, float[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + floats[rowId ++] = src[i]; + } + } + + @Override public void putShorts(int rowId, int count, short[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + shorts[rowId ++] = src[i]; + } + } + + @Override public void putInts(int rowId, int count, int[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + ints[rowId ++] = src[i]; + } + } + + @Override public void putLongs(int rowId, int count, long[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + longs[rowId ++] = src[i]; + } + } + + @Override public void putDoubles(int rowId, int count, double[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + doubles[rowId ++] = src[i]; + } + } + + @Override public void putBytes(int rowId, int count, byte[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + byteArr[rowId ++] = src[i]; + } + } + + } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/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 index cc3a03c..c8fd573 100644 --- 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 @@ -51,4 +51,7 @@ public class CarbonDictionaryImpl implements CarbonDictionary { return dictionary[index]; } + @Override public byte[][] getAllDictionaryValues() { + return dictionary; + } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/core/src/main/java/org/apache/carbondata/core/scan/scanner/impl/BlockletFullScanner.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/scan/scanner/impl/BlockletFullScanner.java b/core/src/main/java/org/apache/carbondata/core/scan/scanner/impl/BlockletFullScanner.java index 4ec8cb6..62674bc 100644 --- a/core/src/main/java/org/apache/carbondata/core/scan/scanner/impl/BlockletFullScanner.java +++ b/core/src/main/java/org/apache/carbondata/core/scan/scanner/impl/BlockletFullScanner.java @@ -123,7 +123,9 @@ public class BlockletFullScanner implements BlockletScanner { } } scannedResult.setPageFilteredRowCount(numberOfRows); - scannedResult.fillDataChunks(); + if (!blockExecutionInfo.isDirectVectorFill()) { + scannedResult.fillDataChunks(); + } // adding statistics for carbon scan time QueryStatistic scanTime = queryStatisticsModel.getStatisticsTypeAndObjMap() .get(QueryStatisticsConstants.SCAN_BLOCKlET_TIME); http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/core/src/main/java/org/apache/carbondata/core/stats/QueryStatisticsModel.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/stats/QueryStatisticsModel.java b/core/src/main/java/org/apache/carbondata/core/stats/QueryStatisticsModel.java index 9635896..2a294ae 100644 --- a/core/src/main/java/org/apache/carbondata/core/stats/QueryStatisticsModel.java +++ b/core/src/main/java/org/apache/carbondata/core/stats/QueryStatisticsModel.java @@ -20,11 +20,20 @@ package org.apache.carbondata.core.stats; import java.util.HashMap; import java.util.Map; +import org.apache.carbondata.core.constants.CarbonCommonConstants; +import org.apache.carbondata.core.util.CarbonProperties; + public class QueryStatisticsModel { + private QueryStatisticsRecorder recorder; + private Map<String, QueryStatistic> statisticsTypeAndObjMap = new HashMap<String, QueryStatistic>(); + private boolean isEnabled = Boolean.parseBoolean(CarbonProperties.getInstance() + .getProperty(CarbonCommonConstants.ENABLE_QUERY_STATISTICS, + CarbonCommonConstants.ENABLE_QUERY_STATISTICS_DEFAULT)); + public QueryStatisticsRecorder getRecorder() { return recorder; } @@ -36,4 +45,8 @@ public class QueryStatisticsModel { public Map<String, QueryStatistic> getStatisticsTypeAndObjMap() { return statisticsTypeAndObjMap; } + + public boolean isEnabled() { + return isEnabled; + } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/core/src/main/java/org/apache/carbondata/core/util/ByteUtil.java ---------------------------------------------------------------------- diff --git a/core/src/main/java/org/apache/carbondata/core/util/ByteUtil.java b/core/src/main/java/org/apache/carbondata/core/util/ByteUtil.java index 596d1dd..6188948 100644 --- a/core/src/main/java/org/apache/carbondata/core/util/ByteUtil.java +++ b/core/src/main/java/org/apache/carbondata/core/util/ByteUtil.java @@ -733,4 +733,12 @@ public final class ByteUtil { public static float toXorFloat(byte[] value, int offset, int length) { return Float.intBitsToFloat(toXorInt(value, offset, length)); } + + public static int[] toIntArray(byte[] data, int size) { + int[] ints = new int[size]; + for (int i = 0; i < ints.length; i++) { + ints[i] = ByteUtil.valueOf3Bytes(data, i * 3); + } + return ints; + } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/core/src/test/java/org/apache/carbondata/core/scan/filter/executer/IncludeFilterExecuterImplTest.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/carbondata/core/scan/filter/executer/IncludeFilterExecuterImplTest.java b/core/src/test/java/org/apache/carbondata/core/scan/filter/executer/IncludeFilterExecuterImplTest.java index b9e90d6..2662cee 100644 --- a/core/src/test/java/org/apache/carbondata/core/scan/filter/executer/IncludeFilterExecuterImplTest.java +++ b/core/src/test/java/org/apache/carbondata/core/scan/filter/executer/IncludeFilterExecuterImplTest.java @@ -270,11 +270,11 @@ public class IncludeFilterExecuterImplTest extends TestCase { long newTime = 0; long start; long end; - + // dimension's data number in a blocklet, usually default is 32000 - int dataChunkSize = 32000; + int dataChunkSize = 32000; // repeat query times in the test - int queryTimes = 10000; + int queryTimes = 10000; // repeated times for a dictionary value int repeatTimes = 200; //filtered value count in a blocklet http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/core/src/test/java/org/apache/carbondata/core/util/CarbonUtilTest.java ---------------------------------------------------------------------- diff --git a/core/src/test/java/org/apache/carbondata/core/util/CarbonUtilTest.java b/core/src/test/java/org/apache/carbondata/core/util/CarbonUtilTest.java index a4abc61..ecd61bd 100644 --- a/core/src/test/java/org/apache/carbondata/core/util/CarbonUtilTest.java +++ b/core/src/test/java/org/apache/carbondata/core/util/CarbonUtilTest.java @@ -807,7 +807,7 @@ public class CarbonUtilTest { .getFirstIndexUsingBinarySearch(fixedLengthDimensionDataChunk, 1, 3, compareValue, true); assertEquals(2, result); } - + @Test public void testBinaryRangeSearch() { http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/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 b843709..7d6eda0 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 @@ -150,24 +150,24 @@ public class CarbonColumnVectorWrapper implements CarbonColumnVector { } } - @Override public void putBytes(int rowId, byte[] value) { + @Override public void putByteArray(int rowId, byte[] value) { if (!filteredRows[rowId]) { - columnVector.putBytes(counter++, value); + columnVector.putByteArray(counter++, value); } } @Override public void putBytes(int rowId, int count, byte[] value) { for (int i = 0; i < count; i++) { if (!filteredRows[rowId]) { - columnVector.putBytes(counter++, value); + columnVector.putByteArray(counter++, value); } rowId++; } } - @Override public void putBytes(int rowId, int offset, int length, byte[] value) { + @Override public void putByteArray(int rowId, int offset, int length, byte[] value) { if (!filteredRows[rowId]) { - columnVector.putBytes(counter++, offset, length, value); + columnVector.putByteArray(counter++, offset, length, value); } } @@ -246,4 +246,59 @@ public class CarbonColumnVectorWrapper implements CarbonColumnVector { return this.columnVector; } + @Override public void putFloats(int rowId, int count, float[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + columnVector.putFloat(counter++, src[i]); + } + rowId++; + } + } + + @Override public void putShorts(int rowId, int count, short[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + columnVector.putShort(counter++, src[i]); + } + rowId++; + } + } + + @Override public void putInts(int rowId, int count, int[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + columnVector.putInt(counter++, src[i]); + } + rowId++; + } + } + + @Override public void putLongs(int rowId, int count, long[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + columnVector.putLong(counter++, src[i]); + } + rowId++; + } + } + + @Override public void putDoubles(int rowId, int count, double[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + columnVector.putDouble(counter++, src[i]); + } + rowId++; + } + } + + @Override public void putBytes(int rowId, int count, byte[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + columnVector.putByte(counter++, src[i]); + } + rowId++; + } + } + + } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/integration/presto/src/main/java/org/apache/carbondata/presto/readers/SliceStreamReader.java ---------------------------------------------------------------------- diff --git a/integration/presto/src/main/java/org/apache/carbondata/presto/readers/SliceStreamReader.java b/integration/presto/src/main/java/org/apache/carbondata/presto/readers/SliceStreamReader.java index 39fd19a..ab270fc 100644 --- a/integration/presto/src/main/java/org/apache/carbondata/presto/readers/SliceStreamReader.java +++ b/integration/presto/src/main/java/org/apache/carbondata/presto/readers/SliceStreamReader.java @@ -71,11 +71,11 @@ public class SliceStreamReader extends CarbonColumnVectorImpl implements PrestoV values[rowId] = value; } - @Override public void putBytes(int rowId, byte[] value) { + @Override public void putByteArray(int rowId, byte[] value) { type.writeSlice(builder, wrappedBuffer(value)); } - @Override public void putBytes(int rowId, int offset, int length, byte[] value) { + @Override public void putByteArray(int rowId, int offset, int length, byte[] value) { byte[] byteArr = new byte[length]; System.arraycopy(value, offset, byteArr, 0, length); type.writeSlice(builder, wrappedBuffer(byteArr)); http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/integration/spark-common-test/src/test/scala/org/apache/carbondata/spark/testsuite/filterexpr/AllDataTypesTestCaseFilter.scala ---------------------------------------------------------------------- diff --git a/integration/spark-common-test/src/test/scala/org/apache/carbondata/spark/testsuite/filterexpr/AllDataTypesTestCaseFilter.scala b/integration/spark-common-test/src/test/scala/org/apache/carbondata/spark/testsuite/filterexpr/AllDataTypesTestCaseFilter.scala index 73786c8..c96e643 100644 --- a/integration/spark-common-test/src/test/scala/org/apache/carbondata/spark/testsuite/filterexpr/AllDataTypesTestCaseFilter.scala +++ b/integration/spark-common-test/src/test/scala/org/apache/carbondata/spark/testsuite/filterexpr/AllDataTypesTestCaseFilter.scala @@ -57,15 +57,25 @@ class AllDataTypesTestCaseFilter extends QueryTest with BeforeAndAfterAll { test("verify like query ends with filter push down") { val df = sql("select * from alldatatypestableFilter where empname like '%nandh'").queryExecution .sparkPlan - assert(df.asInstanceOf[CarbonDataSourceScan].metadata - .get("PushedFilters").get.contains("CarbonEndsWith")) + if (df.isInstanceOf[CarbonDataSourceScan]) { + assert(df.asInstanceOf[CarbonDataSourceScan].metadata + .get("PushedFilters").get.contains("CarbonEndsWith")) + } else { + assert(df.children.head.asInstanceOf[CarbonDataSourceScan].metadata + .get("PushedFilters").get.contains("CarbonEndsWith")) + } } test("verify like query contains with filter push down") { val df = sql("select * from alldatatypestableFilter where empname like '%nand%'").queryExecution .sparkPlan - assert(df.asInstanceOf[CarbonDataSourceScan].metadata - .get("PushedFilters").get.contains("CarbonContainsWith")) + if (df.isInstanceOf[CarbonDataSourceScan]) { + assert(df.asInstanceOf[CarbonDataSourceScan].metadata + .get("PushedFilters").get.contains("CarbonContainsWith")) + } else { + assert(df.children.head.asInstanceOf[CarbonDataSourceScan].metadata + .get("PushedFilters").get.contains("CarbonContainsWith")) + } } override def afterAll { http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapper.java ---------------------------------------------------------------------- diff --git a/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapper.java b/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapper.java index 5121027..a605134 100644 --- a/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapper.java +++ b/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapper.java @@ -29,7 +29,9 @@ import org.apache.spark.sql.types.Decimal; class ColumnarVectorWrapper implements CarbonColumnVector { - private CarbonVectorProxy sparkColumnVectorProxy; + private CarbonVectorProxy.ColumnVectorProxy sparkColumnVectorProxy; + + private CarbonVectorProxy carbonVectorProxy; private boolean[] filteredRows; @@ -47,8 +49,9 @@ class ColumnarVectorWrapper implements CarbonColumnVector { ColumnarVectorWrapper(CarbonVectorProxy writableColumnVector, boolean[] filteredRows, int ordinal) { - this.sparkColumnVectorProxy = writableColumnVector; + this.sparkColumnVectorProxy = writableColumnVector.getColumnVector(ordinal); this.filteredRows = filteredRows; + this.carbonVectorProxy = writableColumnVector; this.ordinal = ordinal; } @@ -167,7 +170,7 @@ class ColumnarVectorWrapper implements CarbonColumnVector { } } - @Override public void putBytes(int rowId, byte[] value) { + @Override public void putByteArray(int rowId, byte[] value) { if (!filteredRows[rowId]) { sparkColumnVectorProxy.putByteArray(counter++, value, ordinal); } @@ -182,7 +185,7 @@ class ColumnarVectorWrapper implements CarbonColumnVector { } } - @Override public void putBytes(int rowId, int offset, int length, byte[] value) { + @Override public void putByteArray(int rowId, int offset, int length, byte[] value) { if (!filteredRows[rowId]) { sparkColumnVectorProxy.putByteArray(counter++, value, offset, length, ordinal); } @@ -276,12 +279,67 @@ class ColumnarVectorWrapper implements CarbonColumnVector { } public void reserveDictionaryIds() { - sparkColumnVectorProxy.reserveDictionaryIds(sparkColumnVectorProxy.numRows(), ordinal); - dictionaryVector = new ColumnarVectorWrapper(sparkColumnVectorProxy, filteredRows, ordinal); + sparkColumnVectorProxy.reserveDictionaryIds(carbonVectorProxy.numRows(), ordinal); + dictionaryVector = new ColumnarVectorWrapper(carbonVectorProxy, filteredRows, ordinal); ((ColumnarVectorWrapper) dictionaryVector).isDictionary = true; } @Override public CarbonColumnVector getDictionaryVector() { return dictionaryVector; } + + @Override public void putFloats(int rowId, int count, float[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + sparkColumnVectorProxy.putFloat(counter++, src[i], ordinal); + } + rowId++; + } + } + + @Override public void putShorts(int rowId, int count, short[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + sparkColumnVectorProxy.putShort(counter++, src[i], ordinal); + } + rowId++; + } + } + + @Override public void putInts(int rowId, int count, int[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + sparkColumnVectorProxy.putInt(counter++, src[i], ordinal); + } + rowId++; + } + } + + @Override public void putLongs(int rowId, int count, long[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + sparkColumnVectorProxy.putLong(counter++, src[i], ordinal); + } + rowId++; + } + } + + @Override public void putDoubles(int rowId, int count, double[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + sparkColumnVectorProxy.putDouble(counter++, src[i], ordinal); + } + rowId++; + } + } + + @Override public void putBytes(int rowId, int count, byte[] src, int srcIndex) { + for (int i = srcIndex; i < count; i++) { + if (!filteredRows[rowId]) { + sparkColumnVectorProxy.putByte(counter++, src[i], ordinal); + } + rowId++; + } + } + } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapperDirect.java ---------------------------------------------------------------------- diff --git a/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapperDirect.java b/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapperDirect.java new file mode 100644 index 0000000..b55749e --- /dev/null +++ b/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/ColumnarVectorWrapperDirect.java @@ -0,0 +1,229 @@ +/* + * 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 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.spark.sql.CarbonVectorProxy; +import org.apache.spark.sql.carbondata.execution.datasources.CarbonSparkDataSourceUtil; +import org.apache.spark.sql.types.Decimal; + +/** + * Fills the vector directly with out considering any deleted rows. + */ +class ColumnarVectorWrapperDirect implements CarbonColumnVector { + + /** + * It is one column vector adapter class. + */ + protected CarbonVectorProxy.ColumnVectorProxy sparkColumnVectorProxy; + + /** + * It is adapter class of complete ColumnarBatch. + */ + protected CarbonVectorProxy carbonVectorProxy; + + protected int ordinal; + + protected boolean isDictionary; + + private DataType blockDataType; + + private CarbonColumnVector dictionaryVector; + + ColumnarVectorWrapperDirect(CarbonVectorProxy writableColumnVector, int ordinal) { + this.sparkColumnVectorProxy = writableColumnVector.getColumnVector(ordinal); + this.carbonVectorProxy = writableColumnVector; + this.ordinal = ordinal; + } + + @Override public void putBoolean(int rowId, boolean value) { + sparkColumnVectorProxy.putBoolean(rowId, value, ordinal); + } + + @Override public void putFloat(int rowId, float value) { + sparkColumnVectorProxy.putFloat(rowId, value, ordinal); + } + + @Override public void putShort(int rowId, short value) { + sparkColumnVectorProxy.putShort(rowId, value, ordinal); + } + + @Override public void putShorts(int rowId, int count, short value) { + sparkColumnVectorProxy.putShorts(rowId, count, value, ordinal); + } + + @Override public void putInt(int rowId, int value) { + if (isDictionary) { + sparkColumnVectorProxy.putDictionaryInt(rowId, value, ordinal); + } else { + sparkColumnVectorProxy.putInt(rowId, value, ordinal); + } + } + + @Override public void putInts(int rowId, int count, int value) { + sparkColumnVectorProxy.putInts(rowId, count, value, ordinal); + } + + @Override public void putLong(int rowId, long value) { + sparkColumnVectorProxy.putLong(rowId, value, ordinal); + } + + @Override public void putLongs(int rowId, int count, long value) { + sparkColumnVectorProxy.putLongs(rowId, count, value, ordinal); + } + + @Override public void putDecimal(int rowId, BigDecimal value, int precision) { + Decimal toDecimal = Decimal.apply(value); + sparkColumnVectorProxy.putDecimal(rowId, toDecimal, precision, ordinal); + } + + @Override public void putDecimals(int rowId, int count, BigDecimal value, int precision) { + Decimal decimal = Decimal.apply(value); + for (int i = 0; i < count; i++) { + sparkColumnVectorProxy.putDecimal(rowId, decimal, precision, ordinal); + rowId++; + } + } + + @Override public void putDouble(int rowId, double value) { + sparkColumnVectorProxy.putDouble(rowId, value, ordinal); + } + + @Override public void putDoubles(int rowId, int count, double value) { + sparkColumnVectorProxy.putDoubles(rowId, count, value, ordinal); + } + + @Override public void putByteArray(int rowId, byte[] value) { + sparkColumnVectorProxy.putByteArray(rowId, value, ordinal); + } + + @Override + public void putBytes(int rowId, int count, byte[] value) { + for (int i = 0; i < count; i++) { + sparkColumnVectorProxy.putByteArray(rowId, value, ordinal); + rowId++; + } + } + + @Override public void putByteArray(int rowId, int offset, int length, byte[] value) { + sparkColumnVectorProxy.putByteArray(rowId, value, offset, length, ordinal); + } + + @Override public void putNull(int rowId) { + sparkColumnVectorProxy.putNull(rowId, ordinal); + } + + @Override public void putNulls(int rowId, int count) { + sparkColumnVectorProxy.putNulls(rowId, count, ordinal); + } + + @Override public void putNotNull(int rowId) { + sparkColumnVectorProxy.putNotNull(rowId, ordinal); + } + + @Override public void putNotNull(int rowId, int count) { + sparkColumnVectorProxy.putNotNulls(rowId, count, ordinal); + } + + @Override public boolean isNull(int rowId) { + return sparkColumnVectorProxy.isNullAt(rowId, ordinal); + } + + @Override public void putObject(int rowId, Object obj) { + //TODO handle complex types + } + + @Override public Object getData(int rowId) { + //TODO handle complex types + return null; + } + + @Override public void reset() { + if (null != dictionaryVector) { + dictionaryVector.reset(); + } + } + + @Override public DataType getType() { + return CarbonSparkDataSourceUtil + .convertSparkToCarbonDataType(sparkColumnVectorProxy.dataType(ordinal)); + } + + @Override public DataType getBlockDataType() { + return blockDataType; + } + + @Override public void setBlockDataType(DataType blockDataType) { + this.blockDataType = blockDataType; + } + + @Override public void setDictionary(CarbonDictionary dictionary) { + sparkColumnVectorProxy.setDictionary(dictionary, ordinal); + } + + @Override public boolean hasDictionary() { + return sparkColumnVectorProxy.hasDictionary(ordinal); + } + + public void reserveDictionaryIds() { + sparkColumnVectorProxy.reserveDictionaryIds(carbonVectorProxy.numRows(), ordinal); + dictionaryVector = new ColumnarVectorWrapperDirect(carbonVectorProxy, ordinal); + ((ColumnarVectorWrapperDirect) dictionaryVector).isDictionary = true; + } + + @Override public CarbonColumnVector getDictionaryVector() { + return dictionaryVector; + } + + @Override public void putByte(int rowId, byte value) { + sparkColumnVectorProxy.putByte(rowId, value, ordinal); + } + + @Override public void setFilteredRowsExist(boolean filteredRowsExist) { + + } + + @Override public void putFloats(int rowId, int count, float[] src, int srcIndex) { + sparkColumnVectorProxy.putFloats(rowId, count, src, srcIndex); + } + + @Override public void putShorts(int rowId, int count, short[] src, int srcIndex) { + sparkColumnVectorProxy.putShorts(rowId, count, src, srcIndex); + } + + @Override public void putInts(int rowId, int count, int[] src, int srcIndex) { + sparkColumnVectorProxy.putInts(rowId, count, src, srcIndex); + } + + @Override public void putLongs(int rowId, int count, long[] src, int srcIndex) { + sparkColumnVectorProxy.putLongs(rowId, count, src, srcIndex); + } + + @Override public void putDoubles(int rowId, int count, double[] src, int srcIndex) { + sparkColumnVectorProxy.putDoubles(rowId, count, src, srcIndex); + } + + @Override public void putBytes(int rowId, int count, byte[] src, int srcIndex) { + sparkColumnVectorProxy.putBytes(rowId, count, src, srcIndex); + } +} http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/VectorizedCarbonRecordReader.java ---------------------------------------------------------------------- diff --git a/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/VectorizedCarbonRecordReader.java b/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/VectorizedCarbonRecordReader.java index 839a8a0..45686ea 100644 --- a/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/VectorizedCarbonRecordReader.java +++ b/integration/spark-datasource/src/main/scala/org/apache/carbondata/spark/vectorreader/VectorizedCarbonRecordReader.java @@ -26,6 +26,7 @@ import java.util.Map; import org.apache.log4j.Logger; import org.apache.carbondata.common.logging.LogServiceFactory; import org.apache.carbondata.core.cache.dictionary.Dictionary; +import org.apache.carbondata.core.constants.CarbonV3DataFormatConstants; import org.apache.carbondata.core.datastore.block.TableBlockInfo; import org.apache.carbondata.core.keygenerator.directdictionary.DirectDictionaryGenerator; import org.apache.carbondata.core.keygenerator.directdictionary.DirectDictionaryKeyGeneratorFactory; @@ -281,7 +282,12 @@ public class VectorizedCarbonRecordReader extends AbstractRecordReader<Object> { schema = schema.add(field); } } - vectorProxy = new CarbonVectorProxy(DEFAULT_MEMORY_MODE,schema,DEFAULT_BATCH_SIZE); + short batchSize = DEFAULT_BATCH_SIZE; + if (queryModel.isDirectVectorFill()) { + batchSize = CarbonV3DataFormatConstants.NUMBER_OF_ROWS_PER_BLOCKLET_COLUMN_PAGE_DEFAULT; + } + vectorProxy = new CarbonVectorProxy(DEFAULT_MEMORY_MODE, schema, batchSize); + if (partitionColumns != null) { int partitionIdx = fields.length; for (int i = 0; i < partitionColumns.fields().length; i++) { @@ -290,12 +296,22 @@ public class VectorizedCarbonRecordReader extends AbstractRecordReader<Object> { } } CarbonColumnVector[] vectors = new CarbonColumnVector[fields.length]; - boolean[] filteredRows = new boolean[vectorProxy.numRows()]; - for (int i = 0; i < fields.length; i++) { - vectors[i] = new ColumnarVectorWrapper(vectorProxy, filteredRows, i); - if (isNoDictStringField[i]) { - if (vectors[i] instanceof ColumnarVectorWrapper) { - ((ColumnarVectorWrapper) vectors[i]).reserveDictionaryIds(); + boolean[] filteredRows = null; + if (queryModel.isDirectVectorFill()) { + for (int i = 0; i < fields.length; i++) { + vectors[i] = new ColumnarVectorWrapperDirect(vectorProxy, i); + if (isNoDictStringField[i]) { + ((ColumnarVectorWrapperDirect) vectors[i]).reserveDictionaryIds(); + } + } + } else { + filteredRows = new boolean[vectorProxy.numRows()]; + for (int i = 0; i < fields.length; i++) { + vectors[i] = new ColumnarVectorWrapper(vectorProxy, filteredRows, i); + if (isNoDictStringField[i]) { + if (vectors[i] instanceof ColumnarVectorWrapper) { + ((ColumnarVectorWrapper) vectors[i]).reserveDictionaryIds(); + } } } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/integration/spark-datasource/src/main/spark2.1andspark2.2/org/apache/spark/sql/CarbonVectorProxy.java ---------------------------------------------------------------------- diff --git a/integration/spark-datasource/src/main/spark2.1andspark2.2/org/apache/spark/sql/CarbonVectorProxy.java b/integration/spark-datasource/src/main/spark2.1andspark2.2/org/apache/spark/sql/CarbonVectorProxy.java index 80e6dbd..03466cc 100644 --- a/integration/spark-datasource/src/main/spark2.1andspark2.2/org/apache/spark/sql/CarbonVectorProxy.java +++ b/integration/spark-datasource/src/main/spark2.1andspark2.2/org/apache/spark/sql/CarbonVectorProxy.java @@ -45,6 +45,8 @@ public class CarbonVectorProxy { private ColumnarBatch columnarBatch; + private ColumnVectorProxy[] columnVectorProxies; + /** * Adapter class which handles the columnar vector reading of the carbondata * based on the spark ColumnVector and ColumnarBatch API. This proxy class @@ -57,14 +59,22 @@ public class CarbonVectorProxy { */ public CarbonVectorProxy(MemoryMode memMode, int rowNum, StructField[] structFileds) { columnarBatch = ColumnarBatch.allocate(new StructType(structFileds), memMode, rowNum); + columnVectorProxies = new ColumnVectorProxy[columnarBatch.numCols()]; + for (int i = 0; i < columnVectorProxies.length; i++) { + columnVectorProxies[i] = new ColumnVectorProxy(columnarBatch, i); + } } public CarbonVectorProxy(MemoryMode memMode, StructType outputSchema, int rowNum) { columnarBatch = ColumnarBatch.allocate(outputSchema, memMode, rowNum); + columnVectorProxies = new ColumnVectorProxy[columnarBatch.numCols()]; + for (int i = 0; i < columnVectorProxies.length; i++) { + columnVectorProxies[i] = new ColumnVectorProxy(columnarBatch, i); + } } - public ColumnVector getColumnVector(int ordinal) { - return columnarBatch.column(ordinal); + public ColumnVectorProxy getColumnVector(int ordinal) { + return columnVectorProxies[ordinal]; } /** @@ -74,9 +84,6 @@ public class CarbonVectorProxy { columnarBatch.setNumRows(numRows); } - public Object reserveDictionaryIds(int capacity , int ordinal) { - return columnarBatch.column(ordinal).reserveDictionaryIds(capacity); - } /** * Returns the number of rows for read, including filtered rows. @@ -85,22 +92,6 @@ public class CarbonVectorProxy { return columnarBatch.capacity(); } - public void setDictionary(CarbonDictionary dictionary, int ordinal) { - if (null != dictionary) { - columnarBatch.column(ordinal) - .setDictionary(new CarbonDictionaryWrapper(Encoding.PLAIN, dictionary)); - } else { - columnarBatch.column(ordinal).setDictionary(null); - } - } - - public void putNull(int rowId, int ordinal) { - columnarBatch.column(ordinal).putNull(rowId); - } - - public void putNulls(int rowId, int count, int ordinal) { - columnarBatch.column(ordinal).putNulls(rowId, count); - } /** * Called to close all the columns in this batch. It is not valid to access the data after @@ -139,9 +130,7 @@ public class CarbonVectorProxy { return columnarBatch.column(ordinal); } - public boolean hasDictionary(int ordinal) { - return columnarBatch.column(ordinal).hasDictionary(); - } + /** * Resets this column for writing. The currently stored values are no longer accessible. @@ -150,127 +139,187 @@ public class CarbonVectorProxy { columnarBatch.reset(); } - public void putRowToColumnBatch(int rowId, Object value, int offset) { - org.apache.spark.sql.types.DataType t = dataType(offset); - if (null == value) { - putNull(rowId, offset); - } else { - if (t == org.apache.spark.sql.types.DataTypes.BooleanType) { - putBoolean(rowId, (boolean) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.ByteType) { - putByte(rowId, (byte) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.ShortType) { - putShort(rowId, (short) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.IntegerType) { - putInt(rowId, (int) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.LongType) { - putLong(rowId, (long) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.FloatType) { - putFloat(rowId, (float) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.DoubleType) { - putDouble(rowId, (double) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.StringType) { - UTF8String v = (UTF8String) value; - putByteArray(rowId, v.getBytes(), offset); - } else if (t instanceof org.apache.spark.sql.types.DecimalType) { - DecimalType dt = (DecimalType) t; - Decimal d = Decimal.fromDecimal(value); - if (dt.precision() <= Decimal.MAX_INT_DIGITS()) { - putInt(rowId, (int) d.toUnscaledLong(), offset); - } else if (dt.precision() <= Decimal.MAX_LONG_DIGITS()) { - putLong(rowId, d.toUnscaledLong(), offset); - } else { - final BigInteger integer = d.toJavaBigDecimal().unscaledValue(); - byte[] bytes = integer.toByteArray(); - putByteArray(rowId, bytes, 0, bytes.length, offset); + + public static class ColumnVectorProxy { + + private ColumnVector vector; + + public ColumnVectorProxy(ColumnarBatch columnarBatch, int ordinal) { + this.vector = columnarBatch.column(ordinal); + } + + public void putRowToColumnBatch(int rowId, Object value, int offset) { + org.apache.spark.sql.types.DataType t = dataType(offset); + if (null == value) { + putNull(rowId, offset); + } else { + if (t == org.apache.spark.sql.types.DataTypes.BooleanType) { + putBoolean(rowId, (boolean) value, offset); + } else if (t == org.apache.spark.sql.types.DataTypes.ByteType) { + putByte(rowId, (byte) value, offset); + } else if (t == org.apache.spark.sql.types.DataTypes.ShortType) { + putShort(rowId, (short) value, offset); + } else if (t == org.apache.spark.sql.types.DataTypes.IntegerType) { + putInt(rowId, (int) value, offset); + } else if (t == org.apache.spark.sql.types.DataTypes.LongType) { + putLong(rowId, (long) value, offset); + } else if (t == org.apache.spark.sql.types.DataTypes.FloatType) { + putFloat(rowId, (float) value, offset); + } else if (t == org.apache.spark.sql.types.DataTypes.DoubleType) { + putDouble(rowId, (double) value, offset); + } else if (t == org.apache.spark.sql.types.DataTypes.StringType) { + UTF8String v = (UTF8String) value; + putByteArray(rowId, v.getBytes(), offset); + } else if (t instanceof org.apache.spark.sql.types.DecimalType) { + DecimalType dt = (DecimalType) t; + Decimal d = Decimal.fromDecimal(value); + if (dt.precision() <= Decimal.MAX_INT_DIGITS()) { + putInt(rowId, (int) d.toUnscaledLong(), offset); + } else if (dt.precision() <= Decimal.MAX_LONG_DIGITS()) { + putLong(rowId, d.toUnscaledLong(), offset); + } else { + final BigInteger integer = d.toJavaBigDecimal().unscaledValue(); + byte[] bytes = integer.toByteArray(); + putByteArray(rowId, bytes, 0, bytes.length, offset); + } + } else if (t instanceof CalendarIntervalType) { + CalendarInterval c = (CalendarInterval) value; + vector.getChildColumn(0).putInt(rowId, c.months); + vector.getChildColumn(1).putLong(rowId, c.microseconds); + } else if (t instanceof org.apache.spark.sql.types.DateType) { + putInt(rowId, (int) value, offset); + } else if (t instanceof org.apache.spark.sql.types.TimestampType) { + putLong(rowId, (long) value, offset); } - } else if (t instanceof CalendarIntervalType) { - CalendarInterval c = (CalendarInterval) value; - columnarBatch.column(offset).getChildColumn(0).putInt(rowId, c.months); - columnarBatch.column(offset).getChildColumn(1).putLong(rowId, c.microseconds); - } else if (t instanceof org.apache.spark.sql.types.DateType) { - putInt(rowId, (int) value, offset); - } else if (t instanceof org.apache.spark.sql.types.TimestampType) { - putLong(rowId, (long) value, offset); } } - } - public void putBoolean(int rowId, boolean value, int ordinal) { - columnarBatch.column(ordinal).putBoolean(rowId, (boolean) value); - } + public void putBoolean(int rowId, boolean value, int ordinal) { + vector.putBoolean(rowId, value); + } - public void putByte(int rowId, byte value, int ordinal) { - columnarBatch.column(ordinal).putByte(rowId, (byte) value); - } + public void putByte(int rowId, byte value, int ordinal) { + vector.putByte(rowId, value); + } - public void putShort(int rowId, short value, int ordinal) { - columnarBatch.column(ordinal).putShort(rowId, (short) value); - } + public void putBytes(int rowId, int count, byte[] src, int srcIndex) { + vector.putBytes(rowId, count, src, srcIndex); + } - public void putInt(int rowId, int value, int ordinal) { - columnarBatch.column(ordinal).putInt(rowId, (int) value); - } + public void putShort(int rowId, short value, int ordinal) { + vector.putShort(rowId, value); + } - public void putFloat(int rowId, float value, int ordinal) { - columnarBatch.column(ordinal).putFloat(rowId, (float) value); - } + public void putInt(int rowId, int value, int ordinal) { + vector.putInt(rowId, value); + } - public void putLong(int rowId, long value, int ordinal) { - columnarBatch.column(ordinal).putLong(rowId, (long) value); - } + public void putFloat(int rowId, float value, int ordinal) { + vector.putFloat(rowId, value); + } - public void putDouble(int rowId, double value, int ordinal) { - columnarBatch.column(ordinal).putDouble(rowId, (double) value); - } + public void putFloats(int rowId, int count, float[] src, int srcIndex) { + vector.putFloats(rowId, count, src, srcIndex); + } - public void putByteArray(int rowId, byte[] value, int ordinal) { - columnarBatch.column(ordinal).putByteArray(rowId, (byte[]) value); - } + public void putLong(int rowId, long value, int ordinal) { + vector.putLong(rowId, value); + } - public void putInts(int rowId, int count, int value, int ordinal) { - columnarBatch.column(ordinal).putInts(rowId, count, value); - } + public void putDouble(int rowId, double value, int ordinal) { + vector.putDouble(rowId, value); + } - public void putShorts(int rowId, int count, short value, int ordinal) { - columnarBatch.column(ordinal).putShorts(rowId, count, value); - } + public void putByteArray(int rowId, byte[] value, int ordinal) { + vector.putByteArray(rowId, value); + } - public void putLongs(int rowId, int count, long value, int ordinal) { - columnarBatch.column(ordinal).putLongs(rowId, count, value); - } + public void putInts(int rowId, int count, int value, int ordinal) { + vector.putInts(rowId, count, value); + } - public void putDecimal(int rowId, Decimal value, int precision, int ordinal) { - columnarBatch.column(ordinal).putDecimal(rowId, value, precision); + public void putInts(int rowId, int count, int[] src, int srcIndex) { + vector.putInts(rowId, count, src, srcIndex); + } - } + public void putShorts(int rowId, int count, short value, int ordinal) { + vector.putShorts(rowId, count, value); + } - public void putDoubles(int rowId, int count, double value, int ordinal) { - columnarBatch.column(ordinal).putDoubles(rowId, count, value); - } + public void putShorts(int rowId, int count, short[] src, int srcIndex) { + vector.putShorts(rowId, count, src, srcIndex); + } - public void putByteArray(int rowId, byte[] value, int offset, int length, int ordinal) { - columnarBatch.column(ordinal).putByteArray(rowId, (byte[]) value, offset, length); - } + public void putLongs(int rowId, int count, long value, int ordinal) { + vector.putLongs(rowId, count, value); + } - public boolean isNullAt(int rowId, int ordinal) { - return columnarBatch - .column(ordinal).isNullAt(rowId); - } + public void putLongs(int rowId, int count, long[] src, int srcIndex) { + vector.putLongs(rowId, count, src, srcIndex); + } - public DataType dataType(int ordinal) { - return columnarBatch.column(ordinal).dataType(); - } + public void putDecimal(int rowId, Decimal value, int precision, int ordinal) { + vector.putDecimal(rowId, value, precision); - public void putNotNull(int rowId, int ordinal) { - columnarBatch.column(ordinal).putNotNull(rowId); - } + } - public void putNotNulls(int rowId, int count, int ordinal) { - columnarBatch.column(ordinal).putNotNulls(rowId, count); - } + public void putDoubles(int rowId, int count, double value, int ordinal) { + vector.putDoubles(rowId, count, value); + } + + public void putDoubles(int rowId, int count, double[] src, int srcIndex) { + vector.putDoubles(rowId, count, src, srcIndex); + } + + public void putByteArray(int rowId, byte[] value, int offset, int length, int ordinal) { + vector.putByteArray(rowId, value, offset, length); + } + + public boolean isNullAt(int rowId, int ordinal) { + return vector.isNullAt(rowId); + } + + public DataType dataType(int ordinal) { + return vector.dataType(); + } + + public void putNotNull(int rowId, int ordinal) { + vector.putNotNull(rowId); + } + + public void putNotNulls(int rowId, int count, int ordinal) { + vector.putNotNulls(rowId, count); + } + + public void putDictionaryInt(int rowId, int value, int ordinal) { + vector.getDictionaryIds().putInt(rowId, value); + } + + public void setDictionary(CarbonDictionary dictionary, int ordinal) { + if (null != dictionary) { + vector.setDictionary(new CarbonDictionaryWrapper(Encoding.PLAIN, dictionary)); + } else { + vector.setDictionary(null); + } + } + + public void putNull(int rowId, int ordinal) { + vector.putNull(rowId); + } + + public void putNulls(int rowId, int count, int ordinal) { + vector.putNulls(rowId, count); + } + + public boolean hasDictionary(int ordinal) { + return vector.hasDictionary(); + } + + public Object reserveDictionaryIds(int capacity , int ordinal) { + return vector.reserveDictionaryIds(capacity); + } + + - public void putDictionaryInt(int rowId, int value, int ordinal) { - columnarBatch.column(ordinal).getDictionaryIds().putInt(rowId, (int) value); } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/integration/spark-datasource/src/main/spark2.3plus/org/apache/spark/sql/CarbonDictionaryWrapper.java ---------------------------------------------------------------------- diff --git a/integration/spark-datasource/src/main/spark2.3plus/org/apache/spark/sql/CarbonDictionaryWrapper.java b/integration/spark-datasource/src/main/spark2.3plus/org/apache/spark/sql/CarbonDictionaryWrapper.java index 5a99c68..bd8c57c 100644 --- a/integration/spark-datasource/src/main/spark2.3plus/org/apache/spark/sql/CarbonDictionaryWrapper.java +++ b/integration/spark-datasource/src/main/spark2.3plus/org/apache/spark/sql/CarbonDictionaryWrapper.java @@ -28,10 +28,7 @@ public class CarbonDictionaryWrapper implements Dictionary { private byte[][] binaries; CarbonDictionaryWrapper(CarbonDictionary dictionary) { - binaries = new byte[dictionary.getDictionarySize()][]; - for (int i = 0; i < binaries.length; i++) { - binaries[i] = dictionary.getDictionaryValue(i); - } + binaries = dictionary.getAllDictionaryValues(); } @Override public int decodeToInt(int id) { http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/integration/spark-datasource/src/main/spark2.3plus/org/apache/spark/sql/CarbonVectorProxy.java ---------------------------------------------------------------------- diff --git a/integration/spark-datasource/src/main/spark2.3plus/org/apache/spark/sql/CarbonVectorProxy.java b/integration/spark-datasource/src/main/spark2.3plus/org/apache/spark/sql/CarbonVectorProxy.java index 4a0fb9e..bd74b05 100644 --- a/integration/spark-datasource/src/main/spark2.3plus/org/apache/spark/sql/CarbonVectorProxy.java +++ b/integration/spark-datasource/src/main/spark2.3plus/org/apache/spark/sql/CarbonVectorProxy.java @@ -22,7 +22,6 @@ import org.apache.carbondata.core.scan.result.vector.CarbonDictionary; import org.apache.spark.memory.MemoryMode; import org.apache.spark.sql.catalyst.InternalRow; -import org.apache.spark.sql.execution.vectorized.Dictionary; import org.apache.spark.sql.execution.vectorized.WritableColumnVector; import org.apache.spark.sql.types.*; import org.apache.spark.sql.vectorized.ColumnarBatch; @@ -38,7 +37,7 @@ import org.apache.spark.unsafe.types.UTF8String; public class CarbonVectorProxy { private ColumnarBatch columnarBatch; - private WritableColumnVector[] columnVectors; + private ColumnVectorProxy[] columnVectorProxies; /** * Adapter class which handles the columnar vector reading of the carbondata @@ -51,17 +50,25 @@ public class CarbonVectorProxy { * @param structFileds, metadata related to current schema of table. */ public CarbonVectorProxy(MemoryMode memMode, int rowNum, StructField[] structFileds) { - columnVectors = ColumnVectorFactory - .getColumnVector(memMode, new StructType(structFileds), rowNum); + WritableColumnVector[] columnVectors = + ColumnVectorFactory.getColumnVector(memMode, new StructType(structFileds), rowNum); columnarBatch = new ColumnarBatch(columnVectors); columnarBatch.setNumRows(rowNum); + columnVectorProxies = new ColumnVectorProxy[columnarBatch.numCols()]; + for (int i = 0; i < columnVectorProxies.length; i++) { + columnVectorProxies[i] = new ColumnVectorProxy(columnarBatch, i); + } } public CarbonVectorProxy(MemoryMode memMode, StructType outputSchema, int rowNum) { - columnVectors = ColumnVectorFactory + WritableColumnVector[] columnVectors = ColumnVectorFactory .getColumnVector(memMode, outputSchema, rowNum); columnarBatch = new ColumnarBatch(columnVectors); columnarBatch.setNumRows(rowNum); + columnVectorProxies = new ColumnVectorProxy[columnarBatch.numCols()]; + for (int i = 0; i < columnVectorProxies.length; i++) { + columnVectorProxies[i] = new ColumnVectorProxy(columnarBatch, i); + } } /** @@ -71,10 +78,6 @@ public class CarbonVectorProxy { return columnarBatch.numRows(); } - public Object reserveDictionaryIds(int capacity, int ordinal) { - return columnVectors[ordinal].reserveDictionaryIds(capacity); - } - /** * This API will return a columnvector from a batch of column vector rows * based on the ordinal @@ -86,21 +89,20 @@ public class CarbonVectorProxy { return (WritableColumnVector) columnarBatch.column(ordinal); } - public WritableColumnVector getColumnVector(int ordinal) { - return columnVectors[ordinal]; + public ColumnVectorProxy getColumnVector(int ordinal) { + return columnVectorProxies[ordinal]; } - /** * Resets this column for writing. The currently stored values are no longer accessible. */ public void reset() { - for (WritableColumnVector col : columnVectors) { - col.reset(); + for (int i = 0; i < columnarBatch.numCols(); i++) { + ((WritableColumnVector)columnarBatch.column(i)).reset(); } } public void resetDictionaryIds(int ordinal) { - columnVectors[ordinal].getDictionaryIds().reset(); + ((WritableColumnVector)columnarBatch.column(ordinal)).getDictionaryIds().reset(); } /** @@ -133,146 +135,189 @@ public class CarbonVectorProxy { columnarBatch.setNumRows(numRows); } - public void putRowToColumnBatch(int rowId, Object value, int offset) { - org.apache.spark.sql.types.DataType t = dataType(offset); - if (null == value) { - putNull(rowId, offset); - } else { - if (t == org.apache.spark.sql.types.DataTypes.BooleanType) { - putBoolean(rowId, (boolean) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.ByteType) { - putByte(rowId, (byte) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.ShortType) { - putShort(rowId, (short) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.IntegerType) { - putInt(rowId, (int) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.LongType) { - putLong(rowId, (long) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.FloatType) { - putFloat(rowId, (float) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.DoubleType) { - putDouble(rowId, (double) value, offset); - } else if (t == org.apache.spark.sql.types.DataTypes.StringType) { - UTF8String v = (UTF8String) value; - putByteArray(rowId, v.getBytes(), offset); - } else if (t instanceof DecimalType) { - DecimalType dt = (DecimalType) t; - Decimal d = Decimal.fromDecimal(value); - if (dt.precision() <= Decimal.MAX_INT_DIGITS()) { - putInt(rowId, (int) d.toUnscaledLong(), offset); - } else if (dt.precision() <= Decimal.MAX_LONG_DIGITS()) { - putLong(rowId, d.toUnscaledLong(), offset); - } else { - final BigInteger integer = d.toJavaBigDecimal().unscaledValue(); - byte[] bytes = integer.toByteArray(); - putByteArray(rowId, bytes, 0, bytes.length, offset); + + public DataType dataType(int ordinal) { + return columnarBatch.column(ordinal).dataType(); + } + + public static class ColumnVectorProxy { + + private WritableColumnVector vector; + + public ColumnVectorProxy(ColumnarBatch columnarBatch, int ordinal) { + vector = (WritableColumnVector) columnarBatch.column(ordinal); + } + + public void putRowToColumnBatch(int rowId, Object value, int offset) { + DataType t = dataType(offset); + if (null == value) { + putNull(rowId, offset); + } else { + if (t == DataTypes.BooleanType) { + putBoolean(rowId, (boolean) value, offset); + } else if (t == DataTypes.ByteType) { + putByte(rowId, (byte) value, offset); + } else if (t == DataTypes.ShortType) { + putShort(rowId, (short) value, offset); + } else if (t == DataTypes.IntegerType) { + putInt(rowId, (int) value, offset); + } else if (t == DataTypes.LongType) { + putLong(rowId, (long) value, offset); + } else if (t == DataTypes.FloatType) { + putFloat(rowId, (float) value, offset); + } else if (t == DataTypes.DoubleType) { + putDouble(rowId, (double) value, offset); + } else if (t == DataTypes.StringType) { + UTF8String v = (UTF8String) value; + putByteArray(rowId, v.getBytes(), offset); + } else if (t instanceof DecimalType) { + DecimalType dt = (DecimalType) t; + Decimal d = Decimal.fromDecimal(value); + if (dt.precision() <= Decimal.MAX_INT_DIGITS()) { + putInt(rowId, (int) d.toUnscaledLong(), offset); + } else if (dt.precision() <= Decimal.MAX_LONG_DIGITS()) { + putLong(rowId, d.toUnscaledLong(), offset); + } else { + final BigInteger integer = d.toJavaBigDecimal().unscaledValue(); + byte[] bytes = integer.toByteArray(); + putByteArray(rowId, bytes, 0, bytes.length, offset); + } + } else if (t instanceof CalendarIntervalType) { + CalendarInterval c = (CalendarInterval) value; + vector.getChild(0).putInt(rowId, c.months); + vector.getChild(1).putLong(rowId, c.microseconds); + } else if (t instanceof DateType) { + putInt(rowId, (int) value, offset); + } else if (t instanceof TimestampType) { + putLong(rowId, (long) value, offset); } - } else if (t instanceof CalendarIntervalType) { - CalendarInterval c = (CalendarInterval) value; - columnVectors[offset].getChild(0).putInt(rowId, c.months); - columnVectors[offset].getChild(1).putLong(rowId, c.microseconds); - } else if (t instanceof org.apache.spark.sql.types.DateType) { - putInt(rowId, (int) value, offset); - } else if (t instanceof org.apache.spark.sql.types.TimestampType) { - putLong(rowId, (long) value, offset); } } - } - public void putBoolean(int rowId, boolean value, int ordinal) { - columnVectors[ordinal].putBoolean(rowId, (boolean) value); - } + public void putBoolean(int rowId, boolean value, int ordinal) { + vector.putBoolean(rowId, value); + } - public void putByte(int rowId, byte value, int ordinal) { - columnVectors[ordinal].putByte(rowId, (byte) value); - } + public void putByte(int rowId, byte value, int ordinal) { + vector.putByte(rowId, value); + } - public void putShort(int rowId, short value, int ordinal) { - columnVectors[ordinal].putShort(rowId, (short) value); - } + public void putBytes(int rowId, int count, byte[] src, int srcIndex) { + vector.putBytes(rowId, count, src, srcIndex); + } - public void putInt(int rowId, int value, int ordinal) { - columnVectors[ordinal].putInt(rowId, (int) value); - } + public void putShort(int rowId, short value, int ordinal) { + vector.putShort(rowId, value); + } - public void putDictionaryInt(int rowId, int value, int ordinal) { - columnVectors[ordinal].getDictionaryIds().putInt(rowId, (int) value); - } + public void putInt(int rowId, int value, int ordinal) { + vector.putInt(rowId, value); + } - public void putFloat(int rowId, float value, int ordinal) { - columnVectors[ordinal].putFloat(rowId, (float) value); - } + public void putFloat(int rowId, float value, int ordinal) { + vector.putFloat(rowId, value); + } - public void putLong(int rowId, long value, int ordinal) { - columnVectors[ordinal].putLong(rowId, (long) value); - } + public void putFloats(int rowId, int count, float[] src, int srcIndex) { + vector.putFloats(rowId, count, src, srcIndex); + } - public void putDouble(int rowId, double value, int ordinal) { - columnVectors[ordinal].putDouble(rowId, (double) value); - } + public void putLong(int rowId, long value, int ordinal) { + vector.putLong(rowId, value); + } - public void putByteArray(int rowId, byte[] value, int ordinal) { - columnVectors[ordinal].putByteArray(rowId, (byte[]) value); - } + public void putDouble(int rowId, double value, int ordinal) { + vector.putDouble(rowId, value); + } - public void putInts(int rowId, int count, int value, int ordinal) { - columnVectors[ordinal].putInts(rowId, count, value); - } + public void putByteArray(int rowId, byte[] value, int ordinal) { + vector.putByteArray(rowId, value); + } - public void putShorts(int rowId, int count, short value, int ordinal) { - columnVectors[ordinal].putShorts(rowId, count, value); - } + public void putInts(int rowId, int count, int value, int ordinal) { + vector.putInts(rowId, count, value); + } - public void putLongs(int rowId, int count, long value, int ordinal) { - columnVectors[ordinal].putLongs(rowId, count, value); - } + public void putInts(int rowId, int count, int[] src, int srcIndex) { + vector.putInts(rowId, count, src, srcIndex); + } - public void putDecimal(int rowId, Decimal value, int precision, int ordinal) { - columnVectors[ordinal].putDecimal(rowId, value, precision); + public void putShorts(int rowId, int count, short value, int ordinal) { + vector.putShorts(rowId, count, value); + } - } + public void putShorts(int rowId, int count, short[] src, int srcIndex) { + vector.putShorts(rowId, count, src, srcIndex); + } - public void putDoubles(int rowId, int count, double value, int ordinal) { - columnVectors[ordinal].putDoubles(rowId, count, value); - } + public void putLongs(int rowId, int count, long value, int ordinal) { + vector.putLongs(rowId, count, value); + } - public void putByteArray(int rowId, byte[] value, int offset, int length, int ordinal) { - columnVectors[ordinal].putByteArray(rowId, (byte[]) value, offset, length); - } + public void putLongs(int rowId, int count, long[] src, int srcIndex) { + vector.putLongs(rowId, count, src, srcIndex); + } - public void putNull(int rowId, int ordinal) { - columnVectors[ordinal].putNull(rowId); - } + public void putDecimal(int rowId, Decimal value, int precision, int ordinal) { + vector.putDecimal(rowId, value, precision); - public void putNulls(int rowId, int count, int ordinal) { - columnVectors[ordinal].putNulls(rowId, count); - } + } - public void putNotNull(int rowId, int ordinal) { - columnVectors[ordinal].putNotNull(rowId); - } + public void putDoubles(int rowId, int count, double value, int ordinal) { + vector.putDoubles(rowId, count, value); + } - public void putNotNulls(int rowId, int count, int ordinal) { - columnVectors[ordinal].putNotNulls(rowId, count); - } + public void putDoubles(int rowId, int count, double[] src, int srcIndex) { + vector.putDoubles(rowId, count, src, srcIndex); + } - public boolean isNullAt(int rowId, int ordinal) { - return columnVectors[ordinal].isNullAt(rowId); - } + public void putByteArray(int rowId, byte[] value, int offset, int length, int ordinal) { + vector.putByteArray(rowId, value, offset, length); + } - public boolean hasDictionary(int ordinal) { - return columnVectors[ordinal].hasDictionary(); - } + public boolean isNullAt(int rowId, int ordinal) { + return vector.isNullAt(rowId); + } - public void setDictionary(CarbonDictionary dictionary, int ordinal) { + public DataType dataType(int ordinal) { + return vector.dataType(); + } + + public void putNotNull(int rowId, int ordinal) { + vector.putNotNull(rowId); + } + + public void putNotNulls(int rowId, int count, int ordinal) { + vector.putNotNulls(rowId, count); + } + + public void putDictionaryInt(int rowId, int value, int ordinal) { + vector.getDictionaryIds().putInt(rowId, value); + } + + public void setDictionary(CarbonDictionary dictionary, int ordinal) { if (null != dictionary) { - columnVectors[ordinal].setDictionary(new CarbonDictionaryWrapper(dictionary)); + vector.setDictionary(new CarbonDictionaryWrapper(dictionary)); } else { - columnVectors[ordinal].setDictionary(null); + vector.setDictionary(null); + } + } + + public void putNull(int rowId, int ordinal) { + vector.putNull(rowId); + } + + public void putNulls(int rowId, int count, int ordinal) { + vector.putNulls(rowId, count); + } + + public boolean hasDictionary(int ordinal) { + return vector.hasDictionary(); + } + + public Object reserveDictionaryIds(int capacity, int ordinal) { + return vector.reserveDictionaryIds(capacity); } - } - public DataType dataType(int ordinal) { - return columnVectors[ordinal].dataType(); } } http://git-wip-us.apache.org/repos/asf/carbondata/blob/3d3b6ff1/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 6c65285..3330e8b 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 @@ -705,7 +705,7 @@ public class CarbonStreamRecordReader extends RecordReader<Void, Object> { private void putRowToColumnBatch(int rowId) { for (int i = 0; i < projection.length; i++) { Object value = outputValues[i]; - vectorProxy.putRowToColumnBatch(rowId,value,i); + vectorProxy.getColumnVector(i).putRowToColumnBatch(rowId,value,i); } }
