diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/CodecFactory.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/CodecFactory.java index c9391201f4..1ecfc66cae 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/CodecFactory.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/CodecFactory.java @@ -385,7 +385,12 @@ private CompressionCodec getCodecAtLevel(CompressionCodecName codecName, int lev break; case GZIP: validateGzipLevel(level); - levelConf.setEnum(GZIP_COMPRESS_LEVEL, zlibCompressionLevel(level)); + // Store the enum constant name rather than using Configuration#setEnum, which persists + // Enum#toString(). Hadoop reads this back via Configuration#getEnum -> Enum#valueOf, which + // requires the constant name. In some Hadoop builds ZlibCompressor.CompressionLevel + // overrides toString() to return the numeric level (e.g. "5"), so setEnum would write a + // value that valueOf cannot resolve, throwing "No enum constant ...CompressionLevel.5". + levelConf.set(GZIP_COMPRESS_LEVEL, zlibCompressionLevel(level).name()); break; case BROTLI: validateBrotliLevel(level); @@ -464,7 +469,12 @@ private String cacheKey(CompressionCodecName codecName) { private String cacheKey(CompressionCodecName codecName, int level) { String codecClass = codecName.getHadoopCompressionCodecClassName(); - return (codecClass == null ? codecName.name() : codecClass) + ":" + level; + // Use a distinct namespace ("#level=") for the leveled path so that this key can never + // collide with the no-level cacheKey(codecName), which appends the raw configuration value + // (e.g. ":5" when zlib.compress.level=5 is set directly). Both maps that use + // these keys (the per-factory compressors map and the shared static CODEC_BY_NAME) would + // otherwise return a codec built from the raw config instead of the level-configured one. + return (codecClass == null ? codecName.name() : codecClass) + "#level=" + level; } @Override diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileReader.java b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileReader.java index 9af4b4ac60..b970b37602 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileReader.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/hadoop/ParquetFileReader.java @@ -1190,10 +1190,18 @@ private ColumnChunkPageReadStore internalReadRowGroup(int blockIndex) throws IOE } // actually read all the chunks ChunkListBuilder builder = new ChunkListBuilder(block.getRowCount()); - readAllPartsVectoredOrNormal(allParts, builder); - rowGroup.setReleaser(builder.releaser); - for (Chunk chunk : builder.build()) { - readChunkPages(chunk, block, rowGroup); + try { + readAllPartsVectoredOrNormal(allParts, builder); + rowGroup.setReleaser(builder.releaser); + for (Chunk chunk : builder.build()) { + readChunkPages(chunk, block, rowGroup); + } + } catch (RuntimeException | IOException e) { + // If we fail before the releaser is transferred to the row group (e.g. a vectored range + // times out after earlier ranges already registered their buffers), release any buffers + // that were registered so far so that partially-read row groups do not leak. + builder.releaser.close(); + throw e; } return rowGroup; @@ -1464,10 +1472,18 @@ private ColumnChunkPageReadStore internalReadFilteredRowGroup( } } } - readAllPartsVectoredOrNormal(allParts, builder); - rowGroup.setReleaser(builder.releaser); - for (Chunk chunk : builder.build()) { - readChunkPages(chunk, block, rowGroup); + try { + readAllPartsVectoredOrNormal(allParts, builder); + rowGroup.setReleaser(builder.releaser); + for (Chunk chunk : builder.build()) { + readChunkPages(chunk, block, rowGroup); + } + } catch (RuntimeException | IOException e) { + // If we fail before the releaser is transferred to the row group (e.g. a vectored range + // times out after earlier ranges already registered their buffers), release any buffers + // that were registered so far so that partially-read row groups do not leak. + builder.releaser.close(); + throw e; } return rowGroup; @@ -2368,6 +2384,10 @@ public void readFromVectoredRange(ParquetFileRange currRange, ChunkListBuilder b LOG.error(error, e); throw new IOException(error, e); } + // Release the vectored-read buffer back to the allocator when the row group is closed. + // Requires fs.file.checksum.verify=false so the returned buffer is the allocator buffer + // rather than a sliced subset (see Hadoop's fs.file.checksum.verify docs). + builder.addBuffersToRelease(Collections.singletonList(buffer)); ByteBufferInputStream stream = ByteBufferInputStream.wrap(buffer); for (ChunkDescriptor descriptor : chunks) { builder.add(descriptor, stream.sliceBuffers(descriptor.size), f); diff --git a/parquet-hadoop/src/test/resources/core-site.xml b/parquet-hadoop/src/test/resources/core-site.xml new file mode 100644 index 0000000000..a6344b87c9 --- /dev/null +++ b/parquet-hadoop/src/test/resources/core-site.xml @@ -0,0 +1,33 @@ + + + + + + fs.file.checksum.verify + false + + Disable checksum verification on the local file system used by tests. + Hadoop's ChecksumFileSystem.readVectored allocates checksum buffers via + the caller-supplied ByteBufferAllocator without releasing them; the + leaked buffers trip TrackingByteBufferAllocator leak detection in tests. + Turning off checksum verification skips the checksum read path entirely + and avoids the leak. + + + diff --git a/pom.xml b/pom.xml index 62bdb32623..2fb1a60215 100644 --- a/pom.xml +++ b/pom.xml @@ -83,7 +83,7 @@ 2.46.1 shaded.parquet - 3.3.0 + 3.4.2 2.13.0 1.17.0 thrift