From c255acc8d035be29fd708edac4822740cca7de69 Mon Sep 17 00:00:00 2001 From: Cheng Pan Date: Tue, 2 Jun 2026 11:47:23 +0800 Subject: [PATCH 1/2] HADOOP-19908. ZStandardDecompressor should reset compressedDirectBuf correctly --- .../io/compress/zstd/ZStandardCompressor.java | 2 +- .../compress/zstd/ZStandardDecompressor.java | 4 +- .../TestZStandardCompressorDecompressor.java | 58 +++++++++++++++++++ 3 files changed, 62 insertions(+), 2 deletions(-) diff --git a/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardCompressor.java b/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardCompressor.java index 8d989ec96db555..b029e3e7652c21 100644 --- a/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardCompressor.java +++ b/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardCompressor.java @@ -301,7 +301,7 @@ public void reset() { finished = false; bytesRead = 0; bytesWritten = 0; - uncompressedDirectBuf.rewind(); + uncompressedDirectBuf.clear(); uncompressedDirectBufOff = 0; uncompressedDirectBufLen = 0; keepUncompressedBuf = false; diff --git a/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardDecompressor.java b/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardDecompressor.java index c7134e7ff5213c..28eea96455c980 100644 --- a/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardDecompressor.java +++ b/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardDecompressor.java @@ -92,7 +92,7 @@ private void setInputFromSavedData() { bytesInCompressedBuffer = directBufferSize; } - compressedDirectBuf.rewind(); + compressedDirectBuf.clear(); compressedDirectBuf.put( userBuf, userBufOff, bytesInCompressedBuffer); @@ -225,6 +225,8 @@ public void reset() { finished = false; compressedDirectBufOff = 0; bytesInCompressedBuffer = 0; + compressedDirectBuf.limit(directBufferSize); + compressedDirectBuf.position(0); uncompressedDirectBuf.limit(directBufferSize); uncompressedDirectBuf.position(directBufferSize); userBufOff = 0; diff --git a/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/io/compress/zstd/TestZStandardCompressorDecompressor.java b/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/io/compress/zstd/TestZStandardCompressorDecompressor.java index e1546cdbbfa23b..41f63418b88ff4 100644 --- a/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/io/compress/zstd/TestZStandardCompressorDecompressor.java +++ b/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/io/compress/zstd/TestZStandardCompressorDecompressor.java @@ -557,6 +557,64 @@ public void testDecompressReturnsWhenNothingToDecompress() throws Exception { assertEquals(0, result); } + /** + * Verify that {@code setInput()} does not throw {@code BufferOverflowException} + * after a previous {@code decompress()} call threw an exception. + * + *

When {@code decompress()} processes compressed data, it sets + * {@code compressedDirectBuf.limit(bytesInCompressedBuffer)} — a value that + * may be smaller than {@code directBufferSize}. If {@code decompressDirectByteBufferStream} + * throws (e.g. on corrupted input), the limit is never restored. A subsequent + * {@code reset()} also does not restore {@code compressedDirectBuf.limit}. + * So the next {@code setInput()} call will hit {@code BufferOverflowException} + * because {@code setInputFromSavedData()} tries to {@code put()} more bytes + * than the current limit allows.

+ * + *

This scenario occurs in practice when reading multiple zstd-compressed + * files from a directory: a corrupted file causes an exception mid-decompress, + * the decompressor is returned to the pool and reset, but the limit stays + * small. The next file's {@code setInput()} then fails.

+ */ + @Test + public void testSetInputAfterDecompressThrowsOnCorruptedData() throws Exception { + byte[] rawData = generate(400); + int bufSize = IO_FILE_BUFFER_SIZE_DEFAULT; + + ByteArrayOutputStream baos = new ByteArrayOutputStream(); + try (CompressionOutputStream cos = new CompressorStream(baos, + new ZStandardCompressor(), bufSize)) { + cos.write(rawData); + } + byte[] compressed = baos.toByteArray(); + + // Corrupt the compressed data by dropping the first 10 bytes. + byte[] corrupted = new byte[compressed.length - 10]; + System.arraycopy(compressed, 10, corrupted, 0, corrupted.length); + + ZStandardDecompressor decompressor = new ZStandardDecompressor(bufSize); + byte[] out = new byte[bufSize]; + + // Feed corrupted data — decompress() sets limit to corrupted.length, then throws. + decompressor.setInput(corrupted, 0, corrupted.length); + try { + decompressor.decompress(out, 0, out.length); + } catch (Exception e) { + // Expected: corrupted data causes an exception. + } + + // Reset the decompressor (as the codec pool would). + decompressor.reset(); + + // Feed valid data — this must NOT throw BufferOverflowException. + decompressor.setInput(compressed, 0, compressed.length); + int n = decompressor.decompress(out, 0, out.length); + assertTrue(n >= 0, "decompress should return >= 0 after reset"); + + while (!decompressor.finished()) { + decompressor.decompress(out, 0, out.length); + } + } + // workers > 0 should produce data that round-trips correctly through the // decompressor, matching the bytes produced with the default workers=0. @Test From 13e44c8cef3d1513c5ecf3b12538afd162bde4d5 Mon Sep 17 00:00:00 2001 From: Cheng Pan Date: Tue, 2 Jun 2026 15:28:48 +0800 Subject: [PATCH 2/2] add comments and enhance tests --- .../hadoop/io/compress/zstd/ZStandardDecompressor.java | 3 +++ .../compress/zstd/TestZStandardCompressorDecompressor.java | 5 ++++- 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardDecompressor.java b/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardDecompressor.java index 28eea96455c980..e927f1233c4ced 100644 --- a/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardDecompressor.java +++ b/hadoop-common-project/hadoop-common/src/main/java/org/apache/hadoop/io/compress/zstd/ZStandardDecompressor.java @@ -188,6 +188,9 @@ public int decompress(byte[] b, int off, int len) } // Restore limit so setInputFromSavedData() can rewind+put on next call. + // There is possible that `decompressDirectByteBufferStream` throws exception + // on decompressing corrupt data, code won't reach here, the caller must + // call `reset` to clear the state before resuing this decompressor. compressedDirectBuf.limit(directBufferSize); } else { n = 0; diff --git a/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/io/compress/zstd/TestZStandardCompressorDecompressor.java b/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/io/compress/zstd/TestZStandardCompressorDecompressor.java index 41f63418b88ff4..176196f10a5be0 100644 --- a/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/io/compress/zstd/TestZStandardCompressorDecompressor.java +++ b/hadoop-common-project/hadoop-common/src/test/java/org/apache/hadoop/io/compress/zstd/TestZStandardCompressorDecompressor.java @@ -15,6 +15,7 @@ */ package org.apache.hadoop.io.compress.zstd; +import com.github.luben.zstd.ZstdException; import org.apache.commons.io.FileUtils; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.io.DataInputBuffer; @@ -50,6 +51,7 @@ import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.junit.jupiter.api.Assertions.fail; public class TestZStandardCompressorDecompressor { private final static char[] HEX_ARRAY = "0123456789ABCDEF".toCharArray(); @@ -598,7 +600,8 @@ public void testSetInputAfterDecompressThrowsOnCorruptedData() throws Exception decompressor.setInput(corrupted, 0, corrupted.length); try { decompressor.decompress(out, 0, out.length); - } catch (Exception e) { + fail("decompress should throw exception on corrupted data"); + } catch (ZstdException e) { // Expected: corrupted data causes an exception. }