From 382ccb0f97cdaa3fabc53236e90e6e229eb31adf Mon Sep 17 00:00:00 2001 From: "nitin.goyal" Date: Mon, 19 Oct 2015 10:11:49 +0530 Subject: [PATCH 1/3] PARQUET-353: Recycle compressors in parquet write path --- .../org/apache/parquet/hadoop/ParquetFileWriter.java | 12 ++++++++++++ .../apache/parquet/hadoop/ParquetOutputFormat.java | 3 +-- .../org/apache/parquet/hadoop/ParquetWriter.java | 3 +-- 3 files changed, 14 insertions(+), 4 deletions(-) diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java index 664ee9d87d..f7187f1db2 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java @@ -171,6 +171,8 @@ private final STATE error() throws IOException { private STATE state = STATE.NOT_STARTED; + private final CodecFactory codecFactory; + /** * @param configuration Hadoop configuration * @param schema the schema of the data @@ -226,6 +228,8 @@ public ParquetFileWriter(Configuration configuration, MessageType schema, this.alignment = NoAlignment.get(rowGroupSize); this.out = fs.create(file, overwriteFlag); } + + this.codecFactory = new CodecFactory(configuration); } /** @@ -469,6 +473,7 @@ public void end(Map extraMetaData) throws IOException { ParquetMetadata footer = new ParquetMetadata(new FileMetaData(schema, extraMetaData, Version.FULL_VERSION), blocks); serializeFooter(footer, out); out.close(); + codecFactory.release(); } private static void serializeFooter(ParquetMetadata footer, FSDataOutputStream out) throws IOException { @@ -588,6 +593,13 @@ public long getPos() throws IOException { return out.getPos(); } + /** + * @return The codec factory which should be used to get compressors + */ + public CodecFactory getCodecFactory() { + return codecFactory; + } + public long getNextRowGroupSize() throws IOException { return alignment.nextRowGroupSize(out); } 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 ad6c0344c1..995af032ba 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 @@ -341,7 +341,6 @@ public RecordWriter getRecordWriter(Configuration conf, Path file, Comp throws IOException, InterruptedException { final WriteSupport writeSupport = getWriteSupport(conf); - CodecFactory codecFactory = new CodecFactory(conf); long blockSize = getLongBlockSize(conf); if (INFO) LOG.info("Parquet block size to " + blockSize); int pageSize = getPageSize(conf); @@ -379,7 +378,7 @@ public RecordWriter getRecordWriter(Configuration conf, Path file, Comp init.getSchema(), init.getExtraMetaData(), blockSize, pageSize, - codecFactory.getCompressor(codec, pageSize), + w.getCodecFactory().getCompressor(codec, pageSize), dictionaryPageSize, enableDictionary, validating, 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 e3b7953760..4e1c9403c4 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 @@ -267,8 +267,7 @@ public ParquetWriter(Path file, Configuration conf, WriteSupport writeSupport conf, schema, file, mode, blockSize, maxPaddingSize); fileWriter.start(); - CodecFactory codecFactory = new CodecFactory(conf); - CodecFactory.BytesCompressor compressor = codecFactory.getCompressor(compressionCodecName, 0); + CodecFactory.BytesCompressor compressor = fileWriter.getCodecFactory().getCompressor(compressionCodecName, 0); this.writer = new InternalParquetRecordWriter( fileWriter, writeSupport, From 801f18eb61d4139ec1a33d51b93f8b88f30c5081 Mon Sep 17 00:00:00 2001 From: "nitin.goyal" Date: Mon, 19 Oct 2015 12:07:00 +0530 Subject: [PATCH 2/3] PARQUET-353: Recycle compressors in parquet write path --- .../main/java/org/apache/parquet/hadoop/ParquetFileWriter.java | 1 + 1 file changed, 1 insertion(+) diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java index f7187f1db2..6f9aac73c5 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileWriter.java @@ -250,6 +250,7 @@ public ParquetFileWriter(Configuration configuration, MessageType schema, rowAndBlockSize, rowAndBlockSize, maxPaddingSize); this.out = fs.create(file, true, DFS_BUFFER_SIZE_DEFAULT, fs.getDefaultReplication(file), rowAndBlockSize); + this.codecFactory = new CodecFactory(configuration); } /** From 29cc207940399451ad633d721d518f1433c78f5e Mon Sep 17 00:00:00 2001 From: "nitin.goyal" Date: Mon, 26 Oct 2015 10:32:52 +0530 Subject: [PATCH 3/3] Dummy checkin