From e7cd54b033cc5fd9c1693222f5ff1ac378092bb9 Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Fri, 24 Jul 2015 16:44:46 -0700 Subject: [PATCH 1/9] Added property to toggle page size check estimation and initial row size checking --- .../parquet/column/ParquetProperties.java | 16 ++++++-- .../column/impl/ColumnWriteStoreV1.java | 13 +++++- .../parquet/column/impl/ColumnWriterV1.java | 17 +++++--- .../column/impl/TestColumnReaderImpl.java | 6 ++- .../hadoop/InternalParquetRecordWriter.java | 39 +++++++++++++++++- .../parquet/hadoop/ParquetOutputFormat.java | 17 ++++++++ .../parquet/hadoop/ParquetRecordWriter.java | 41 ++++++++++++++++++- 7 files changed, 135 insertions(+), 14 deletions(-) diff --git a/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java b/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java index df44c4be4e..a54cf3959f 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java @@ -51,6 +51,9 @@ */ public class ParquetProperties { + public static final int INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK = 100; + public static final boolean DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK = true; + public enum WriterVersion { PARQUET_1_0 ("v1"), PARQUET_2_0 ("v2"); @@ -74,11 +77,15 @@ public static WriterVersion fromString(String name) { private final int dictionaryPageSizeThreshold; private final WriterVersion writerVersion; private final boolean enableDictionary; + private final int initialRowCountForSizeCheck; + private final boolean estimateNextSizeCheck; - public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean enableDict) { + public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean enableDict, int initialRowCountForSizeCheck, boolean estimateNextSizeCheck) { this.dictionaryPageSizeThreshold = dictPageSize; this.writerVersion = writerVersion; this.enableDictionary = enableDict; + this.initialRowCountForSizeCheck = initialRowCountForSizeCheck; + this.estimateNextSizeCheck = estimateNextSizeCheck; } public static ValuesWriter getColumnDescriptorValuesWriter(int maxLevel, int initialSizePerCol, int pageSize) { @@ -228,13 +235,16 @@ public ColumnWriteStore newColumnWriteStore( pageStore, pageSize, dictionaryPageSizeThreshold, - enableDictionary, writerVersion); + enableDictionary, + initialRowCountForSizeCheck, + estimateNextSizeCheck, + writerVersion); case PARQUET_2_0: return new ColumnWriteStoreV2( schema, pageStore, pageSize, - new ParquetProperties(dictionaryPageSizeThreshold, writerVersion, enableDictionary)); + new ParquetProperties(dictionaryPageSizeThreshold, writerVersion, enableDictionary, initialRowCountForSizeCheck, estimateNextSizeCheck)); default: throw new IllegalArgumentException("unknown version " + writerVersion); } diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java index a72b6f7f08..ae6044918a 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java @@ -32,6 +32,9 @@ import org.apache.parquet.column.page.PageWriteStore; import org.apache.parquet.column.page.PageWriter; +import static org.apache.parquet.column.ParquetProperties.INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK; +import static org.apache.parquet.column.ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK; + public class ColumnWriteStoreV1 implements ColumnWriteStore { private final Map columns = new TreeMap(); @@ -39,14 +42,22 @@ public class ColumnWriteStoreV1 implements ColumnWriteStore { private final int pageSizeThreshold; private final int dictionaryPageSizeThreshold; private final boolean enableDictionary; + private final int initialRowCountForSizeCheck; + private final boolean estimateNextSizeCheck; private final WriterVersion writerVersion; public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, WriterVersion writerVersion) { + this(pageWriteStore, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, writerVersion); + } + + public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, int initialRowCountForSizeCheck, boolean estimateNextSizeCheck, WriterVersion writerVersion) { super(); this.pageWriteStore = pageWriteStore; this.pageSizeThreshold = pageSizeThreshold; this.dictionaryPageSizeThreshold = dictionaryPageSizeThreshold; this.enableDictionary = enableDictionary; + this.initialRowCountForSizeCheck = initialRowCountForSizeCheck; + this.estimateNextSizeCheck = estimateNextSizeCheck; this.writerVersion = writerVersion; } @@ -65,7 +76,7 @@ public Set getColumnDescriptors() { private ColumnWriterV1 newMemColumn(ColumnDescriptor path) { PageWriter pageWriter = pageWriteStore.getPageWriter(path); - return new ColumnWriterV1(path, pageWriter, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, writerVersion); + return new ColumnWriterV1(path, pageWriter, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, initialRowCountForSizeCheck, estimateNextSizeCheck, writerVersion); } @Override diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java index f4079c761c..929ac7d8c2 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java @@ -46,7 +46,6 @@ final class ColumnWriterV1 implements ColumnWriter { private static final Log LOG = Log.getLog(ColumnWriterV1.class); private static final boolean DEBUG = Log.DEBUG; - private static final int INITIAL_COUNT_FOR_SIZE_CHECK = 100; private static final int MIN_SLAB_SIZE = 64; private final ColumnDescriptor path; @@ -57,6 +56,7 @@ final class ColumnWriterV1 implements ColumnWriter { private ValuesWriter dataColumn; private int valueCount; private int valueCountForNextSizeCheck; + private boolean estimateNextSizeCheck; private Statistics statistics; @@ -66,15 +66,20 @@ public ColumnWriterV1( int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, + int initialRowCountForSizeCheck, + boolean estimateNextSizeCheck, WriterVersion writerVersion) { this.path = path; this.pageWriter = pageWriter; this.pageSizeThreshold = pageSizeThreshold; // initial check of memory usage. So that we have enough data to make an initial prediction - this.valueCountForNextSizeCheck = INITIAL_COUNT_FOR_SIZE_CHECK; + this.valueCountForNextSizeCheck = initialRowCountForSizeCheck; + // Do not attempt to predict next size check. Prevents issues with rows that vary significantly in size. + this.estimateNextSizeCheck = estimateNextSizeCheck; + resetStatistics(); - ParquetProperties parquetProps = new ParquetProperties(dictionaryPageSizeThreshold, writerVersion, enableDictionary); + ParquetProperties parquetProps = new ParquetProperties(dictionaryPageSizeThreshold, writerVersion, enableDictionary, initialRowCountForSizeCheck, estimateNextSizeCheck); this.repetitionLevelColumn = ParquetProperties.getColumnDescriptorValuesWriter(path.getMaxRepetitionLevel(), MIN_SLAB_SIZE, pageSizeThreshold); this.definitionLevelColumn = ParquetProperties.getColumnDescriptorValuesWriter(path.getMaxDefinitionLevel(), MIN_SLAB_SIZE, pageSizeThreshold); @@ -109,9 +114,11 @@ private void accountForValueWritten() { + dataColumn.getBufferedSize(); if (memSize > pageSizeThreshold) { // we will write the current page and check again the size at the predicted middle of next page - valueCountForNextSizeCheck = valueCount / 2; + if(!estimateNextSizeCheck) { + valueCountForNextSizeCheck = valueCount / 2; + } writePage(); - } else { + } else if (!estimateNextSizeCheck) { // not reached the threshold, will check again midway valueCountForNextSizeCheck = (int)(valueCount + ((float)valueCount * pageSizeThreshold / memSize)) / 2 + 1; } diff --git a/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java b/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java index a1820e654e..bf5a1beb3d 100644 --- a/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java +++ b/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java @@ -19,7 +19,9 @@ package org.apache.parquet.column.impl; import static junit.framework.Assert.assertEquals; +import static org.apache.parquet.column.ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK; import static org.apache.parquet.column.ParquetProperties.WriterVersion.PARQUET_2_0; +import static org.apache.parquet.column.ParquetProperties.INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK; import java.util.List; @@ -58,7 +60,7 @@ public void test() throws Exception { MessageType schema = MessageTypeParser.parseMessageType("message test { required binary foo; }"); ColumnDescriptor col = schema.getColumns().get(0); MemPageWriter pageWriter = new MemPageWriter(); - ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true), 2048); + ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK), 2048); for (int i = 0; i < rows; i++) { columnWriterV2.write(Binary.fromString("bar" + i % 10), 0, 0); if ((i + 1) % 1000 == 0) { @@ -93,7 +95,7 @@ public void testOptional() throws Exception { MessageType schema = MessageTypeParser.parseMessageType("message test { optional binary foo; }"); ColumnDescriptor col = schema.getColumns().get(0); MemPageWriter pageWriter = new MemPageWriter(); - ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true), 2048); + ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK), 2048); for (int i = 0; i < rows; i++) { columnWriterV2.writeNull(0, 0); if ((i + 1) % 1000 == 0) { diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java index 37e8db5b80..6def8c6512 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java @@ -65,6 +65,41 @@ class InternalParquetRecordWriter { private ColumnChunkPageWriteStore pageStore; + /** + * @param parquetFileWriter the file to write to + * @param writeSupport the class to convert incoming records + * @param schema the schema of the records + * @param extraMetaData extra meta data to write in the footer of the file + * @param rowGroupSize the size of a block in the file (this will be approximate) + * @param compressor the codec used to compress + */ + public InternalParquetRecordWriter( + ParquetFileWriter parquetFileWriter, + WriteSupport writeSupport, + MessageType schema, + Map extraMetaData, + long rowGroupSize, + int pageSize, + BytesCompressor compressor, + int dictionaryPageSize, + boolean enableDictionary, + boolean validating, + WriterVersion writerVersion) { + this(parquetFileWriter, + writeSupport, + schema, + extraMetaData, + rowGroupSize, + pageSize, + compressor, + dictionaryPageSize, + enableDictionary, + ParquetProperties.INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, + false, + validating, + writerVersion); + } + /** * @param parquetFileWriter the file to write to * @param writeSupport the class to convert incoming records @@ -83,6 +118,8 @@ public InternalParquetRecordWriter( BytesCompressor compressor, int dictionaryPageSize, boolean enableDictionary, + int initialRowCountForSizeCheck, + boolean constantNextSizeCheck, boolean validating, WriterVersion writerVersion) { this.parquetFileWriter = parquetFileWriter; @@ -95,7 +132,7 @@ public InternalParquetRecordWriter( this.pageSize = pageSize; this.compressor = compressor; this.validating = validating; - this.parquetProperties = new ParquetProperties(dictionaryPageSize, writerVersion, enableDictionary); + this.parquetProperties = new ParquetProperties(dictionaryPageSize, writerVersion, enableDictionary, initialRowCountForSizeCheck, constantNextSizeCheck); initStore(); } diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java index e075db38c4..35476c268a 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java @@ -37,6 +37,7 @@ import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import org.apache.parquet.Log; +import org.apache.parquet.column.ParquetProperties; import org.apache.parquet.column.ParquetProperties.WriterVersion; import org.apache.parquet.hadoop.ParquetFileWriter.Mode; import org.apache.parquet.hadoop.api.WriteSupport; @@ -115,6 +116,8 @@ public class ParquetOutputFormat extends FileOutputFormat { public static final String MEMORY_POOL_RATIO = "parquet.memory.pool.ratio"; public static final String MIN_MEMORY_ALLOCATION = "parquet.memory.min.chunk.size"; public static final String MAX_PADDING_BYTES = "parquet.writer.max-padding"; + public static final String INITIAL_ROW_COUNT_SIZE_CHECK = "parquet.page.size.initial.row.check"; + public static final String ESTIMATE_PAGE_SIZE_CHECK = "parquet.page.size.check.estimate"; // default to no padding for now private static final int DEFAULT_MAX_PADDING_SIZE = 0; @@ -192,6 +195,14 @@ public static boolean getEnableDictionary(Configuration configuration) { return configuration.getBoolean(ENABLE_DICTIONARY, true); } + public static int getInitialRowCountForPageSizeCheck(Configuration configuration) { + return configuration.getInt(INITIAL_ROW_COUNT_SIZE_CHECK, ParquetProperties.INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK); + } + + public static boolean getEstimatePageSizeCheck(Configuration configuration) { + return configuration.getBoolean(ESTIMATE_PAGE_SIZE_CHECK, ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK); + } + @Deprecated public static int getBlockSize(Configuration configuration) { return configuration.getInt(BLOCK_SIZE, DEFAULT_BLOCK_SIZE); @@ -307,6 +318,10 @@ public RecordWriter getRecordWriter(Configuration conf, Path file, Comp if (INFO) LOG.info("Writer version is: " + writerVersion); int maxPaddingSize = getMaxPaddingSize(conf); if (INFO) LOG.info("Maximum row group padding size is " + maxPaddingSize + " bytes"); + boolean estimateNextSizeCheck = getEstimatePageSizeCheck(conf); + if(INFO) LOG.info("Page size checking is: " + (estimateNextSizeCheck ? "estimated" : "constant")); + int initialRowCountForPageSizeCheck = getInitialRowCountForPageSizeCheck(conf); + if(INFO) LOG.info("Initial row count for page size check is: " + initialRowCountForPageSizeCheck); WriteContext init = writeSupport.init(conf); ParquetFileWriter w = new ParquetFileWriter( @@ -333,6 +348,8 @@ public RecordWriter getRecordWriter(Configuration conf, Path file, Comp codecFactory.getCompressor(codec, pageSize), dictionaryPageSize, enableDictionary, + initialRowCountForPageSizeCheck, + estimateNextSizeCheck, validating, writerVersion, memoryManager); diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java index 2449192205..f9d1bcb41a 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java @@ -29,6 +29,8 @@ import org.apache.parquet.schema.MessageType; import static org.apache.parquet.Preconditions.checkNotNull; +import static org.apache.parquet.column.ParquetProperties.INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK; +import static org.apache.parquet.column.ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK; /** * Writes records to a Parquet file @@ -70,9 +72,41 @@ public ParquetRecordWriter( WriterVersion writerVersion) { internalWriter = new InternalParquetRecordWriter(w, writeSupport, schema, extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, - validating, writerVersion); + 100, false, validating, writerVersion); } + /** + * + * @param w the file to write to + * @param writeSupport the class to convert incoming records + * @param schema the schema of the records + * @param extraMetaData extra meta data to write in the footer of the file + * @param blockSize the size of a block in the file (this will be approximate) + * @param compressor the compressor used to compress the pages + * @param dictionaryPageSize the threshold for dictionary size + * @param enableDictionary to enable the dictionary + * @param validating if schema validation should be turned on + */ + public ParquetRecordWriter( + ParquetFileWriter w, + WriteSupport writeSupport, + MessageType schema, + Map extraMetaData, + long blockSize, int pageSize, + BytesCompressor compressor, + int dictionaryPageSize, + boolean enableDictionary, + boolean validating, + WriterVersion writerVersion, + MemoryManager memoryManager) { + internalWriter = new InternalParquetRecordWriter(w, writeSupport, schema, + extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, + INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, + validating, writerVersion); + this.memoryManager = checkNotNull(memoryManager, "memoryManager"); + memoryManager.addWriter(internalWriter, blockSize); + } + /** * * @param w the file to write to @@ -94,11 +128,14 @@ public ParquetRecordWriter( BytesCompressor compressor, int dictionaryPageSize, boolean enableDictionary, + int initialRowCountForSizeCheck, + boolean estimateNextPageSizeCheck, boolean validating, WriterVersion writerVersion, MemoryManager memoryManager) { internalWriter = new InternalParquetRecordWriter(w, writeSupport, schema, - extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, + extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, initialRowCountForSizeCheck, + estimateNextPageSizeCheck, validating, writerVersion); this.memoryManager = checkNotNull(memoryManager, "memoryManager"); memoryManager.addWriter(internalWriter, blockSize); From a057f4616c045e98ea834208924527962d615358 Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Fri, 24 Jul 2015 16:57:12 -0700 Subject: [PATCH 2/9] Fixed inverted property logic --- .../java/org/apache/parquet/column/impl/ColumnWriterV1.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java index 929ac7d8c2..8858fea841 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java @@ -114,11 +114,11 @@ private void accountForValueWritten() { + dataColumn.getBufferedSize(); if (memSize > pageSizeThreshold) { // we will write the current page and check again the size at the predicted middle of next page - if(!estimateNextSizeCheck) { + if(estimateNextSizeCheck) { valueCountForNextSizeCheck = valueCount / 2; } writePage(); - } else if (!estimateNextSizeCheck) { + } else if (estimateNextSizeCheck) { // not reached the threshold, will check again midway valueCountForNextSizeCheck = (int)(valueCount + ((float)valueCount * pageSizeThreshold / memSize)) / 2 + 1; } From b49f03c2581639854f3b89f3c3a1d3d0e0abc0c3 Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Wed, 29 Jul 2015 13:50:35 -0700 Subject: [PATCH 3/9] Fixed reset of nextSizeCheck --- .../java/org/apache/parquet/column/impl/ColumnWriterV1.java | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java index 8858fea841..6c6c9c2467 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java @@ -55,6 +55,7 @@ final class ColumnWriterV1 implements ColumnWriter { private ValuesWriter definitionLevelColumn; private ValuesWriter dataColumn; private int valueCount; + private int initialRowCountForSizeCheck; private int valueCountForNextSizeCheck; private boolean estimateNextSizeCheck; @@ -73,6 +74,7 @@ public ColumnWriterV1( this.pageWriter = pageWriter; this.pageSizeThreshold = pageSizeThreshold; // initial check of memory usage. So that we have enough data to make an initial prediction + this.initialRowCountForSizeCheck = initialRowCountForSizeCheck; this.valueCountForNextSizeCheck = initialRowCountForSizeCheck; // Do not attempt to predict next size check. Prevents issues with rows that vary significantly in size. this.estimateNextSizeCheck = estimateNextSizeCheck; @@ -116,11 +118,15 @@ private void accountForValueWritten() { // we will write the current page and check again the size at the predicted middle of next page if(estimateNextSizeCheck) { valueCountForNextSizeCheck = valueCount / 2; + } else { + valueCountForNextSizeCheck = initialRowCountForSizeCheck; } writePage(); } else if (estimateNextSizeCheck) { // not reached the threshold, will check again midway valueCountForNextSizeCheck = (int)(valueCount + ((float)valueCount * pageSizeThreshold / memSize)) / 2 + 1; + } else { + valueCountForNextSizeCheck += initialRowCountForSizeCheck; } } } From 3f7870c4e1783f1e4e31d5e71ef77096a43cf454 Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Mon, 30 Nov 2015 11:27:01 -0800 Subject: [PATCH 4/9] Rebase to resolve byte buffer conflicts --- .../java/org/apache/parquet/column/ParquetProperties.java | 7 ++++++- .../org/apache/parquet/column/impl/ColumnWriteStoreV1.java | 6 +++--- .../apache/parquet/column/impl/TestColumnReaderImpl.java | 4 ++-- .../apache/parquet/hadoop/InternalParquetRecordWriter.java | 6 ++++-- .../org/apache/parquet/hadoop/ParquetRecordWriter.java | 2 +- 5 files changed, 16 insertions(+), 9 deletions(-) diff --git a/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java b/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java index 7706d77ff8..20ddfaeb7a 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java @@ -90,6 +90,11 @@ public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean } public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean enableDict, ByteBufferAllocator allocator) { + this(dictPageSize, writerVersion, enableDict, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, + DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, allocator); + } + + public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean enableDict, int initialRowCountForSizeCheck, boolean estimateNextSizeCheck, ByteBufferAllocator allocator) { this.dictionaryPageSizeThreshold = dictPageSize; this.writerVersion = writerVersion; this.enableDictionary = enableDict; @@ -104,7 +109,7 @@ public ValuesWriter getColumnDescriptorValuesWriter(int maxLevel, int initialSiz return new DevNullValuesWriter(); } else { return new RunLengthBitPackingHybridValuesWriter( - getWidthFromMaxInt(maxLevel), initialSizePerCol, pageSize); + getWidthFromMaxInt(maxLevel), initialSizePerCol, pageSize, this.allocator); } } diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java index a6cabc8b5a..90ac9408ae 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java @@ -48,11 +48,11 @@ public class ColumnWriteStoreV1 implements ColumnWriteStore { private final WriterVersion writerVersion; private final ByteBufferAllocator allocator; - public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, WriterVersion writerVersion) { - this(pageWriteStore, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, writerVersion); + public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, WriterVersion writerVersion, ByteBufferAllocator allocator) { + this(pageWriteStore, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, writerVersion, allocator); } - public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, int initialRowCountForSizeCheck, boolean estimateNextSizeCheck, WriterVersion writerVersion) { + public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, int initialRowCountForSizeCheck, boolean estimateNextSizeCheck, WriterVersion writerVersion, ByteBufferAllocator allocator) { super(); this.pageWriteStore = pageWriteStore; this.pageSizeThreshold = pageSizeThreshold; diff --git a/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java b/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java index 5c297325b3..caf3e961fd 100644 --- a/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java +++ b/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java @@ -61,7 +61,7 @@ public void test() throws Exception { MessageType schema = MessageTypeParser.parseMessageType("message test { required binary foo; }"); ColumnDescriptor col = schema.getColumns().get(0); MemPageWriter pageWriter = new MemPageWriter(); - ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK), 2048); + ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, new HeapByteBufferAllocator()), 2048); for (int i = 0; i < rows; i++) { columnWriterV2.write(Binary.fromString("bar" + i % 10), 0, 0); if ((i + 1) % 1000 == 0) { @@ -96,7 +96,7 @@ public void testOptional() throws Exception { MessageType schema = MessageTypeParser.parseMessageType("message test { optional binary foo; }"); ColumnDescriptor col = schema.getColumns().get(0); MemPageWriter pageWriter = new MemPageWriter(); - ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK), 2048); + ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, new HeapByteBufferAllocator()), 2048); for (int i = 0; i < rows; i++) { columnWriterV2.writeNull(0, 0); if ((i + 1) % 1000 == 0) { diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java index 0b625ff1f7..92640402c4 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java @@ -87,7 +87,8 @@ public InternalParquetRecordWriter( int dictionaryPageSize, boolean enableDictionary, boolean validating, - WriterVersion writerVersion) { + WriterVersion writerVersion, + ByteBufferAllocator allocator) { this(parquetFileWriter, writeSupport, schema, @@ -100,7 +101,8 @@ public InternalParquetRecordWriter( ParquetProperties.INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, false, validating, - writerVersion); + writerVersion, + allocator); } /** diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java index 62dda8dfdf..9cf0b73f84 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java @@ -138,7 +138,7 @@ public ParquetRecordWriter( internalWriter = new InternalParquetRecordWriter(w, writeSupport, schema, extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, initialRowCountForSizeCheck, estimateNextPageSizeCheck, - validating, writerVersion); + validating, writerVersion, new HeapByteBufferAllocator()); this.memoryManager = checkNotNull(memoryManager, "memoryManager"); memoryManager.addWriter(internalWriter, blockSize); } From 5d99072bea638118740247c408f4065288d70cdf Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Tue, 1 Dec 2015 10:23:04 -0800 Subject: [PATCH 5/9] Update page size checking for v2 writer --- .../parquet/column/ParquetProperties.java | 37 +++++++-- .../column/impl/ColumnWriteStoreV1.java | 12 +-- .../column/impl/ColumnWriteStoreV2.java | 22 ++++-- .../parquet/column/impl/ColumnWriterV1.java | 12 +-- .../column/impl/TestColumnReaderImpl.java | 11 ++- .../hadoop/InternalParquetRecordWriter.java | 9 ++- .../parquet/hadoop/ParquetOutputFormat.java | 20 +++-- .../parquet/hadoop/ParquetRecordWriter.java | 75 ++++++++++--------- 8 files changed, 124 insertions(+), 74 deletions(-) diff --git a/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java b/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java index 20ddfaeb7a..5b670f0c63 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java @@ -55,8 +55,9 @@ */ public class ParquetProperties { - public static final int INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK = 100; public static final boolean DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK = true; + public static final int DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK = 100; + public static final int DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK = 10000; public enum WriterVersion { PARQUET_1_0 ("v1"), @@ -81,7 +82,8 @@ public static WriterVersion fromString(String name) { private final int dictionaryPageSizeThreshold; private final WriterVersion writerVersion; private final boolean enableDictionary; - private final int initialRowCountForSizeCheck; + private final int minRowCountForPageSizeCheck; + private final int maxRowCountForPageSizeCheck; private final boolean estimateNextSizeCheck; private final ByteBufferAllocator allocator; @@ -90,15 +92,17 @@ public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean } public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean enableDict, ByteBufferAllocator allocator) { - this(dictPageSize, writerVersion, enableDict, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, + this(dictPageSize, writerVersion, enableDict, DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, allocator); } - public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean enableDict, int initialRowCountForSizeCheck, boolean estimateNextSizeCheck, ByteBufferAllocator allocator) { + public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean enableDict, int minRowCountForPageSizeCheck, + int maxRowCountForPageSizeCheck, boolean estimateNextSizeCheck, ByteBufferAllocator allocator) { this.dictionaryPageSizeThreshold = dictPageSize; this.writerVersion = writerVersion; this.enableDictionary = enableDict; - this.initialRowCountForSizeCheck = initialRowCountForSizeCheck; + this.minRowCountForPageSizeCheck = minRowCountForPageSizeCheck; + this.maxRowCountForPageSizeCheck = maxRowCountForPageSizeCheck; this.estimateNextSizeCheck = estimateNextSizeCheck; Preconditions.checkNotNull(allocator, "ByteBufferAllocator"); this.allocator = allocator; @@ -257,7 +261,7 @@ public ColumnWriteStore newColumnWriteStore( pageSize, dictionaryPageSizeThreshold, enableDictionary, - initialRowCountForSizeCheck, + minRowCountForPageSizeCheck, estimateNextSizeCheck, writerVersion, allocator); @@ -266,9 +270,28 @@ public ColumnWriteStore newColumnWriteStore( schema, pageStore, pageSize, - new ParquetProperties(dictionaryPageSizeThreshold, writerVersion, enableDictionary, initialRowCountForSizeCheck, estimateNextSizeCheck, allocator)); + new ParquetProperties( + dictionaryPageSizeThreshold, + writerVersion, + enableDictionary, + minRowCountForPageSizeCheck, + maxRowCountForPageSizeCheck, + estimateNextSizeCheck, + allocator)); default: throw new IllegalArgumentException("unknown version " + writerVersion); } } + + public int getMinRowCountForPageSizeCheck() { + return minRowCountForPageSizeCheck; + } + + public int getMaxRowCountForPageSizeCheck() { + return maxRowCountForPageSizeCheck; + } + + public boolean isEstimateNextSizeCheck() { + return estimateNextSizeCheck; + } } diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java index 90ac9408ae..6222b4380d 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java @@ -33,7 +33,7 @@ import org.apache.parquet.column.page.PageWriteStore; import org.apache.parquet.column.page.PageWriter; -import static org.apache.parquet.column.ParquetProperties.INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK; +import static org.apache.parquet.column.ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK; import static org.apache.parquet.column.ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK; public class ColumnWriteStoreV1 implements ColumnWriteStore { @@ -43,22 +43,22 @@ public class ColumnWriteStoreV1 implements ColumnWriteStore { private final int pageSizeThreshold; private final int dictionaryPageSizeThreshold; private final boolean enableDictionary; - private final int initialRowCountForSizeCheck; + private final int initialRowCountForPageSizeCheck; private final boolean estimateNextSizeCheck; private final WriterVersion writerVersion; private final ByteBufferAllocator allocator; public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, WriterVersion writerVersion, ByteBufferAllocator allocator) { - this(pageWriteStore, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, writerVersion, allocator); + this(pageWriteStore, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, writerVersion, allocator); } - public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, int initialRowCountForSizeCheck, boolean estimateNextSizeCheck, WriterVersion writerVersion, ByteBufferAllocator allocator) { + public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, int maxRowCountForPageSizeCheck, boolean estimateNextSizeCheck, WriterVersion writerVersion, ByteBufferAllocator allocator) { super(); this.pageWriteStore = pageWriteStore; this.pageSizeThreshold = pageSizeThreshold; this.dictionaryPageSizeThreshold = dictionaryPageSizeThreshold; this.enableDictionary = enableDictionary; - this.initialRowCountForSizeCheck = initialRowCountForSizeCheck; + this.initialRowCountForPageSizeCheck = maxRowCountForPageSizeCheck; this.estimateNextSizeCheck = estimateNextSizeCheck; this.writerVersion = writerVersion; this.allocator = allocator; @@ -79,7 +79,7 @@ public Set getColumnDescriptors() { private ColumnWriterV1 newMemColumn(ColumnDescriptor path) { PageWriter pageWriter = pageWriteStore.getPageWriter(path); - return new ColumnWriterV1(path, pageWriter, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, initialRowCountForSizeCheck, estimateNextSizeCheck, writerVersion, allocator); + return new ColumnWriterV1(path, pageWriter, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, initialRowCountForPageSizeCheck, estimateNextSizeCheck, writerVersion, allocator); } @Override diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java index 4126004109..64ea40abe0 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java @@ -40,16 +40,14 @@ public class ColumnWriteStoreV2 implements ColumnWriteStore { - // will wait for at least that many records before checking again - private static final int MINIMUM_RECORD_COUNT_FOR_CHECK = 100; - private static final int MAXIMUM_RECORD_COUNT_FOR_CHECK = 10000; // will flush even if size bellow the threshold by this much to facilitate page alignment private static final float THRESHOLD_TOLERANCE_RATIO = 0.1f; // 10 % private final Map columns; private final Collection writers; + private final ParquetProperties parquetProperties; private long rowCount; - private long rowCountForNextSizeCheck = MINIMUM_RECORD_COUNT_FOR_CHECK; + private long rowCountForNextSizeCheck; private final long thresholdTolerance; private final ByteBufferAllocator allocator; @@ -61,6 +59,7 @@ public ColumnWriteStoreV2( int pageSizeThreshold, ParquetProperties parquetProps) { super(); + this.parquetProperties = parquetProps; this.pageSizeThreshold = pageSizeThreshold; this.thresholdTolerance = (long)(pageSizeThreshold * THRESHOLD_TOLERANCE_RATIO); this.allocator = parquetProps.getAllocator(); @@ -71,6 +70,8 @@ public ColumnWriteStoreV2( } this.columns = unmodifiableMap(mcolumns); this.writers = this.columns.values(); + + this.rowCountForNextSizeCheck = parquetProperties.getMinRowCountForPageSizeCheck(); } public ColumnWriter getColumnWriter(ColumnDescriptor path) { @@ -147,6 +148,11 @@ public void endRecord() { } private void sizeCheck() { + if(!parquetProperties.isEstimateNextSizeCheck()) { + rowCountForNextSizeCheck = rowCount + parquetProperties.getMinRowCountForPageSizeCheck(); + return; + } + long minRecordToWait = Long.MAX_VALUE; for (ColumnWriterV2 writer : writers) { long usedMem = writer.getCurrentPageBufferedSize(); @@ -158,20 +164,20 @@ private void sizeCheck() { } long rowsToFillPage = usedMem == 0 ? - MAXIMUM_RECORD_COUNT_FOR_CHECK + parquetProperties.getMaxRowCountForPageSizeCheck() : (long)((float)rows) / usedMem * remainingMem; if (rowsToFillPage < minRecordToWait) { minRecordToWait = rowsToFillPage; } } if (minRecordToWait == Long.MAX_VALUE) { - minRecordToWait = MINIMUM_RECORD_COUNT_FOR_CHECK; + minRecordToWait = parquetProperties.getMinRowCountForPageSizeCheck(); } // will check again halfway rowCountForNextSizeCheck = rowCount + min( - max(minRecordToWait / 2, MINIMUM_RECORD_COUNT_FOR_CHECK), // no less than MINIMUM_RECORD_COUNT_FOR_CHECK - MAXIMUM_RECORD_COUNT_FOR_CHECK); // no more than MAXIMUM_RECORD_COUNT_FOR_CHECK + max(minRecordToWait / 2, parquetProperties.getMinRowCountForPageSizeCheck()), // no less than MINIMUM_RECORD_COUNT_FOR_CHECK + parquetProperties.getMaxRowCountForPageSizeCheck()); // no more than MAXIMUM_RECORD_COUNT_FOR_CHECK } } diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java index 6c24bc5a40..d2ec54b4e2 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java @@ -56,7 +56,7 @@ final class ColumnWriterV1 implements ColumnWriter { private ValuesWriter definitionLevelColumn; private ValuesWriter dataColumn; private int valueCount; - private int initialRowCountForSizeCheck; + private int initialRowCountForPageSizeCheck; private int valueCountForNextSizeCheck; private boolean estimateNextSizeCheck; @@ -68,7 +68,7 @@ public ColumnWriterV1( int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, - int initialRowCountForSizeCheck, + int initialRowCountForPageSizeCheck, boolean estimateNextSizeCheck, WriterVersion writerVersion, ByteBufferAllocator allocator) { @@ -76,8 +76,8 @@ public ColumnWriterV1( this.pageWriter = pageWriter; this.pageSizeThreshold = pageSizeThreshold; // initial check of memory usage. So that we have enough data to make an initial prediction - this.initialRowCountForSizeCheck = initialRowCountForSizeCheck; - this.valueCountForNextSizeCheck = initialRowCountForSizeCheck; + this.initialRowCountForPageSizeCheck = initialRowCountForPageSizeCheck; + this.valueCountForNextSizeCheck = initialRowCountForPageSizeCheck; // Do not attempt to predict next size check. Prevents issues with rows that vary significantly in size. this.estimateNextSizeCheck = estimateNextSizeCheck; @@ -121,14 +121,14 @@ private void accountForValueWritten() { if(estimateNextSizeCheck) { valueCountForNextSizeCheck = valueCount / 2; } else { - valueCountForNextSizeCheck = initialRowCountForSizeCheck; + valueCountForNextSizeCheck = initialRowCountForPageSizeCheck; } writePage(); } else if (estimateNextSizeCheck) { // not reached the threshold, will check again midway valueCountForNextSizeCheck = (int)(valueCount + ((float)valueCount * pageSizeThreshold / memSize)) / 2 + 1; } else { - valueCountForNextSizeCheck += initialRowCountForSizeCheck; + valueCountForNextSizeCheck += initialRowCountForPageSizeCheck; } } } diff --git a/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java b/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java index caf3e961fd..521b668c4a 100644 --- a/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java +++ b/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java @@ -21,7 +21,8 @@ import static junit.framework.Assert.assertEquals; import static org.apache.parquet.column.ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK; import static org.apache.parquet.column.ParquetProperties.WriterVersion.PARQUET_2_0; -import static org.apache.parquet.column.ParquetProperties.INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK; +import static org.apache.parquet.column.ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK; +import static org.apache.parquet.column.ParquetProperties.DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK; import java.util.List; @@ -61,7 +62,9 @@ public void test() throws Exception { MessageType schema = MessageTypeParser.parseMessageType("message test { required binary foo; }"); ColumnDescriptor col = schema.getColumns().get(0); MemPageWriter pageWriter = new MemPageWriter(); - ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, new HeapByteBufferAllocator()), 2048); + ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, + DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, + DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, new HeapByteBufferAllocator()), 2048); for (int i = 0; i < rows; i++) { columnWriterV2.write(Binary.fromString("bar" + i % 10), 0, 0); if ((i + 1) % 1000 == 0) { @@ -96,7 +99,9 @@ public void testOptional() throws Exception { MessageType schema = MessageTypeParser.parseMessageType("message test { optional binary foo; }"); ColumnDescriptor col = schema.getColumns().get(0); MemPageWriter pageWriter = new MemPageWriter(); - ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, new HeapByteBufferAllocator()), 2048); + ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, + DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, + DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, new HeapByteBufferAllocator()), 2048); for (int i = 0; i < rows; i++) { columnWriterV2.writeNull(0, 0); if ((i + 1) % 1000 == 0) { diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java index 92640402c4..0730f2a578 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java @@ -98,7 +98,8 @@ public InternalParquetRecordWriter( compressor, dictionaryPageSize, enableDictionary, - ParquetProperties.INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, + ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, + ParquetProperties.DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, false, validating, writerVersion, @@ -123,7 +124,8 @@ public InternalParquetRecordWriter( BytesCompressor compressor, int dictionaryPageSize, boolean enableDictionary, - int initialRowCountForSizeCheck, + int minRowCountForSizeCheck, + int maxRowCountForSizeCheck, boolean constantNextSizeCheck, boolean validating, WriterVersion writerVersion, @@ -138,7 +140,8 @@ public InternalParquetRecordWriter( this.pageSize = pageSize; this.compressor = compressor; this.validating = validating; - this.parquetProperties = new ParquetProperties(dictionaryPageSize, writerVersion, enableDictionary, initialRowCountForSizeCheck, constantNextSizeCheck, allocator); + this.parquetProperties = new ParquetProperties(dictionaryPageSize, writerVersion, enableDictionary, + minRowCountForSizeCheck, maxRowCountForSizeCheck, constantNextSizeCheck, allocator); initStore(); } diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java index e33793d740..e9d701dc60 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java @@ -142,7 +142,8 @@ public static enum JobSummaryLevel { public static final String MEMORY_POOL_RATIO = "parquet.memory.pool.ratio"; public static final String MIN_MEMORY_ALLOCATION = "parquet.memory.min.chunk.size"; public static final String MAX_PADDING_BYTES = "parquet.writer.max-padding"; - public static final String INITIAL_ROW_COUNT_SIZE_CHECK = "parquet.page.size.initial.row.check"; + public static final String MIN_ROW_COUNT_FOR_PAGE_SIZE_CHECK = "parquet.page.size.row.check.min"; + public static final String MAX_ROW_COUNT_FOR_PAGE_SIZE_CHECK = "parquet.page.size.row.check.max"; public static final String ESTIMATE_PAGE_SIZE_CHECK = "parquet.page.size.check.estimate"; // default to no padding for now @@ -244,8 +245,12 @@ public static boolean getEnableDictionary(Configuration configuration) { return configuration.getBoolean(ENABLE_DICTIONARY, true); } - public static int getInitialRowCountForPageSizeCheck(Configuration configuration) { - return configuration.getInt(INITIAL_ROW_COUNT_SIZE_CHECK, ParquetProperties.INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK); + public static int getMinRowCountForPageSizeCheck(Configuration configuration) { + return configuration.getInt(MIN_ROW_COUNT_FOR_PAGE_SIZE_CHECK, ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK); + } + + public static int getMaxRowCountForPageSizeCheck(Configuration configuration) { + return configuration.getInt(MAX_ROW_COUNT_FOR_PAGE_SIZE_CHECK, ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK); } public static boolean getEstimatePageSizeCheck(Configuration configuration) { @@ -368,8 +373,10 @@ public RecordWriter getRecordWriter(Configuration conf, Path file, Comp if (INFO) LOG.info("Maximum row group padding size is " + maxPaddingSize + " bytes"); boolean estimateNextSizeCheck = getEstimatePageSizeCheck(conf); if(INFO) LOG.info("Page size checking is: " + (estimateNextSizeCheck ? "estimated" : "constant")); - int initialRowCountForPageSizeCheck = getInitialRowCountForPageSizeCheck(conf); - if(INFO) LOG.info("Initial row count for page size check is: " + initialRowCountForPageSizeCheck); + int minRowCountForPageSizeCheck = getMinRowCountForPageSizeCheck(conf); + if(INFO) LOG.info("Min row count for page size check is: " + minRowCountForPageSizeCheck); + int maxRowCountForPageSizeCheck = getMaxRowCountForPageSizeCheck(conf); + if(INFO) LOG.info("Min row count for page size check is: " + maxRowCountForPageSizeCheck); CodecFactory codecFactory = new CodecFactory(conf, pageSize); @@ -398,7 +405,8 @@ public RecordWriter getRecordWriter(Configuration conf, Path file, Comp codecFactory.getCompressor(codec), dictionaryPageSize, enableDictionary, - initialRowCountForPageSizeCheck, + minRowCountForPageSizeCheck, + maxRowCountForPageSizeCheck, estimateNextSizeCheck, validating, writerVersion, diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java index 9cf0b73f84..3b08e53b4a 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java @@ -31,7 +31,8 @@ import org.apache.parquet.schema.MessageType; import static org.apache.parquet.Preconditions.checkNotNull; -import static org.apache.parquet.column.ParquetProperties.INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK; +import static org.apache.parquet.column.ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK; +import static org.apache.parquet.column.ParquetProperties.DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK; import static org.apache.parquet.column.ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK; /** @@ -74,39 +75,39 @@ public ParquetRecordWriter( WriterVersion writerVersion) { internalWriter = new InternalParquetRecordWriter(w, writeSupport, schema, extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, - 100, false, validating, writerVersion, new HeapByteBufferAllocator()); + DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, + DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, validating, writerVersion, new HeapByteBufferAllocator()); } - /** - * - * @param w the file to write to - * @param writeSupport the class to convert incoming records - * @param schema the schema of the records - * @param extraMetaData extra meta data to write in the footer of the file - * @param blockSize the size of a block in the file (this will be approximate) - * @param compressor the compressor used to compress the pages - * @param dictionaryPageSize the threshold for dictionary size - * @param enableDictionary to enable the dictionary - * @param validating if schema validation should be turned on - */ - public ParquetRecordWriter( - ParquetFileWriter w, - WriteSupport writeSupport, - MessageType schema, - Map extraMetaData, - long blockSize, int pageSize, - BytesCompressor compressor, - int dictionaryPageSize, - boolean enableDictionary, - boolean validating, - WriterVersion writerVersion, - MemoryManager memoryManager) { - internalWriter = new InternalParquetRecordWriter(w, writeSupport, schema, - extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, - INITIAL_ROW_COUNT_FOR_PAGE_SIZE_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, - validating, writerVersion, new HeapByteBufferAllocator()); - this.memoryManager = checkNotNull(memoryManager, "memoryManager"); - memoryManager.addWriter(internalWriter, blockSize); + /** + * + * @param w the file to write to + * @param writeSupport the class to convert incoming records + * @param schema the schema of the records + * @param extraMetaData extra meta data to write in the footer of the file + * @param blockSize the size of a block in the file (this will be approximate) + * @param compressor the compressor used to compress the pages + * @param dictionaryPageSize the threshold for dictionary size + * @param enableDictionary to enable the dictionary + * @param validating if schema validation should be turned on + */ + public ParquetRecordWriter( + ParquetFileWriter w, + WriteSupport writeSupport, + MessageType schema, + Map extraMetaData, + long blockSize, int pageSize, + BytesCompressor compressor, + int dictionaryPageSize, + boolean enableDictionary, + boolean validating, + WriterVersion writerVersion, + MemoryManager memoryManager) { + this(w, writeSupport, schema, + extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, + DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, + DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, + validating, writerVersion, memoryManager); } /** @@ -119,6 +120,9 @@ public ParquetRecordWriter( * @param compressor the compressor used to compress the pages * @param dictionaryPageSize the threshold for dictionary size * @param enableDictionary to enable the dictionary + * @param minRowCountForSizeCheck min number of rows before page size check + * @param maxRowCountForSizeCheck max number of rows before page size check + * @param estimateNextPageSizeCheck dynamically estimate next page size check (min will be used if false) * @param validating if schema validation should be turned on */ public ParquetRecordWriter( @@ -130,14 +134,15 @@ public ParquetRecordWriter( BytesCompressor compressor, int dictionaryPageSize, boolean enableDictionary, - int initialRowCountForSizeCheck, + int minRowCountForSizeCheck, + int maxRowCountForSizeCheck, boolean estimateNextPageSizeCheck, boolean validating, WriterVersion writerVersion, MemoryManager memoryManager) { internalWriter = new InternalParquetRecordWriter(w, writeSupport, schema, - extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, initialRowCountForSizeCheck, - estimateNextPageSizeCheck, + extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, minRowCountForSizeCheck, + maxRowCountForSizeCheck, estimateNextPageSizeCheck, validating, writerVersion, new HeapByteBufferAllocator()); this.memoryManager = checkNotNull(memoryManager, "memoryManager"); memoryManager.addWriter(internalWriter, blockSize); From 71336eefde68b60f198e762a87605e3dafe45fa8 Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Wed, 2 Dec 2015 10:00:24 -0800 Subject: [PATCH 6/9] Fixed param name --- .../org/apache/parquet/column/impl/ColumnWriteStoreV1.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java index 6222b4380d..ad9946c6fc 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java @@ -52,13 +52,13 @@ public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, this(pageWriteStore, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, writerVersion, allocator); } - public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, int maxRowCountForPageSizeCheck, boolean estimateNextSizeCheck, WriterVersion writerVersion, ByteBufferAllocator allocator) { + public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, int initialRowCountForPageSizeCheck, boolean estimateNextSizeCheck, WriterVersion writerVersion, ByteBufferAllocator allocator) { super(); this.pageWriteStore = pageWriteStore; this.pageSizeThreshold = pageSizeThreshold; this.dictionaryPageSizeThreshold = dictionaryPageSizeThreshold; this.enableDictionary = enableDictionary; - this.initialRowCountForPageSizeCheck = maxRowCountForPageSizeCheck; + this.initialRowCountForPageSizeCheck = initialRowCountForPageSizeCheck; this.estimateNextSizeCheck = estimateNextSizeCheck; this.writerVersion = writerVersion; this.allocator = allocator; From 2090719c76eaa6ca66cb8fd2ef62a5744f00baff Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Wed, 2 Dec 2015 10:10:58 -0800 Subject: [PATCH 7/9] Update sizeCheck to write page properly if estimating is turned off --- .../column/impl/ColumnWriteStoreV2.java | 20 +++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java index 64ea40abe0..197bc95e0e 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java @@ -148,11 +148,6 @@ public void endRecord() { } private void sizeCheck() { - if(!parquetProperties.isEstimateNextSizeCheck()) { - rowCountForNextSizeCheck = rowCount + parquetProperties.getMinRowCountForPageSizeCheck(); - return; - } - long minRecordToWait = Long.MAX_VALUE; for (ColumnWriterV2 writer : writers) { long usedMem = writer.getCurrentPageBufferedSize(); @@ -173,11 +168,16 @@ private void sizeCheck() { if (minRecordToWait == Long.MAX_VALUE) { minRecordToWait = parquetProperties.getMinRowCountForPageSizeCheck(); } - // will check again halfway - rowCountForNextSizeCheck = rowCount + - min( - max(minRecordToWait / 2, parquetProperties.getMinRowCountForPageSizeCheck()), // no less than MINIMUM_RECORD_COUNT_FOR_CHECK - parquetProperties.getMaxRowCountForPageSizeCheck()); // no more than MAXIMUM_RECORD_COUNT_FOR_CHECK + + if(parquetProperties.isEstimateNextSizeCheck()) { + // will check again halfway + rowCountForNextSizeCheck = rowCount + + min( + max(minRecordToWait / 2, parquetProperties.getMinRowCountForPageSizeCheck()), // no less than MINIMUM_RECORD_COUNT_FOR_CHECK + parquetProperties.getMaxRowCountForPageSizeCheck()); // no more than MAXIMUM_RECORD_COUNT_FOR_CHECK + } else { + rowCountForNextSizeCheck = rowCount + parquetProperties.getMinRowCountForPageSizeCheck(); + } } } From 18f8d3aa72b6058cebbe9a5f1f4d70e131cc0223 Mon Sep 17 00:00:00 2001 From: Daniel Weeks Date: Wed, 2 Dec 2015 10:14:13 -0800 Subject: [PATCH 8/9] Spacing --- .../java/org/apache/parquet/hadoop/ParquetOutputFormat.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java index e9d701dc60..6e28266753 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java @@ -372,11 +372,11 @@ public RecordWriter getRecordWriter(Configuration conf, Path file, Comp int maxPaddingSize = getMaxPaddingSize(conf); if (INFO) LOG.info("Maximum row group padding size is " + maxPaddingSize + " bytes"); boolean estimateNextSizeCheck = getEstimatePageSizeCheck(conf); - if(INFO) LOG.info("Page size checking is: " + (estimateNextSizeCheck ? "estimated" : "constant")); + if (INFO) LOG.info("Page size checking is: " + (estimateNextSizeCheck ? "estimated" : "constant")); int minRowCountForPageSizeCheck = getMinRowCountForPageSizeCheck(conf); - if(INFO) LOG.info("Min row count for page size check is: " + minRowCountForPageSizeCheck); + if (INFO) LOG.info("Min row count for page size check is: " + minRowCountForPageSizeCheck); int maxRowCountForPageSizeCheck = getMaxRowCountForPageSizeCheck(conf); - if(INFO) LOG.info("Min row count for page size check is: " + maxRowCountForPageSizeCheck); + if (INFO) LOG.info("Min row count for page size check is: " + maxRowCountForPageSizeCheck); CodecFactory codecFactory = new CodecFactory(conf, pageSize); From c93b73e6531766333efc172fc28097c660ad2e2d Mon Sep 17 00:00:00 2001 From: Ryan Blue Date: Wed, 2 Dec 2015 11:41:51 -0800 Subject: [PATCH 9/9] PARQUET-99: Use ParquetProperties to carry encoding config. --- .../parquet/column/ParquetProperties.java | 239 +++++++++++++----- .../column/impl/ColumnWriteStoreV1.java | 30 +-- .../column/impl/ColumnWriteStoreV2.java | 40 ++- .../parquet/column/impl/ColumnWriterV1.java | 51 ++-- .../parquet/column/impl/ColumnWriterV2.java | 15 +- .../column/impl/TestColumnReaderImpl.java | 18 +- .../impl/TestCorruptDeltaByteArrays.java | 6 +- .../parquet/column/mem/TestMemColumn.java | 9 +- .../java/org/apache/parquet/io/PerfTest.java | 10 +- .../org/apache/parquet/io/TestColumnIO.java | 10 +- .../org/apache/parquet/io/TestFiltered.java | 10 +- .../hadoop/InternalParquetRecordWriter.java | 65 +---- .../parquet/hadoop/ParquetOutputFormat.java | 68 ++--- .../parquet/hadoop/ParquetRecordWriter.java | 75 +++--- .../apache/parquet/hadoop/ParquetWriter.java | 55 ++-- .../org/apache/parquet/hadoop/TestUtils.java | 7 +- .../parquet/pig/TupleConsumerPerfTest.java | 9 +- .../thrift/TestParquetReadProtocol.java | 10 +- 18 files changed, 381 insertions(+), 346 deletions(-) diff --git a/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java b/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java index 5b670f0c63..0c07d54ff2 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/ParquetProperties.java @@ -20,6 +20,7 @@ import org.apache.parquet.Preconditions; import org.apache.parquet.bytes.ByteBufferAllocator; +import org.apache.parquet.bytes.CapacityByteArrayOutputStream; import org.apache.parquet.bytes.HeapByteBufferAllocator; import static org.apache.parquet.bytes.BytesUtils.getWidthFromMaxInt; @@ -44,6 +45,7 @@ import org.apache.parquet.column.values.plain.BooleanPlainValuesWriter; import org.apache.parquet.column.values.plain.FixedLenByteArrayPlainValuesWriter; import org.apache.parquet.column.values.plain.PlainValuesWriter; +import org.apache.parquet.column.values.rle.RunLengthBitPackingHybridEncoder; import org.apache.parquet.column.values.rle.RunLengthBitPackingHybridValuesWriter; import org.apache.parquet.schema.MessageType; @@ -55,10 +57,16 @@ */ public class ParquetProperties { + public static final int DEFAULT_PAGE_SIZE = 1024 * 1024; + public static final int DEFAULT_DICTIONARY_PAGE_SIZE = DEFAULT_PAGE_SIZE; + public static final boolean DEFAULT_IS_DICTIONARY_ENABLED = true; + public static final WriterVersion DEFAULT_WRITER_VERSION = WriterVersion.PARQUET_1_0; public static final boolean DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK = true; public static final int DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK = 100; public static final int DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK = 10000; + private static final int MIN_SLAB_SIZE = 64; + public enum WriterVersion { PARQUET_1_0 ("v1"), PARQUET_2_0 ("v2"); @@ -79,6 +87,8 @@ public static WriterVersion fromString(String name) { return WriterVersion.valueOf(name); } } + + private final int pageSizeThreshold; private final int dictionaryPageSizeThreshold; private final WriterVersion writerVersion; private final boolean enableDictionary; @@ -87,56 +97,73 @@ public static WriterVersion fromString(String name) { private final boolean estimateNextSizeCheck; private final ByteBufferAllocator allocator; - public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean enableDict) { - this(dictPageSize, writerVersion, enableDict, new HeapByteBufferAllocator()); - } - - public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean enableDict, ByteBufferAllocator allocator) { - this(dictPageSize, writerVersion, enableDict, DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, - DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, allocator); - } + private final int initialSlabSize; - public ParquetProperties(int dictPageSize, WriterVersion writerVersion, boolean enableDict, int minRowCountForPageSizeCheck, - int maxRowCountForPageSizeCheck, boolean estimateNextSizeCheck, ByteBufferAllocator allocator) { + private ParquetProperties(WriterVersion writerVersion, int pageSize, int dictPageSize, boolean enableDict, int minRowCountForPageSizeCheck, + int maxRowCountForPageSizeCheck, boolean estimateNextSizeCheck, ByteBufferAllocator allocator) { + this.pageSizeThreshold = pageSize; + this.initialSlabSize = CapacityByteArrayOutputStream + .initialSlabSizeHeuristic(MIN_SLAB_SIZE, pageSizeThreshold, 10); this.dictionaryPageSizeThreshold = dictPageSize; this.writerVersion = writerVersion; this.enableDictionary = enableDict; this.minRowCountForPageSizeCheck = minRowCountForPageSizeCheck; this.maxRowCountForPageSizeCheck = maxRowCountForPageSizeCheck; this.estimateNextSizeCheck = estimateNextSizeCheck; - Preconditions.checkNotNull(allocator, "ByteBufferAllocator"); this.allocator = allocator; } - public ValuesWriter getColumnDescriptorValuesWriter(int maxLevel, int initialSizePerCol, int pageSize) { + public ValuesWriter newRepetitionLevelWriter(ColumnDescriptor path) { + return newColumnDescriptorValuesWriter(path.getMaxRepetitionLevel()); + } + + public ValuesWriter newDefinitionLevelWriter(ColumnDescriptor path) { + return newColumnDescriptorValuesWriter(path.getMaxDefinitionLevel()); + } + + private ValuesWriter newColumnDescriptorValuesWriter(int maxLevel) { if (maxLevel == 0) { return new DevNullValuesWriter(); } else { return new RunLengthBitPackingHybridValuesWriter( - getWidthFromMaxInt(maxLevel), initialSizePerCol, pageSize, this.allocator); + getWidthFromMaxInt(maxLevel), MIN_SLAB_SIZE, pageSizeThreshold, allocator); } } - private ValuesWriter plainWriter(ColumnDescriptor path, int initialSizePerCol, int pageSize) { + public RunLengthBitPackingHybridEncoder newRepetitionLevelEncoder(ColumnDescriptor path) { + return newLevelEncoder(path.getMaxRepetitionLevel()); + } + + public RunLengthBitPackingHybridEncoder newDefinitionLevelEncoder(ColumnDescriptor path) { + return newLevelEncoder(path.getMaxDefinitionLevel()); + } + + private RunLengthBitPackingHybridEncoder newLevelEncoder(int maxLevel) { + return new RunLengthBitPackingHybridEncoder( + getWidthFromMaxInt(maxLevel), MIN_SLAB_SIZE, pageSizeThreshold, allocator); + } + + private ValuesWriter plainWriter(ColumnDescriptor path) { switch (path.getType()) { case BOOLEAN: return new BooleanPlainValuesWriter(); case INT96: - return new FixedLenByteArrayPlainValuesWriter(12, initialSizePerCol, pageSize, this.allocator); + return new FixedLenByteArrayPlainValuesWriter(12, initialSlabSize, pageSizeThreshold, allocator); case FIXED_LEN_BYTE_ARRAY: - return new FixedLenByteArrayPlainValuesWriter(path.getTypeLength(), initialSizePerCol, pageSize, this.allocator); + return new FixedLenByteArrayPlainValuesWriter(path.getTypeLength(), initialSlabSize, pageSizeThreshold, allocator); case BINARY: case INT32: case INT64: case DOUBLE: case FLOAT: - return new PlainValuesWriter(initialSizePerCol, pageSize, this.allocator); + return new PlainValuesWriter(initialSlabSize, pageSizeThreshold, allocator); default: throw new IllegalArgumentException("Unknown type " + path.getType()); } } - private DictionaryValuesWriter dictionaryWriter(ColumnDescriptor path, int initialSizePerCol) { + @SuppressWarnings("deprecation") + private DictionaryValuesWriter dictionaryWriter(ColumnDescriptor path) { Encoding encodingForDataPage; Encoding encodingForDictionaryPage; switch(writerVersion) { @@ -173,24 +200,24 @@ private DictionaryValuesWriter dictionaryWriter(ColumnDescriptor path, int initi } } - private ValuesWriter writerToFallbackTo(ColumnDescriptor path, int initialSizePerCol, int pageSize) { + private ValuesWriter writerToFallbackTo(ColumnDescriptor path) { switch(writerVersion) { case PARQUET_1_0: - return plainWriter(path, initialSizePerCol, pageSize); + return plainWriter(path); case PARQUET_2_0: switch (path.getType()) { case BOOLEAN: - return new RunLengthBitPackingHybridValuesWriter(1, initialSizePerCol, pageSize, this.allocator); + return new RunLengthBitPackingHybridValuesWriter(1, initialSlabSize, pageSizeThreshold, allocator); case BINARY: case FIXED_LEN_BYTE_ARRAY: - return new DeltaByteArrayWriter(initialSizePerCol, pageSize,this.allocator); + return new DeltaByteArrayWriter(initialSlabSize, pageSizeThreshold, allocator); case INT32: - return new DeltaBinaryPackingValuesWriter(initialSizePerCol, pageSize, this.allocator); + return new DeltaBinaryPackingValuesWriter(initialSlabSize, pageSizeThreshold, allocator); case INT96: case INT64: case DOUBLE: case FLOAT: - return plainWriter(path, initialSizePerCol, pageSize); + return plainWriter(path); default: throw new IllegalArgumentException("Unknown type " + path.getType()); } @@ -199,27 +226,27 @@ private ValuesWriter writerToFallbackTo(ColumnDescriptor path, int initialSizePe } } - private ValuesWriter dictWriterWithFallBack(ColumnDescriptor path, int initialSizePerCol, int pageSize) { - ValuesWriter writerToFallBackTo = writerToFallbackTo(path, initialSizePerCol, pageSize); + private ValuesWriter dictWriterWithFallBack(ColumnDescriptor path) { + ValuesWriter writerToFallBackTo = writerToFallbackTo(path); if (enableDictionary) { return FallbackValuesWriter.of( - dictionaryWriter(path, initialSizePerCol), + dictionaryWriter(path), writerToFallBackTo); } else { return writerToFallBackTo; } } - public ValuesWriter getValuesWriter(ColumnDescriptor path, int initialSizePerCol, int pageSize) { + public ValuesWriter newValuesWriter(ColumnDescriptor path) { switch (path.getType()) { case BOOLEAN: // no dictionary encoding for boolean - return writerToFallbackTo(path, initialSizePerCol, pageSize); + return writerToFallbackTo(path); case FIXED_LEN_BYTE_ARRAY: // dictionary encoding for that type was not enabled in PARQUET 1.0 if (writerVersion == WriterVersion.PARQUET_2_0) { - return dictWriterWithFallBack(path, initialSizePerCol, pageSize); + return dictWriterWithFallBack(path); } else { - return writerToFallbackTo(path, initialSizePerCol, pageSize); + return writerToFallbackTo(path); } case BINARY: case INT32: @@ -227,12 +254,16 @@ public ValuesWriter getValuesWriter(ColumnDescriptor path, int initialSizePerCol case INT96: case DOUBLE: case FLOAT: - return dictWriterWithFallBack(path, initialSizePerCol, pageSize); + return dictWriterWithFallBack(path); default: throw new IllegalArgumentException("Unknown type " + path.getType()); } } + public int getPageSizeThreshold() { + return pageSizeThreshold; + } + public int getDictionaryPageSizeThreshold() { return dictionaryPageSizeThreshold; } @@ -249,35 +280,13 @@ public ByteBufferAllocator getAllocator() { return allocator; } - public ColumnWriteStore newColumnWriteStore( - MessageType schema, - PageWriteStore pageStore, - int pageSize, - ByteBufferAllocator allocator) { + public ColumnWriteStore newColumnWriteStore(MessageType schema, + PageWriteStore pageStore) { switch (writerVersion) { case PARQUET_1_0: - return new ColumnWriteStoreV1( - pageStore, - pageSize, - dictionaryPageSizeThreshold, - enableDictionary, - minRowCountForPageSizeCheck, - estimateNextSizeCheck, - writerVersion, - allocator); + return new ColumnWriteStoreV1(pageStore, this); case PARQUET_2_0: - return new ColumnWriteStoreV2( - schema, - pageStore, - pageSize, - new ParquetProperties( - dictionaryPageSizeThreshold, - writerVersion, - enableDictionary, - minRowCountForPageSizeCheck, - maxRowCountForPageSizeCheck, - estimateNextSizeCheck, - allocator)); + return new ColumnWriteStoreV2(schema, pageStore, this); default: throw new IllegalArgumentException("unknown version " + writerVersion); } @@ -291,7 +300,119 @@ public int getMaxRowCountForPageSizeCheck() { return maxRowCountForPageSizeCheck; } - public boolean isEstimateNextSizeCheck() { + public boolean estimateNextSizeCheck() { return estimateNextSizeCheck; } + + public static Builder builder() { + return new Builder(); + } + + public static Builder copy(ParquetProperties toCopy) { + return new Builder(toCopy); + } + + public static class Builder { + private int pageSize = DEFAULT_PAGE_SIZE; + private int dictPageSize = DEFAULT_DICTIONARY_PAGE_SIZE; + private boolean enableDict = DEFAULT_IS_DICTIONARY_ENABLED; + private WriterVersion writerVersion = DEFAULT_WRITER_VERSION; + private int minRowCountForPageSizeCheck = DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK; + private int maxRowCountForPageSizeCheck = DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK; + private boolean estimateNextSizeCheck = DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK; + private ByteBufferAllocator allocator = new HeapByteBufferAllocator(); + + private Builder() { + } + + private Builder(ParquetProperties toCopy) { + this.enableDict = toCopy.enableDictionary; + this.dictPageSize = toCopy.dictionaryPageSizeThreshold; + this.writerVersion = toCopy.writerVersion; + this.minRowCountForPageSizeCheck = toCopy.minRowCountForPageSizeCheck; + this.maxRowCountForPageSizeCheck = toCopy.maxRowCountForPageSizeCheck; + this.estimateNextSizeCheck = toCopy.estimateNextSizeCheck; + this.allocator = toCopy.allocator; + } + + /** + * Set the Parquet format page size. + * + * @param pageSize an integer size in bytes + * @return this builder for method chaining. + */ + public Builder withPageSize(int pageSize) { + Preconditions.checkArgument(pageSize > 0, + "Invalid page size (negative): %s", pageSize); + this.pageSize = pageSize; + return this; + } + + /** + * Enable or disable dictionary encoding. + * + * @param enableDictionary whether dictionary encoding should be enabled + * @return this builder for method chaining. + */ + public Builder withDictionaryEncoding(boolean enableDictionary) { + this.enableDict = enableDictionary; + return this; + } + + /** + * Set the Parquet format dictionary page size. + * + * @param dictionaryPageSize an integer size in bytes + * @return this builder for method chaining. + */ + public Builder withDictionaryPageSize(int dictionaryPageSize) { + Preconditions.checkArgument(dictionaryPageSize > 0, + "Invalid dictionary page size (negative): %s", dictionaryPageSize); + this.dictPageSize = dictionaryPageSize; + return this; + } + + /** + * Set the {@link WriterVersion format version}. + * + * @param version a {@code WriterVersion} + * @return this builder for method chaining. + */ + public Builder withWriterVersion(WriterVersion version) { + this.writerVersion = version; + return this; + } + + public Builder withMinRowCountForPageSizeCheck(int min) { + Preconditions.checkArgument(min > 0, + "Invalid row count for page size check (negative): %s", min); + this.minRowCountForPageSizeCheck = min; + return this; + } + + public Builder withMaxRowCountForPageSizeCheck(int max) { + Preconditions.checkArgument(max > 0, + "Invalid row count for page size check (negative): %s", max); + this.maxRowCountForPageSizeCheck = max; + return this; + } + + // Do not attempt to predict next size check. Prevents issues with rows that vary significantly in size. + public Builder estimateRowCountForPageSizeCheck(boolean estimateNextSizeCheck) { + this.estimateNextSizeCheck = estimateNextSizeCheck; + return this; + } + + public Builder withAllocator(ByteBufferAllocator allocator) { + Preconditions.checkNotNull(allocator, "ByteBufferAllocator"); + this.allocator = allocator; + return this; + } + + public ParquetProperties build() { + return new ParquetProperties(writerVersion, pageSize, dictPageSize, + enableDict, minRowCountForPageSizeCheck, maxRowCountForPageSizeCheck, + estimateNextSizeCheck, allocator); + } + } } diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java index ad9946c6fc..93a497fad8 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV1.java @@ -29,39 +29,21 @@ import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.column.ColumnWriteStore; import org.apache.parquet.column.ColumnWriter; +import org.apache.parquet.column.ParquetProperties; import org.apache.parquet.column.ParquetProperties.WriterVersion; import org.apache.parquet.column.page.PageWriteStore; import org.apache.parquet.column.page.PageWriter; -import static org.apache.parquet.column.ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK; -import static org.apache.parquet.column.ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK; - public class ColumnWriteStoreV1 implements ColumnWriteStore { private final Map columns = new TreeMap(); private final PageWriteStore pageWriteStore; - private final int pageSizeThreshold; - private final int dictionaryPageSizeThreshold; - private final boolean enableDictionary; - private final int initialRowCountForPageSizeCheck; - private final boolean estimateNextSizeCheck; - private final WriterVersion writerVersion; - private final ByteBufferAllocator allocator; - - public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, WriterVersion writerVersion, ByteBufferAllocator allocator) { - this(pageWriteStore, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, writerVersion, allocator); - } + private final ParquetProperties props; - public ColumnWriteStoreV1(PageWriteStore pageWriteStore, int pageSizeThreshold, int dictionaryPageSizeThreshold, boolean enableDictionary, int initialRowCountForPageSizeCheck, boolean estimateNextSizeCheck, WriterVersion writerVersion, ByteBufferAllocator allocator) { - super(); + public ColumnWriteStoreV1(PageWriteStore pageWriteStore, + ParquetProperties props) { this.pageWriteStore = pageWriteStore; - this.pageSizeThreshold = pageSizeThreshold; - this.dictionaryPageSizeThreshold = dictionaryPageSizeThreshold; - this.enableDictionary = enableDictionary; - this.initialRowCountForPageSizeCheck = initialRowCountForPageSizeCheck; - this.estimateNextSizeCheck = estimateNextSizeCheck; - this.writerVersion = writerVersion; - this.allocator = allocator; + this.props = props; } public ColumnWriter getColumnWriter(ColumnDescriptor path) { @@ -79,7 +61,7 @@ public Set getColumnDescriptors() { private ColumnWriterV1 newMemColumn(ColumnDescriptor path) { PageWriter pageWriter = pageWriteStore.getPageWriter(path); - return new ColumnWriterV1(path, pageWriter, pageSizeThreshold, dictionaryPageSizeThreshold, enableDictionary, initialRowCountForPageSizeCheck, estimateNextSizeCheck, writerVersion, allocator); + return new ColumnWriterV1(path, pageWriter, props); } @Override diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java index 197bc95e0e..7574cedf75 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriteStoreV2.java @@ -29,7 +29,6 @@ import java.util.Set; import java.util.TreeMap; -import org.apache.parquet.bytes.ByteBufferAllocator; import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.column.ColumnWriteStore; import org.apache.parquet.column.ColumnWriter; @@ -45,33 +44,26 @@ public class ColumnWriteStoreV2 implements ColumnWriteStore { private final Map columns; private final Collection writers; - private final ParquetProperties parquetProperties; + private final ParquetProperties props; + private final long thresholdTolerance; private long rowCount; private long rowCountForNextSizeCheck; - private final long thresholdTolerance; - private final ByteBufferAllocator allocator; - - private int pageSizeThreshold; public ColumnWriteStoreV2( MessageType schema, PageWriteStore pageWriteStore, - int pageSizeThreshold, - ParquetProperties parquetProps) { - super(); - this.parquetProperties = parquetProps; - this.pageSizeThreshold = pageSizeThreshold; - this.thresholdTolerance = (long)(pageSizeThreshold * THRESHOLD_TOLERANCE_RATIO); - this.allocator = parquetProps.getAllocator(); + ParquetProperties props) { + this.props = props; + this.thresholdTolerance = (long)(props.getPageSizeThreshold() * THRESHOLD_TOLERANCE_RATIO); Map mcolumns = new TreeMap(); for (ColumnDescriptor path : schema.getColumns()) { PageWriter pageWriter = pageWriteStore.getPageWriter(path); - mcolumns.put(path, new ColumnWriterV2(path, pageWriter, parquetProps, pageSizeThreshold)); + mcolumns.put(path, new ColumnWriterV2(path, pageWriter, props)); } this.columns = unmodifiableMap(mcolumns); this.writers = this.columns.values(); - this.rowCountForNextSizeCheck = parquetProperties.getMinRowCountForPageSizeCheck(); + this.rowCountForNextSizeCheck = props.getMinRowCountForPageSizeCheck(); } public ColumnWriter getColumnWriter(ColumnDescriptor path) { @@ -152,31 +144,31 @@ private void sizeCheck() { for (ColumnWriterV2 writer : writers) { long usedMem = writer.getCurrentPageBufferedSize(); long rows = rowCount - writer.getRowsWrittenSoFar(); - long remainingMem = pageSizeThreshold - usedMem; + long remainingMem = props.getPageSizeThreshold() - usedMem; if (remainingMem <= thresholdTolerance) { writer.writePage(rowCount); - remainingMem = pageSizeThreshold; + remainingMem = props.getPageSizeThreshold(); } long rowsToFillPage = usedMem == 0 ? - parquetProperties.getMaxRowCountForPageSizeCheck() + props.getMaxRowCountForPageSizeCheck() : (long)((float)rows) / usedMem * remainingMem; if (rowsToFillPage < minRecordToWait) { minRecordToWait = rowsToFillPage; } } if (minRecordToWait == Long.MAX_VALUE) { - minRecordToWait = parquetProperties.getMinRowCountForPageSizeCheck(); + minRecordToWait = props.getMinRowCountForPageSizeCheck(); } - if(parquetProperties.isEstimateNextSizeCheck()) { - // will check again halfway + if(props.estimateNextSizeCheck()) { + // will check again halfway if between min and max rowCountForNextSizeCheck = rowCount + min( - max(minRecordToWait / 2, parquetProperties.getMinRowCountForPageSizeCheck()), // no less than MINIMUM_RECORD_COUNT_FOR_CHECK - parquetProperties.getMaxRowCountForPageSizeCheck()); // no more than MAXIMUM_RECORD_COUNT_FOR_CHECK + max(minRecordToWait / 2, props.getMinRowCountForPageSizeCheck()), + props.getMaxRowCountForPageSizeCheck()); } else { - rowCountForNextSizeCheck = rowCount + parquetProperties.getMinRowCountForPageSizeCheck(); + rowCountForNextSizeCheck = rowCount + props.getMinRowCountForPageSizeCheck(); } } diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java index d2ec54b4e2..dc6ebecb5a 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV1.java @@ -23,12 +23,9 @@ import java.io.IOException; import org.apache.parquet.Log; -import org.apache.parquet.bytes.ByteBufferAllocator; -import org.apache.parquet.bytes.CapacityByteArrayOutputStream; import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.column.ColumnWriter; import org.apache.parquet.column.ParquetProperties; -import org.apache.parquet.column.ParquetProperties.WriterVersion; import org.apache.parquet.column.page.DictionaryPage; import org.apache.parquet.column.page.PageWriter; import org.apache.parquet.column.statistics.Statistics; @@ -47,49 +44,33 @@ final class ColumnWriterV1 implements ColumnWriter { private static final Log LOG = Log.getLog(ColumnWriterV1.class); private static final boolean DEBUG = Log.DEBUG; - private static final int MIN_SLAB_SIZE = 64; private final ColumnDescriptor path; private final PageWriter pageWriter; - private final long pageSizeThreshold; + private final ParquetProperties props; + private ValuesWriter repetitionLevelColumn; private ValuesWriter definitionLevelColumn; private ValuesWriter dataColumn; private int valueCount; - private int initialRowCountForPageSizeCheck; private int valueCountForNextSizeCheck; - private boolean estimateNextSizeCheck; private Statistics statistics; - public ColumnWriterV1( - ColumnDescriptor path, - PageWriter pageWriter, - int pageSizeThreshold, - int dictionaryPageSizeThreshold, - boolean enableDictionary, - int initialRowCountForPageSizeCheck, - boolean estimateNextSizeCheck, - WriterVersion writerVersion, - ByteBufferAllocator allocator) { + public ColumnWriterV1(ColumnDescriptor path, PageWriter pageWriter, + ParquetProperties props) { this.path = path; this.pageWriter = pageWriter; - this.pageSizeThreshold = pageSizeThreshold; + this.props = props; + // initial check of memory usage. So that we have enough data to make an initial prediction - this.initialRowCountForPageSizeCheck = initialRowCountForPageSizeCheck; - this.valueCountForNextSizeCheck = initialRowCountForPageSizeCheck; - // Do not attempt to predict next size check. Prevents issues with rows that vary significantly in size. - this.estimateNextSizeCheck = estimateNextSizeCheck; + this.valueCountForNextSizeCheck = props.getMinRowCountForPageSizeCheck(); resetStatistics(); - ParquetProperties parquetProps = new ParquetProperties(dictionaryPageSizeThreshold, writerVersion, enableDictionary, allocator); - - this.repetitionLevelColumn = parquetProps.getColumnDescriptorValuesWriter(path.getMaxRepetitionLevel(), MIN_SLAB_SIZE, pageSizeThreshold); - this.definitionLevelColumn = parquetProps.getColumnDescriptorValuesWriter(path.getMaxDefinitionLevel(), MIN_SLAB_SIZE, pageSizeThreshold); - - int initialSlabSize = CapacityByteArrayOutputStream.initialSlabSizeHeuristic(MIN_SLAB_SIZE, pageSizeThreshold, 10); - this.dataColumn = parquetProps.getValuesWriter(path, initialSlabSize, pageSizeThreshold); + this.repetitionLevelColumn = props.newRepetitionLevelWriter(path); + this.definitionLevelColumn = props.newDefinitionLevelWriter(path); + this.dataColumn = props.newValuesWriter(path); } private void log(Object value, int r, int d) { @@ -116,19 +97,19 @@ private void accountForValueWritten() { long memSize = repetitionLevelColumn.getBufferedSize() + definitionLevelColumn.getBufferedSize() + dataColumn.getBufferedSize(); - if (memSize > pageSizeThreshold) { + if (memSize > props.getPageSizeThreshold()) { // we will write the current page and check again the size at the predicted middle of next page - if(estimateNextSizeCheck) { + if (props.estimateNextSizeCheck()) { valueCountForNextSizeCheck = valueCount / 2; } else { - valueCountForNextSizeCheck = initialRowCountForPageSizeCheck; + valueCountForNextSizeCheck = props.getMinRowCountForPageSizeCheck(); } writePage(); - } else if (estimateNextSizeCheck) { + } else if (props.estimateNextSizeCheck()) { // not reached the threshold, will check again midway - valueCountForNextSizeCheck = (int)(valueCount + ((float)valueCount * pageSizeThreshold / memSize)) / 2 + 1; + valueCountForNextSizeCheck = (int)(valueCount + ((float)valueCount * props.getPageSizeThreshold() / memSize)) / 2 + 1; } else { - valueCountForNextSizeCheck += initialRowCountForPageSizeCheck; + valueCountForNextSizeCheck += props.getMinRowCountForPageSizeCheck(); } } } diff --git a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV2.java b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV2.java index 8249b72077..396d53a1a5 100644 --- a/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV2.java +++ b/parquet-column/src/main/java/org/apache/parquet/column/impl/ColumnWriterV2.java @@ -25,7 +25,6 @@ import org.apache.parquet.Ints; import org.apache.parquet.Log; -import org.apache.parquet.bytes.ByteBufferAllocator; import org.apache.parquet.bytes.BytesInput; import org.apache.parquet.bytes.CapacityByteArrayOutputStream; import org.apache.parquet.column.ColumnDescriptor; @@ -49,7 +48,6 @@ final class ColumnWriterV2 implements ColumnWriter { private static final Log LOG = Log.getLog(ColumnWriterV2.class); private static final boolean DEBUG = Log.DEBUG; - private static final int MIN_SLAB_SIZE = 64; private final ColumnDescriptor path; private final PageWriter pageWriter; @@ -64,19 +62,14 @@ final class ColumnWriterV2 implements ColumnWriter { public ColumnWriterV2( ColumnDescriptor path, PageWriter pageWriter, - ParquetProperties parquetProps, - int pageSize) { + ParquetProperties props) { this.path = path; this.pageWriter = pageWriter; resetStatistics(); - this.repetitionLevelColumn = new RunLengthBitPackingHybridEncoder( - getWidthFromMaxInt(path.getMaxRepetitionLevel()), MIN_SLAB_SIZE, pageSize, parquetProps.getAllocator()); - this.definitionLevelColumn = new RunLengthBitPackingHybridEncoder( - getWidthFromMaxInt(path.getMaxDefinitionLevel()), MIN_SLAB_SIZE, pageSize, parquetProps.getAllocator()); - - int initialSlabSize = CapacityByteArrayOutputStream.initialSlabSizeHeuristic(MIN_SLAB_SIZE, pageSize, 10); - this.dataColumn = parquetProps.getValuesWriter(path, initialSlabSize, pageSize); + this.repetitionLevelColumn = props.newRepetitionLevelEncoder(path); + this.definitionLevelColumn = props.newDefinitionLevelEncoder(path); + this.dataColumn = props.newValuesWriter(path); } private void log(Object value, int r, int d) { diff --git a/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java b/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java index 521b668c4a..d2d78c43d1 100644 --- a/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java +++ b/parquet-column/src/test/java/org/apache/parquet/column/impl/TestColumnReaderImpl.java @@ -19,16 +19,12 @@ package org.apache.parquet.column.impl; import static junit.framework.Assert.assertEquals; -import static org.apache.parquet.column.ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK; import static org.apache.parquet.column.ParquetProperties.WriterVersion.PARQUET_2_0; -import static org.apache.parquet.column.ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK; -import static org.apache.parquet.column.ParquetProperties.DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK; import java.util.List; import org.apache.parquet.Version; import org.apache.parquet.VersionParser; -import org.apache.parquet.bytes.HeapByteBufferAllocator; import org.junit.Test; import org.apache.parquet.column.ColumnDescriptor; @@ -62,9 +58,10 @@ public void test() throws Exception { MessageType schema = MessageTypeParser.parseMessageType("message test { required binary foo; }"); ColumnDescriptor col = schema.getColumns().get(0); MemPageWriter pageWriter = new MemPageWriter(); - ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, - DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, - DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, new HeapByteBufferAllocator()), 2048); + ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, + ParquetProperties.builder() + .withDictionaryPageSize(1024).withWriterVersion(PARQUET_2_0) + .withPageSize(2048).build()); for (int i = 0; i < rows; i++) { columnWriterV2.write(Binary.fromString("bar" + i % 10), 0, 0); if ((i + 1) % 1000 == 0) { @@ -99,9 +96,10 @@ public void testOptional() throws Exception { MessageType schema = MessageTypeParser.parseMessageType("message test { optional binary foo; }"); ColumnDescriptor col = schema.getColumns().get(0); MemPageWriter pageWriter = new MemPageWriter(); - ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, new ParquetProperties(1024, PARQUET_2_0, true, - DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, - DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, new HeapByteBufferAllocator()), 2048); + ColumnWriterV2 columnWriterV2 = new ColumnWriterV2(col, pageWriter, + ParquetProperties.builder() + .withDictionaryPageSize(1024).withWriterVersion(PARQUET_2_0) + .withPageSize(2048).build()); for (int i = 0; i < rows; i++) { columnWriterV2.writeNull(0, 0); if ((i + 1) % 1000 == 0) { diff --git a/parquet-column/src/test/java/org/apache/parquet/column/impl/TestCorruptDeltaByteArrays.java b/parquet-column/src/test/java/org/apache/parquet/column/impl/TestCorruptDeltaByteArrays.java index 9bb2759a10..1f39d95b82 100644 --- a/parquet-column/src/test/java/org/apache/parquet/column/impl/TestCorruptDeltaByteArrays.java +++ b/parquet-column/src/test/java/org/apache/parquet/column/impl/TestCorruptDeltaByteArrays.java @@ -191,10 +191,12 @@ public void testColumnReaderImplWithCorruptPage() throws Exception { MemPageStore pages = new MemPageStore(0); PageWriter memWriter = pages.getPageWriter(column); - ParquetProperties parquetProps = new ParquetProperties(0, ParquetProperties.WriterVersion.PARQUET_1_0, false, new HeapByteBufferAllocator()); + ParquetProperties parquetProps = ParquetProperties.builder() + .withDictionaryEncoding(false) + .build(); // get generic repetition and definition level bytes to use for pages - ValuesWriter rdValues = parquetProps.getColumnDescriptorValuesWriter(0, 10, 100); + ValuesWriter rdValues = parquetProps.newDefinitionLevelWriter(column); for (int i = 0; i < 10; i += 1) { rdValues.writeInteger(0); } diff --git a/parquet-column/src/test/java/org/apache/parquet/column/mem/TestMemColumn.java b/parquet-column/src/test/java/org/apache/parquet/column/mem/TestMemColumn.java index 044fe2af0f..42c1776cc7 100644 --- a/parquet-column/src/test/java/org/apache/parquet/column/mem/TestMemColumn.java +++ b/parquet-column/src/test/java/org/apache/parquet/column/mem/TestMemColumn.java @@ -20,14 +20,13 @@ import static org.junit.Assert.assertEquals; -import org.apache.parquet.bytes.HeapByteBufferAllocator; +import org.apache.parquet.column.ParquetProperties; import org.junit.Test; import org.apache.parquet.Log; import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.column.ColumnReader; import org.apache.parquet.column.ColumnWriter; -import org.apache.parquet.column.ParquetProperties.WriterVersion; import org.apache.parquet.column.impl.ColumnReadStoreImpl; import org.apache.parquet.column.impl.ColumnWriteStoreV1; import org.apache.parquet.column.page.mem.MemPageStore; @@ -161,6 +160,10 @@ public void testMemColumnSeveralPagesRepeated() throws Exception { } private ColumnWriteStoreV1 newColumnWriteStoreImpl(MemPageStore memPageStore) { - return new ColumnWriteStoreV1(memPageStore, 2048, 2048, false, WriterVersion.PARQUET_1_0, new HeapByteBufferAllocator()); + return new ColumnWriteStoreV1(memPageStore, + ParquetProperties.builder() + .withPageSize(2048) + .withDictionaryEncoding(false) + .build()); } } diff --git a/parquet-column/src/test/java/org/apache/parquet/io/PerfTest.java b/parquet-column/src/test/java/org/apache/parquet/io/PerfTest.java index aff3937b9e..e4687b1c3c 100644 --- a/parquet-column/src/test/java/org/apache/parquet/io/PerfTest.java +++ b/parquet-column/src/test/java/org/apache/parquet/io/PerfTest.java @@ -27,8 +27,7 @@ import java.util.logging.Level; import org.apache.parquet.Log; -import org.apache.parquet.bytes.HeapByteBufferAllocator; -import org.apache.parquet.column.ParquetProperties.WriterVersion; +import org.apache.parquet.column.ParquetProperties; import org.apache.parquet.column.impl.ColumnWriteStoreV1; import org.apache.parquet.column.page.mem.MemPageStore; import org.apache.parquet.example.DummyRecordConverter; @@ -78,7 +77,12 @@ private static void read(MemPageStore memPageStore, MessageType myschema, private static void write(MemPageStore memPageStore) { - ColumnWriteStoreV1 columns = new ColumnWriteStoreV1(memPageStore, 50*1024*1024, 50*1024*1024, false, WriterVersion.PARQUET_1_0, new HeapByteBufferAllocator()); + ColumnWriteStoreV1 columns = new ColumnWriteStoreV1( + memPageStore, + ParquetProperties.builder() + .withPageSize(50*1024*1024) + .withDictionaryEncoding(false) + .build()); MessageColumnIO columnIO = newColumnFactory(schema); GroupWriter groupWriter = new GroupWriter(columnIO.getRecordWriter(columns), schema); diff --git a/parquet-column/src/test/java/org/apache/parquet/io/TestColumnIO.java b/parquet-column/src/test/java/org/apache/parquet/io/TestColumnIO.java index 06f22b6f35..e9e599affe 100644 --- a/parquet-column/src/test/java/org/apache/parquet/io/TestColumnIO.java +++ b/parquet-column/src/test/java/org/apache/parquet/io/TestColumnIO.java @@ -38,7 +38,7 @@ import java.util.Iterator; import java.util.List; -import org.apache.parquet.bytes.HeapByteBufferAllocator; +import org.apache.parquet.column.ParquetProperties; import org.junit.Assert; import org.junit.Test; import org.junit.runner.RunWith; @@ -48,7 +48,6 @@ import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.column.ColumnWriteStore; import org.apache.parquet.column.ColumnWriter; -import org.apache.parquet.column.ParquetProperties.WriterVersion; import org.apache.parquet.column.impl.ColumnWriteStoreV1; import org.apache.parquet.column.page.PageReadStore; import org.apache.parquet.column.page.mem.MemPageStore; @@ -527,7 +526,12 @@ public void testPushParser() { } private ColumnWriteStoreV1 newColumnWriteStore(MemPageStore memPageStore) { - return new ColumnWriteStoreV1(memPageStore, 800, 800, useDictionary, WriterVersion.PARQUET_1_0, new HeapByteBufferAllocator()); + return new ColumnWriteStoreV1(memPageStore, + ParquetProperties.builder() + .withPageSize(800) + .withDictionaryPageSize(800) + .withDictionaryEncoding(useDictionary) + .build()); } @Test diff --git a/parquet-column/src/test/java/org/apache/parquet/io/TestFiltered.java b/parquet-column/src/test/java/org/apache/parquet/io/TestFiltered.java index 25b629b47e..ab5a5750b0 100644 --- a/parquet-column/src/test/java/org/apache/parquet/io/TestFiltered.java +++ b/parquet-column/src/test/java/org/apache/parquet/io/TestFiltered.java @@ -21,11 +21,10 @@ import java.util.ArrayList; import java.util.List; -import org.apache.parquet.bytes.HeapByteBufferAllocator; +import org.apache.parquet.column.ParquetProperties; import org.apache.parquet.io.api.RecordConsumer; import org.junit.Test; -import org.apache.parquet.column.ParquetProperties.WriterVersion; import org.apache.parquet.column.impl.ColumnWriteStoreV1; import org.apache.parquet.column.page.mem.MemPageStore; import org.apache.parquet.example.data.Group; @@ -259,7 +258,12 @@ public void testFilteredNotPaged() { private MemPageStore writeTestRecords(MessageColumnIO columnIO, int number) { MemPageStore memPageStore = new MemPageStore(number * 2); - ColumnWriteStoreV1 columns = new ColumnWriteStoreV1(memPageStore, 800, 800, false, WriterVersion.PARQUET_1_0, new HeapByteBufferAllocator()); + ColumnWriteStoreV1 columns = new ColumnWriteStoreV1( + memPageStore, + ParquetProperties.builder() + .withPageSize(800) + .withDictionaryEncoding(false) + .build()); RecordConsumer recordWriter = columnIO.getRecordWriter(columns); GroupWriter groupWriter = new GroupWriter(recordWriter, schema); diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java index 0730f2a578..c107d67fb7 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/InternalParquetRecordWriter.java @@ -28,11 +28,9 @@ import java.util.HashMap; import java.util.Map; -import org.apache.parquet.bytes.ByteBufferAllocator; import org.apache.parquet.Log; import org.apache.parquet.column.ColumnWriteStore; import org.apache.parquet.column.ParquetProperties; -import org.apache.parquet.column.ParquetProperties.WriterVersion; import org.apache.parquet.hadoop.CodecFactory.BytesCompressor; import org.apache.parquet.hadoop.api.WriteSupport; import org.apache.parquet.hadoop.api.WriteSupport.FinalizedWriteContext; @@ -54,10 +52,9 @@ class InternalParquetRecordWriter { private final long rowGroupSize; private long rowGroupSizeThreshold; private long nextRowGroupSize; - private final int pageSize; private final BytesCompressor compressor; private final boolean validating; - private final ParquetProperties parquetProperties; + private final ParquetProperties props; private long recordCount = 0; private long recordCountForNextMemCheck = MINIMUM_RECORD_COUNT_FOR_CHECK; @@ -67,45 +64,6 @@ class InternalParquetRecordWriter { private ColumnChunkPageWriteStore pageStore; private RecordConsumer recordConsumer; - - /** - * @param parquetFileWriter the file to write to - * @param writeSupport the class to convert incoming records - * @param schema the schema of the records - * @param extraMetaData extra meta data to write in the footer of the file - * @param rowGroupSize the size of a block in the file (this will be approximate) - * @param compressor the codec used to compress - */ - public InternalParquetRecordWriter( - ParquetFileWriter parquetFileWriter, - WriteSupport writeSupport, - MessageType schema, - Map extraMetaData, - long rowGroupSize, - int pageSize, - BytesCompressor compressor, - int dictionaryPageSize, - boolean enableDictionary, - boolean validating, - WriterVersion writerVersion, - ByteBufferAllocator allocator) { - this(parquetFileWriter, - writeSupport, - schema, - extraMetaData, - rowGroupSize, - pageSize, - compressor, - dictionaryPageSize, - enableDictionary, - ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, - ParquetProperties.DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, - false, - validating, - writerVersion, - allocator); - } - /** * @param parquetFileWriter the file to write to * @param writeSupport the class to convert incoming records @@ -120,16 +78,9 @@ public InternalParquetRecordWriter( MessageType schema, Map extraMetaData, long rowGroupSize, - int pageSize, BytesCompressor compressor, - int dictionaryPageSize, - boolean enableDictionary, - int minRowCountForSizeCheck, - int maxRowCountForSizeCheck, - boolean constantNextSizeCheck, boolean validating, - WriterVersion writerVersion, - ByteBufferAllocator allocator) { + ParquetProperties props) { this.parquetFileWriter = parquetFileWriter; this.writeSupport = checkNotNull(writeSupport, "writeSupport"); this.schema = schema; @@ -137,21 +88,15 @@ public InternalParquetRecordWriter( this.rowGroupSize = rowGroupSize; this.rowGroupSizeThreshold = rowGroupSize; this.nextRowGroupSize = rowGroupSizeThreshold; - this.pageSize = pageSize; this.compressor = compressor; this.validating = validating; - this.parquetProperties = new ParquetProperties(dictionaryPageSize, writerVersion, enableDictionary, - minRowCountForSizeCheck, maxRowCountForSizeCheck, constantNextSizeCheck, allocator); + this.props = props; initStore(); } private void initStore() { - pageStore = new ColumnChunkPageWriteStore(compressor, schema, parquetProperties.getAllocator()); - columnStore = parquetProperties.newColumnWriteStore( - schema, - pageStore, - pageSize, - parquetProperties.getAllocator()); + pageStore = new ColumnChunkPageWriteStore(compressor, schema, props.getAllocator()); + columnStore = props.newColumnWriteStore(schema, pageStore); MessageColumnIO columnIO = new ColumnIOFactory(validating).getColumnIO(schema); this.recordConsumer = columnIO.getRecordWriter(columnStore); writeSupport.prepareForWrite(recordConsumer); diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java index 6e28266753..8979dba3f6 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetOutputFormat.java @@ -21,7 +21,6 @@ import static org.apache.parquet.Log.INFO; import static org.apache.parquet.Preconditions.checkNotNull; import static org.apache.parquet.hadoop.ParquetWriter.DEFAULT_BLOCK_SIZE; -import static org.apache.parquet.hadoop.ParquetWriter.DEFAULT_PAGE_SIZE; import static org.apache.parquet.hadoop.util.ContextUtil.getConfiguration; import java.io.IOException; @@ -242,19 +241,23 @@ public static boolean getValidation(JobContext jobContext) { } public static boolean getEnableDictionary(Configuration configuration) { - return configuration.getBoolean(ENABLE_DICTIONARY, true); + return configuration.getBoolean( + ENABLE_DICTIONARY, ParquetProperties.DEFAULT_IS_DICTIONARY_ENABLED); } public static int getMinRowCountForPageSizeCheck(Configuration configuration) { - return configuration.getInt(MIN_ROW_COUNT_FOR_PAGE_SIZE_CHECK, ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK); + return configuration.getInt(MIN_ROW_COUNT_FOR_PAGE_SIZE_CHECK, + ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK); } public static int getMaxRowCountForPageSizeCheck(Configuration configuration) { - return configuration.getInt(MAX_ROW_COUNT_FOR_PAGE_SIZE_CHECK, ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK); + return configuration.getInt(MAX_ROW_COUNT_FOR_PAGE_SIZE_CHECK, + ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK); } public static boolean getEstimatePageSizeCheck(Configuration configuration) { - return configuration.getBoolean(ESTIMATE_PAGE_SIZE_CHECK, ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK); + return configuration.getBoolean(ESTIMATE_PAGE_SIZE_CHECK, + ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK); } @Deprecated @@ -267,15 +270,17 @@ public static long getLongBlockSize(Configuration configuration) { } public static int getPageSize(Configuration configuration) { - return configuration.getInt(PAGE_SIZE, DEFAULT_PAGE_SIZE); + return configuration.getInt(PAGE_SIZE, ParquetProperties.DEFAULT_PAGE_SIZE); } public static int getDictionaryPageSize(Configuration configuration) { - return configuration.getInt(DICTIONARY_PAGE_SIZE, DEFAULT_PAGE_SIZE); + return configuration.getInt( + DICTIONARY_PAGE_SIZE, ParquetProperties.DEFAULT_DICTIONARY_PAGE_SIZE); } public static WriterVersion getWriterVersion(Configuration configuration) { - String writerVersion = configuration.get(WRITER_VERSION, WriterVersion.PARQUET_1_0.toString()); + String writerVersion = configuration.get( + WRITER_VERSION, ParquetProperties.DEFAULT_WRITER_VERSION.toString()); return WriterVersion.fromString(writerVersion); } @@ -357,28 +362,32 @@ public RecordWriter getRecordWriter(Configuration conf, Path file, Comp throws IOException, InterruptedException { final WriteSupport writeSupport = getWriteSupport(conf); + ParquetProperties props = ParquetProperties.builder() + .withPageSize(getPageSize(conf)) + .withDictionaryPageSize(getDictionaryPageSize(conf)) + .withDictionaryEncoding(getEnableDictionary(conf)) + .withWriterVersion(getWriterVersion(conf)) + .estimateRowCountForPageSizeCheck(getEstimatePageSizeCheck(conf)) + .withMinRowCountForPageSizeCheck(getMinRowCountForPageSizeCheck(conf)) + .withMaxRowCountForPageSizeCheck(getMaxRowCountForPageSizeCheck(conf)) + .build(); + long blockSize = getLongBlockSize(conf); - if (INFO) LOG.info("Parquet block size to " + blockSize); - int pageSize = getPageSize(conf); - if (INFO) LOG.info("Parquet page size to " + pageSize); - int dictionaryPageSize = getDictionaryPageSize(conf); - if (INFO) LOG.info("Parquet dictionary page size to " + dictionaryPageSize); - boolean enableDictionary = getEnableDictionary(conf); - if (INFO) LOG.info("Dictionary is " + (enableDictionary ? "on" : "off")); + int maxPaddingSize = getMaxPaddingSize(conf); boolean validating = getValidation(conf); + + if (INFO) LOG.info("Parquet block size to " + blockSize); + if (INFO) LOG.info("Parquet page size to " + props.getPageSizeThreshold()); + if (INFO) LOG.info("Parquet dictionary page size to " + props.getDictionaryPageSizeThreshold()); + if (INFO) LOG.info("Dictionary is " + (props.isEnableDictionary() ? "on" : "off")); if (INFO) LOG.info("Validation is " + (validating ? "on" : "off")); - WriterVersion writerVersion = getWriterVersion(conf); - if (INFO) LOG.info("Writer version is: " + writerVersion); - int maxPaddingSize = getMaxPaddingSize(conf); + if (INFO) LOG.info("Writer version is: " + props.getWriterVersion()); if (INFO) LOG.info("Maximum row group padding size is " + maxPaddingSize + " bytes"); - boolean estimateNextSizeCheck = getEstimatePageSizeCheck(conf); - if (INFO) LOG.info("Page size checking is: " + (estimateNextSizeCheck ? "estimated" : "constant")); - int minRowCountForPageSizeCheck = getMinRowCountForPageSizeCheck(conf); - if (INFO) LOG.info("Min row count for page size check is: " + minRowCountForPageSizeCheck); - int maxRowCountForPageSizeCheck = getMaxRowCountForPageSizeCheck(conf); - if (INFO) LOG.info("Min row count for page size check is: " + maxRowCountForPageSizeCheck); + if (INFO) LOG.info("Page size checking is: " + (props.estimateNextSizeCheck() ? "estimated" : "constant")); + if (INFO) LOG.info("Min row count for page size check is: " + props.getMinRowCountForPageSizeCheck()); + if (INFO) LOG.info("Min row count for page size check is: " + props.getMaxRowCountForPageSizeCheck()); - CodecFactory codecFactory = new CodecFactory(conf, pageSize); + CodecFactory codecFactory = new CodecFactory(conf, props.getPageSizeThreshold()); WriteContext init = writeSupport.init(conf); ParquetFileWriter w = new ParquetFileWriter( @@ -401,15 +410,10 @@ public RecordWriter getRecordWriter(Configuration conf, Path file, Comp writeSupport, init.getSchema(), init.getExtraMetaData(), - blockSize, pageSize, + blockSize, codecFactory.getCompressor(codec), - dictionaryPageSize, - enableDictionary, - minRowCountForPageSizeCheck, - maxRowCountForPageSizeCheck, - estimateNextSizeCheck, validating, - writerVersion, + props, memoryManager); } diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java index 3b08e53b4a..6c94fac5c6 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetRecordWriter.java @@ -24,16 +24,13 @@ import org.apache.hadoop.mapreduce.RecordWriter; import org.apache.hadoop.mapreduce.TaskAttemptContext; -import org.apache.parquet.bytes.HeapByteBufferAllocator; +import org.apache.parquet.column.ParquetProperties; import org.apache.parquet.column.ParquetProperties.WriterVersion; import org.apache.parquet.hadoop.CodecFactory.BytesCompressor; import org.apache.parquet.hadoop.api.WriteSupport; import org.apache.parquet.schema.MessageType; import static org.apache.parquet.Preconditions.checkNotNull; -import static org.apache.parquet.column.ParquetProperties.DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK; -import static org.apache.parquet.column.ParquetProperties.DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK; -import static org.apache.parquet.column.ParquetProperties.DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK; /** * Writes records to a Parquet file @@ -73,10 +70,15 @@ public ParquetRecordWriter( boolean enableDictionary, boolean validating, WriterVersion writerVersion) { + ParquetProperties props = ParquetProperties.builder() + .withPageSize(pageSize) + .withDictionaryPageSize(dictionaryPageSize) + .withDictionaryEncoding(enableDictionary) + .withWriterVersion(writerVersion) + .build(); internalWriter = new InternalParquetRecordWriter(w, writeSupport, schema, - extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, - DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, - DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, validating, writerVersion, new HeapByteBufferAllocator()); + extraMetaData, blockSize, compressor, validating, props); + this.memoryManager = null; } /** @@ -91,24 +93,28 @@ public ParquetRecordWriter( * @param enableDictionary to enable the dictionary * @param validating if schema validation should be turned on */ + @Deprecated public ParquetRecordWriter( - ParquetFileWriter w, - WriteSupport writeSupport, - MessageType schema, - Map extraMetaData, - long blockSize, int pageSize, - BytesCompressor compressor, - int dictionaryPageSize, - boolean enableDictionary, - boolean validating, - WriterVersion writerVersion, - MemoryManager memoryManager) { - this(w, writeSupport, schema, - extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, - DEFAULT_MINIMUM_RECORD_COUNT_FOR_CHECK, DEFAULT_MAXIMUM_RECORD_COUNT_FOR_CHECK, - DEFAULT_ESTIMATE_ROW_COUNT_FOR_PAGE_SIZE_CHECK, - validating, writerVersion, memoryManager); - } + ParquetFileWriter w, + WriteSupport writeSupport, + MessageType schema, + Map extraMetaData, + long blockSize, int pageSize, + BytesCompressor compressor, + int dictionaryPageSize, + boolean enableDictionary, + boolean validating, + WriterVersion writerVersion, + MemoryManager memoryManager) { + this(w, writeSupport, schema, extraMetaData, blockSize, compressor, + validating, ParquetProperties.builder() + .withPageSize(pageSize) + .withDictionaryPageSize(dictionaryPageSize) + .withDictionaryEncoding(enableDictionary) + .withWriterVersion(writerVersion) + .build(), + memoryManager); + } /** * @@ -118,32 +124,21 @@ public ParquetRecordWriter( * @param extraMetaData extra meta data to write in the footer of the file * @param blockSize the size of a block in the file (this will be approximate) * @param compressor the compressor used to compress the pages - * @param dictionaryPageSize the threshold for dictionary size - * @param enableDictionary to enable the dictionary - * @param minRowCountForSizeCheck min number of rows before page size check - * @param maxRowCountForSizeCheck max number of rows before page size check - * @param estimateNextPageSizeCheck dynamically estimate next page size check (min will be used if false) * @param validating if schema validation should be turned on + * @param props parquet encoding properties */ - public ParquetRecordWriter( + ParquetRecordWriter( ParquetFileWriter w, WriteSupport writeSupport, MessageType schema, Map extraMetaData, - long blockSize, int pageSize, + long blockSize, BytesCompressor compressor, - int dictionaryPageSize, - boolean enableDictionary, - int minRowCountForSizeCheck, - int maxRowCountForSizeCheck, - boolean estimateNextPageSizeCheck, boolean validating, - WriterVersion writerVersion, + ParquetProperties props, MemoryManager memoryManager) { internalWriter = new InternalParquetRecordWriter(w, writeSupport, schema, - extraMetaData, blockSize, pageSize, compressor, dictionaryPageSize, enableDictionary, minRowCountForSizeCheck, - maxRowCountForSizeCheck, estimateNextPageSizeCheck, - validating, writerVersion, new HeapByteBufferAllocator()); + extraMetaData, blockSize, compressor, validating, props); this.memoryManager = checkNotNull(memoryManager, "memoryManager"); memoryManager.addWriter(internalWriter, blockSize); } diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetWriter.java index e2521fb090..f214d8d141 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetWriter.java @@ -29,7 +29,6 @@ import org.apache.parquet.hadoop.api.WriteSupport; import org.apache.parquet.hadoop.metadata.CompressionCodecName; import org.apache.parquet.schema.MessageType; -import org.apache.parquet.bytes.HeapByteBufferAllocator; /** * Write records to a Parquet file. @@ -37,13 +36,15 @@ public class ParquetWriter implements Closeable { public static final int DEFAULT_BLOCK_SIZE = 128 * 1024 * 1024; - public static final int DEFAULT_PAGE_SIZE = 1 * 1024 * 1024; + public static final int DEFAULT_PAGE_SIZE = + ParquetProperties.DEFAULT_PAGE_SIZE; public static final CompressionCodecName DEFAULT_COMPRESSION_CODEC_NAME = CompressionCodecName.UNCOMPRESSED; - public static final boolean DEFAULT_IS_DICTIONARY_ENABLED = true; + public static final boolean DEFAULT_IS_DICTIONARY_ENABLED = + ParquetProperties.DEFAULT_IS_DICTIONARY_ENABLED; public static final boolean DEFAULT_IS_VALIDATING_ENABLED = false; public static final WriterVersion DEFAULT_WRITER_VERSION = - WriterVersion.PARQUET_1_0; + ParquetProperties.DEFAULT_WRITER_VERSION; // max size (bytes) to write as padding and the min size of a row group public static final int MAX_PADDING_SIZE_DEFAULT = 0; @@ -215,9 +216,14 @@ public ParquetWriter( boolean validating, WriterVersion writerVersion, Configuration conf) throws IOException { - this(file, mode, writeSupport, compressionCodecName, blockSize, pageSize, - dictionaryPageSize, enableDictionary, validating, writerVersion, conf, - MAX_PADDING_SIZE_DEFAULT); + this(file, mode, writeSupport, compressionCodecName, blockSize, + validating, conf, MAX_PADDING_SIZE_DEFAULT, + ParquetProperties.builder() + .withPageSize(pageSize) + .withDictionaryPageSize(dictionaryPageSize) + .withDictionaryEncoding(enableDictionary) + .withWriterVersion(writerVersion) + .build()); } /** @@ -253,13 +259,10 @@ public ParquetWriter(Path file, Configuration conf, WriteSupport writeSupport WriteSupport writeSupport, CompressionCodecName compressionCodecName, int blockSize, - int pageSize, - int dictionaryPageSize, - boolean enableDictionary, boolean validating, - WriterVersion writerVersion, Configuration conf, - int maxPaddingSize) throws IOException { + int maxPaddingSize, + ParquetProperties encodingProps) throws IOException { WriteSupport.WriteContext writeContext = writeSupport.init(conf); MessageType schema = writeContext.getSchema(); @@ -268,7 +271,7 @@ public ParquetWriter(Path file, Configuration conf, WriteSupport writeSupport conf, schema, file, mode, blockSize, maxPaddingSize); fileWriter.start(); - CodecFactory codecFactory = new CodecFactory(conf, pageSize); + CodecFactory codecFactory = new CodecFactory(conf, encodingProps.getPageSizeThreshold()); CodecFactory.BytesCompressor compressor = codecFactory.getCompressor(compressionCodecName); this.writer = new InternalParquetRecordWriter( fileWriter, @@ -276,13 +279,9 @@ public ParquetWriter(Path file, Configuration conf, WriteSupport writeSupport schema, writeContext.getExtraMetaData(), blockSize, - pageSize, compressor, - dictionaryPageSize, - enableDictionary, validating, - writerVersion, - new HeapByteBufferAllocator()); + encodingProps); } public void write(T object) throws IOException { @@ -324,12 +323,10 @@ public abstract static class Builder> { private ParquetFileWriter.Mode mode; private CompressionCodecName codecName = DEFAULT_COMPRESSION_CODEC_NAME; private int rowGroupSize = DEFAULT_BLOCK_SIZE; - private int pageSize = DEFAULT_PAGE_SIZE; - private int dictionaryPageSize = DEFAULT_PAGE_SIZE; private int maxPaddingSize = MAX_PADDING_SIZE_DEFAULT; - private boolean enableDictionary = DEFAULT_IS_DICTIONARY_ENABLED; private boolean enableValidation = DEFAULT_IS_VALIDATING_ENABLED; - private WriterVersion writerVersion = DEFAULT_WRITER_VERSION; + private ParquetProperties.Builder encodingPropsBuilder = + ParquetProperties.builder(); protected Builder(Path file) { this.file = file; @@ -398,7 +395,7 @@ public SELF withRowGroupSize(int rowGroupSize) { * @return this builder for method chaining. */ public SELF withPageSize(int pageSize) { - this.pageSize = pageSize; + encodingPropsBuilder.withPageSize(pageSize); return self(); } @@ -410,7 +407,7 @@ public SELF withPageSize(int pageSize) { * @return this builder for method chaining. */ public SELF withDictionaryPageSize(int dictionaryPageSize) { - this.dictionaryPageSize = dictionaryPageSize; + encodingPropsBuilder.withDictionaryPageSize(dictionaryPageSize); return self(); } @@ -433,7 +430,7 @@ public SELF withMaxPaddingSize(int maxPaddingSize) { * @return this builder for method chaining. */ public SELF enableDictionaryEncoding() { - this.enableDictionary = true; + encodingPropsBuilder.withDictionaryEncoding(true); return self(); } @@ -444,7 +441,7 @@ public SELF enableDictionaryEncoding() { * @return this builder for method chaining. */ public SELF withDictionaryEncoding(boolean enableDictionary) { - this.enableDictionary = enableDictionary; + encodingPropsBuilder.withDictionaryEncoding(enableDictionary); return self(); } @@ -477,7 +474,7 @@ public SELF withValidation(boolean enableValidation) { * @return this builder for method chaining. */ public SELF withWriterVersion(WriterVersion version) { - this.writerVersion = version; + encodingPropsBuilder.withWriterVersion(version); return self(); } @@ -489,8 +486,8 @@ public SELF withWriterVersion(WriterVersion version) { */ public ParquetWriter build() throws IOException { return new ParquetWriter(file, mode, getWriteSupport(conf), codecName, - rowGroupSize, pageSize, dictionaryPageSize, enableDictionary, - enableValidation, writerVersion, conf, maxPaddingSize); + rowGroupSize, enableValidation, conf, maxPaddingSize, + encodingPropsBuilder.build()); } } } diff --git a/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/TestUtils.java b/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/TestUtils.java index 7c5a186817..e53ac785a0 100644 --- a/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/TestUtils.java +++ b/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/TestUtils.java @@ -53,7 +53,12 @@ public static void assertThrows( Assert.fail("No exception was thrown (" + message + "), expected: " + expected.getName()); } catch (Exception actual) { - Assert.assertEquals(message, expected, actual.getClass()); + try { + Assert.assertEquals(message, expected, actual.getClass()); + } catch (AssertionError e) { + e.addSuppressed(actual); + throw e; + } } } } diff --git a/parquet-pig/src/test/java/org/apache/parquet/pig/TupleConsumerPerfTest.java b/parquet-pig/src/test/java/org/apache/parquet/pig/TupleConsumerPerfTest.java index c050922847..ff192e2c90 100644 --- a/parquet-pig/src/test/java/org/apache/parquet/pig/TupleConsumerPerfTest.java +++ b/parquet-pig/src/test/java/org/apache/parquet/pig/TupleConsumerPerfTest.java @@ -22,7 +22,7 @@ import java.util.Map; import java.util.logging.Level; -import org.apache.parquet.bytes.HeapByteBufferAllocator; +import org.apache.parquet.column.ParquetProperties; import org.apache.pig.backend.executionengine.ExecException; import org.apache.pig.data.DataBag; import org.apache.pig.data.NonSpillableDataBag; @@ -32,7 +32,6 @@ import org.apache.pig.parser.ParserException; import org.apache.parquet.Log; -import org.apache.parquet.column.ParquetProperties.WriterVersion; import org.apache.parquet.column.impl.ColumnWriteStoreV1; import org.apache.parquet.column.page.PageReadStore; import org.apache.parquet.column.page.mem.MemPageStore; @@ -60,7 +59,11 @@ public static void main(String[] args) throws Exception { MessageType schema = new PigSchemaConverter().convert(Utils.getSchemaFromString(pigSchema)); MemPageStore memPageStore = new MemPageStore(0); - ColumnWriteStoreV1 columns = new ColumnWriteStoreV1(memPageStore, 50*1024*1024, 50*1024*1024, false, WriterVersion.PARQUET_1_0, new HeapByteBufferAllocator()); + ColumnWriteStoreV1 columns = new ColumnWriteStoreV1( + memPageStore, ParquetProperties.builder() + .withPageSize(50*1024*1024) + .withDictionaryEncoding(false) + .build()); write(memPageStore, columns, schema, pigSchema); columns.flush(); read(memPageStore, pigSchema, pigSchemaProjected, pigSchemaNoString); diff --git a/parquet-thrift/src/test/java/org/apache/parquet/thrift/TestParquetReadProtocol.java b/parquet-thrift/src/test/java/org/apache/parquet/thrift/TestParquetReadProtocol.java index f954e4c258..97e0054b8a 100644 --- a/parquet-thrift/src/test/java/org/apache/parquet/thrift/TestParquetReadProtocol.java +++ b/parquet-thrift/src/test/java/org/apache/parquet/thrift/TestParquetReadProtocol.java @@ -31,7 +31,7 @@ import java.util.Map; import java.util.Set; -import org.apache.parquet.bytes.HeapByteBufferAllocator; +import org.apache.parquet.column.ParquetProperties; import thrift.test.OneOfEach; import org.apache.thrift.TBase; @@ -39,7 +39,6 @@ import org.junit.Test; import org.apache.parquet.Log; -import org.apache.parquet.column.ParquetProperties.WriterVersion; import org.apache.parquet.column.impl.ColumnWriteStoreV1; import org.apache.parquet.column.page.mem.MemPageStore; import org.apache.parquet.io.ColumnIOFactory; @@ -149,8 +148,11 @@ private > void validate(T expected) throws TException { final MessageType schema = schemaConverter.convert(thriftClass); LOG.info(schema); final MessageColumnIO columnIO = new ColumnIOFactory(true).getColumnIO(schema); - final ColumnWriteStoreV1 columns = new ColumnWriteStoreV1(memPageStore, 10000, 10000, false, - WriterVersion.PARQUET_1_0, new HeapByteBufferAllocator()); + final ColumnWriteStoreV1 columns = new ColumnWriteStoreV1(memPageStore, + ParquetProperties.builder() + .withPageSize(10000) + .withDictionaryEncoding(false) + .build()); final RecordConsumer recordWriter = columnIO.getRecordWriter(columns); final StructType thriftType = schemaConverter.toStructType(thriftClass); ParquetWriteProtocol parquetWriteProtocol = new ParquetWriteProtocol(recordWriter, columnIO, thriftType);