diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/ozone/client/io/ECBlockReconstructedStripeInputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/ozone/client/io/ECBlockReconstructedStripeInputStream.java index c71db0e41ed4..672894f186f1 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/ozone/client/io/ECBlockReconstructedStripeInputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/ozone/client/io/ECBlockReconstructedStripeInputStream.java @@ -727,6 +727,11 @@ protected int readWithStrategy(ByteReaderStrategy strategy) { public synchronized void close() { super.close(); freeBuffers(); + if (decoder != null) { + // Release the decoder's native resources (e.g. ISA-L coder context) + // that were allocated lazily in init(). + decoder.release(); + } } @Override diff --git a/hadoop-hdds/client/src/test/java/org/apache/hadoop/ozone/client/io/TestECBlockReconstructedStripeInputStream.java b/hadoop-hdds/client/src/test/java/org/apache/hadoop/ozone/client/io/TestECBlockReconstructedStripeInputStream.java index 8589a2b73f40..f6892940a069 100644 --- a/hadoop-hdds/client/src/test/java/org/apache/hadoop/ozone/client/io/TestECBlockReconstructedStripeInputStream.java +++ b/hadoop-hdds/client/src/test/java/org/apache/hadoop/ozone/client/io/TestECBlockReconstructedStripeInputStream.java @@ -26,6 +26,10 @@ 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.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mockStatic; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; import com.google.common.collect.ImmutableSet; import java.io.IOException; @@ -52,11 +56,15 @@ import org.apache.hadoop.io.ElasticByteBufferPool; import org.apache.hadoop.ozone.client.io.ECStreamTestUtil.TestBlockInputStream; import org.apache.hadoop.ozone.client.io.ECStreamTestUtil.TestBlockInputStreamFactory; +import org.apache.ozone.erasurecode.rawcoder.RSRawErasureCoderFactory; +import org.apache.ozone.erasurecode.rawcoder.RawErasureDecoder; +import org.apache.ozone.erasurecode.rawcoder.util.CodecUtil; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.MethodSource; +import org.mockito.MockedStatic; /** * Test for the ECBlockReconstructedStripeInputStream. @@ -807,6 +815,40 @@ public void testFailedLocationsAreNotRead() throws IOException { } } + @Test + public void testDecoderReleasedOnClose() throws IOException { + // The stream must release its RawErasureDecoder on close so the decoder's + // native resources (e.g. ISA-L coder context) are not leaked. The decoder + // is created lazily in init() on the first read. + int chunkSize = repConfig.getEcChunkSize(); + int blockLength = chunkSize * repConfig.getData(); + ByteBuffer[] dataBufs = allocateBuffers(repConfig.getData(), chunkSize); + ECStreamTestUtil.randomFill(dataBufs, chunkSize, dataGen, blockLength); + ByteBuffer[] parity = generateParity(dataBufs, repConfig); + addDataStreamsToFactory(dataBufs, parity); + + // Data index 1 is missing, so a read must build the decoder to reconstruct. + Map dnMap = + ECStreamTestUtil.createIndexMap(2, 3, 4, 5); + BlockLocationInfo keyInfo = + ECStreamTestUtil.createKeyInfo(repConfig, blockLength, dnMap); + streamFactory.setCurrentPipeline(keyInfo.getPipeline()); + + RawErasureDecoder spyDecoder = + spy(new RSRawErasureCoderFactory().createDecoder(repConfig)); + ByteBuffer[] bufs = allocateByteBuffers(repConfig); + try (MockedStatic mocked = mockStatic(CodecUtil.class)) { + mocked.when(() -> CodecUtil.createRawDecoderWithFallback(any())) + .thenReturn(spyDecoder); + try (ECBlockReconstructedStripeInputStream ecBlockReconstructedStripeInputStream = + createInputStream(keyInfo)) { + // The read triggers init(), which creates the (spied) decoder. + ecBlockReconstructedStripeInputStream.read(bufs); + } + } + verify(spyDecoder).release(); + } + private ECBlockReconstructedStripeInputStream createInputStream( BlockLocationInfo keyInfo) { OzoneClientConfig clientConfig = conf.getObject(OzoneClientConfig.class); diff --git a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/ECKeyOutputStream.java b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/ECKeyOutputStream.java index 9f94384d4df2..d757256aeef7 100644 --- a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/ECKeyOutputStream.java +++ b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/io/ECKeyOutputStream.java @@ -499,6 +499,9 @@ public void close() throws IOException { } finally { closeCurrentStreamEntry(); blockOutputStreamEntryPool.cleanup(); + // Release the encoder's native resources (e.g. ISA-L coder context) + // that were allocated in the constructor. + encoder.release(); } } diff --git a/hadoop-ozone/client/src/test/java/org/apache/hadoop/ozone/client/TestOzoneECClient.java b/hadoop-ozone/client/src/test/java/org/apache/hadoop/ozone/client/TestOzoneECClient.java index 6317baac800d..502e11a30950 100644 --- a/hadoop-ozone/client/src/test/java/org/apache/hadoop/ozone/client/TestOzoneECClient.java +++ b/hadoop-ozone/client/src/test/java/org/apache/hadoop/ozone/client/TestOzoneECClient.java @@ -22,6 +22,10 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mockStatic; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; import java.io.IOException; import java.nio.ByteBuffer; @@ -64,11 +68,13 @@ import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos; import org.apache.ozone.erasurecode.rawcoder.RSRawErasureCoderFactory; import org.apache.ozone.erasurecode.rawcoder.RawErasureEncoder; +import org.apache.ozone.erasurecode.rawcoder.util.CodecUtil; import org.apache.ozone.test.GenericTestUtils; import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.mockito.MockedStatic; /** * Real unit test for OzoneECClient. @@ -604,6 +610,21 @@ public void test10D4PConfigWithPartialStripe() } } + @Test + public void testEncoderReleasedOnClose() throws IOException { + // The ECKeyOutputStream must release its RawErasureEncoder on close so the + // encoder's native resources (e.g. ISA-L coder context) are not leaked. + RawErasureEncoder spyEncoder = spy(new RSRawErasureCoderFactory() + .createEncoder(new ECReplicationConfig(dataBlocks, parityBlocks))); + try (MockedStatic mocked = mockStatic(CodecUtil.class)) { + mocked.when(() -> CodecUtil.createRawEncoderWithFallback(any())) + .thenReturn(spyEncoder); + // Creates the ECKeyOutputStream (picks up the spy encoder) and closes it. + writeIntoECKey(inputChunks, keyName, null); + } + verify(spyEncoder).release(); + } + @Test public void testWriteShouldFailIfMoreThanParityNodesFail() throws Exception {