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]