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..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 @@ -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); } /** @@ -246,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); } /** @@ -469,6 +474,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 +594,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,