This is an automated email from the ASF dual-hosted git repository.

ZanderXu pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/hadoop.git


The following commit(s) were added to refs/heads/trunk by this push:
     new 574159caa96 HDFS-17921. StripedBlockWriter should free up buffer if 
init fails (#8498)
574159caa96 is described below

commit 574159caa96aa88feb83b0a8d4f874b95ec45293
Author: Felix Nguyen <[email protected]>
AuthorDate: Mon Jun 1 10:34:27 2026 +0700

    HDFS-17921. StripedBlockWriter should free up buffer if init fails (#8498)
---
 .../server/datanode/DataNodeFaultInjector.java     |  8 +-
 .../datanode/erasurecode/StripedBlockWriter.java   |  8 +-
 .../server/datanode/erasurecode/StripedWriter.java |  5 +-
 .../hadoop/hdfs/TestReconstructStripedFile.java    | 92 ++++++++++++++++++++++
 4 files changed, 106 insertions(+), 7 deletions(-)

diff --git 
a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/DataNodeFaultInjector.java
 
b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/DataNodeFaultInjector.java
index 9e046cc3600..d54e1d9ce61 100644
--- 
a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/DataNodeFaultInjector.java
+++ 
b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/DataNodeFaultInjector.java
@@ -42,8 +42,6 @@ public static void set(DataNodeFaultInjector injector) {
     instance = injector;
   }
 
-  public void getHdfsBlocksMetadata() {}
-
   public void writeBlockAfterFlush() throws IOException {}
 
   public void sendShortCircuitShmResponse() throws IOException {}
@@ -113,6 +111,12 @@ public void throwTooManyOpenFiles() throws 
FileNotFoundException {
    */
   public void stripedBlockReconstruction() throws IOException {}
 
+  /**
+   * Used as a hook to inject failure when initializing a striped block writer.
+   */
+  public void stripedBlockWriterInit(ByteBuffer targetBuffer) throws 
IOException {
+  }
+
   /**
    * Used as a hook to inject failure in erasure coding checksum reconstruction
    * process.
diff --git 
a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/erasurecode/StripedBlockWriter.java
 
b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/erasurecode/StripedBlockWriter.java
index 24c1d613822..5b8b22c3c11 100644
--- 
a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/erasurecode/StripedBlockWriter.java
+++ 
b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/erasurecode/StripedBlockWriter.java
@@ -31,6 +31,7 @@
 import 
org.apache.hadoop.hdfs.protocol.datatransfer.sasl.DataEncryptionKeyFactory;
 import org.apache.hadoop.hdfs.security.token.block.BlockTokenIdentifier;
 import org.apache.hadoop.hdfs.server.datanode.DataNode;
+import org.apache.hadoop.hdfs.server.datanode.DataNodeFaultInjector;
 import org.apache.hadoop.io.ByteBufferPool;
 import org.apache.hadoop.io.ElasticByteBufferPool;
 import org.apache.hadoop.io.IOUtils;
@@ -94,7 +95,10 @@ ByteBuffer getTargetBuffer() {
   }
 
   void freeTargetBuffer() {
-    targetBuffer = null;
+    if (targetBuffer != null) {
+      stripedWriter.getReconstructor().freeBuffer(targetBuffer);
+      targetBuffer = null;
+    }
   }
 
   /**
@@ -115,6 +119,7 @@ private void init() throws IOException {
       socket.setTcpNoDelay(
           datanode.getDnConf().getDataTransferServerTcpNoDelay());
       socket.setSoTimeout(datanode.getDnConf().getSocketTimeout());
+      DataNodeFaultInjector.get().stripedBlockWriterInit(targetBuffer);
 
       Token<BlockTokenIdentifier> blockToken =
           datanode.getBlockAccessToken(block,
@@ -151,6 +156,7 @@ private void init() throws IOException {
       success = true;
     } finally {
       if (!success) {
+        freeTargetBuffer();
         IOUtils.closeStream(out);
         IOUtils.closeStream(in);
         IOUtils.closeStream(socket);
diff --git 
a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/erasurecode/StripedWriter.java
 
b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/erasurecode/StripedWriter.java
index 00be1279c81..a795c8fd5ab 100644
--- 
a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/erasurecode/StripedWriter.java
+++ 
b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/erasurecode/StripedWriter.java
@@ -309,10 +309,7 @@ void clearBuffers() {
 
   void close() {
     for (StripedBlockWriter writer : writers) {
-      ByteBuffer targetBuffer =
-          writer != null ? writer.getTargetBuffer() : null;
-      if (targetBuffer != null) {
-        reconstructor.freeBuffer(targetBuffer);
+      if (writer != null) {
         writer.freeTargetBuffer();
       }
     }
diff --git 
a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/TestReconstructStripedFile.java
 
b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/TestReconstructStripedFile.java
index 3f4eca21bc8..fcbfc44872f 100644
--- 
a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/TestReconstructStripedFile.java
+++ 
b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/TestReconstructStripedFile.java
@@ -40,6 +40,9 @@
 import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.hadoop.fs.StorageType;
 import org.apache.hadoop.hdfs.protocol.LocatedBlock;
 import 
org.apache.hadoop.hdfs.server.datanode.erasurecode.ErasureCodingTestHelper;
 import org.apache.hadoop.io.ElasticByteBufferPool;
@@ -51,6 +54,7 @@
 import org.apache.hadoop.hdfs.client.HdfsClientConfigKeys;
 import org.apache.hadoop.hdfs.protocol.DatanodeID;
 import org.apache.hadoop.hdfs.protocol.DatanodeInfo;
+import org.apache.hadoop.hdfs.protocol.DatanodeInfo.DatanodeInfoBuilder;
 import org.apache.hadoop.hdfs.protocol.ErasureCodingPolicy;
 import org.apache.hadoop.hdfs.protocol.ExtendedBlock;
 import org.apache.hadoop.hdfs.protocol.LocatedBlocks;
@@ -63,6 +67,7 @@
 import org.apache.hadoop.hdfs.server.namenode.FSNamesystem;
 import org.apache.hadoop.hdfs.server.protocol.DatanodeStorage;
 import 
org.apache.hadoop.hdfs.server.protocol.BlockECReconstructionCommand.BlockECReconstructionInfo;
+import org.apache.hadoop.hdfs.server.protocol.StorageReport;
 import org.apache.hadoop.hdfs.util.StripedBlockUtil;
 import org.apache.hadoop.io.erasurecode.CodecUtil;
 import org.apache.hadoop.io.erasurecode.ErasureCodeNative;
@@ -495,6 +500,93 @@ public void 
testProcessErasureCodingTasksSubmitionShouldSucceed()
     dataNode.getErasureCodingWorker().processErasureCodingTasks(ecTasks);
   }
 
+  @Test
+  @Timeout(value = 120)
+  public void testTargetBufferLeakWhenBlockWriterInitFails() throws Exception {
+    int fileLen = 1;
+    String testFile = "/buffer-leak";
+    Path path = new Path(testFile);
+    writeFile(fs, testFile, fileLen);
+
+    LocatedBlocks locatedBlocks = StripedFileTestUtil.getLocatedBlocks(path, 
fs);
+    LocatedStripedBlock lastBlock = (LocatedStripedBlock) 
locatedBlocks.getLastLocatedBlock();
+    DatanodeInfo[] storageInfos = lastBlock.getLocations();
+    byte[] indices = lastBlock.getBlockIndices();
+
+    DatanodeInfo[] sources = new DatanodeInfo[storageInfos.length - 1];
+    byte[] liveBlockIndices = new byte[indices.length - 1];
+    for (int i = 1; i < storageInfos.length; i++) {
+      sources[i - 1] = storageInfos[i];
+      liveBlockIndices[i - 1] = indices[i];
+    }
+
+    int targetDN = dnMap.get(storageInfos[0]);
+    DataNode targetDataNode = cluster.getDataNodes().get(targetDN);
+    DatanodeInfo target =
+        new 
DatanodeInfoBuilder().setNodeID(targetDataNode.getDatanodeId()).build();
+    StorageReport[] storageReports =
+        
targetDataNode.getFSDataset().getStorageReports(lastBlock.getBlock().getBlockPoolId());
+    assertTrue(storageReports.length > 0);
+    DatanodeStorage targetStorage = storageReports[0].getStorage();
+
+    ElasticByteBufferPool bufferPool =
+        (ElasticByteBufferPool) ErasureCodingTestHelper.getBufferPool();
+    emptyBufferPool(bufferPool, true);
+    emptyBufferPool(bufferPool, false);
+
+    DataNodeFaultInjector oldInjector = DataNodeFaultInjector.get();
+    AtomicReference<ByteBuffer> bufferRef = new AtomicReference<>();
+    DataNodeFaultInjector injector = new DataNodeFaultInjector() {
+      @Override
+      public void stripedBlockWriterInit(ByteBuffer targetBuffer) throws 
IOException {
+        bufferRef.set(targetBuffer);
+        throw new IOException("Failed to init writer");
+      }
+    };
+
+    DataNode dataNode = cluster.getDataNode(sources[0].getIpcPort());
+    BlockECReconstructionInfo reconstructionInfo =
+        new BlockECReconstructionInfo(lastBlock.getBlock(), sources,
+            new DatanodeInfo[] {target},
+            new String[] {targetStorage.getStorageID()},
+            new StorageType[] {targetStorage.getStorageType()},
+            liveBlockIndices, new byte[0], ecPolicy);
+    List<BlockECReconstructionInfo> ecTasks = new ArrayList<>();
+    ecTasks.add(reconstructionInfo);
+
+    DataNodeFaultInjector.set(injector);
+    try {
+      dataNode.getErasureCodingWorker().processErasureCodingTasks(ecTasks);
+      GenericTestUtils.waitFor(() -> {
+        ByteBuffer expectedBuffer = bufferRef.get();
+        if (expectedBuffer == null) {
+          return false;
+        }
+
+        List<ByteBuffer> buffs = new ArrayList<>();
+        try {
+          boolean isDirectBuff = expectedBuffer.isDirect();
+          while (bufferPool.size(isDirectBuff) > 0) {
+            ByteBuffer buff = bufferPool.getBuffer(isDirectBuff, 1);
+            buffs.add(buff);
+            if (buff == expectedBuffer) {
+              // The allocated buffer has been returned to the buffer pool
+              return true;
+            }
+          }
+          return false;
+        } finally {
+          // Put the buffers back
+          for (ByteBuffer buff : buffs) {
+            bufferPool.putBuffer(buff);
+          }
+        }
+      }, 100, 30000);
+    } finally {
+      DataNodeFaultInjector.set(oldInjector);
+    }
+  }
+
   // HDFS-12044
   @Test
   @Timeout(value = 120)


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to