From ac4a2ce429840304ca6f8848c66b705feec36989 Mon Sep 17 00:00:00 2001 From: Jay Zou Date: Sun, 30 Aug 2026 13:58:07 +0800 Subject: [PATCH] HDFS-17970. Exclude failed EC checksum target from reconstruction sources --- .../server/datanode/BlockChecksumHelper.java | 55 ++++++++++- .../apache/hadoop/hdfs/TestFileChecksum.java | 98 +++++++++++++++++++ 2 files changed, 150 insertions(+), 3 deletions(-) diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/BlockChecksumHelper.java b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/BlockChecksumHelper.java index 0955a3ff5cc28a..2370260ad3b893 100644 --- a/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/BlockChecksumHelper.java +++ b/hadoop-hdfs-project/hadoop-hdfs/src/main/java/org/apache/hadoop/hdfs/server/datanode/BlockChecksumHelper.java @@ -56,6 +56,7 @@ import java.io.IOException; import java.io.InputStream; import java.security.MessageDigest; +import java.util.BitSet; import java.util.HashMap; import java.util.Map; @@ -697,12 +698,60 @@ private void recalculateChecksum(int errBlkIndex, long blockLength) throws IOException { LOG.debug("Recalculate checksum for the missing/failed block index {}", errBlkIndex); - byte[] errIndices = new byte[1]; - errIndices[0] = (byte) errBlkIndex; + byte errIndex = (byte) errBlkIndex; + + // A failed checksum RPC does not remove the block from the location + // list. Exclude the target and place duplicate replicas after the + // independent sources, while retaining them as read fallbacks. + BitSet seenIndices = new BitSet(); + int numSources = 0; + for (byte blockIndex : blockIndices) { + if (blockIndex != errIndex) { + numSources++; + seenIndices.set(blockIndex); + } + } + + int cellsNum = (int) ((blockGroup.getNumBytes() - 1) + / ecPolicy.getCellSize() + 1); + int minRequiredSources = Math.min( + cellsNum, ecPolicy.getNumDataUnits()); + int numUniqueSources = seenIndices.cardinality(); + if (numUniqueSources < minRequiredSources) { + throw new IOException(String.format( + "Not enough unique sources to reconstruct block index %d " + + "in block group %s: required=%d, available=%d", + errBlkIndex, blockGroup, minRequiredSources, numUniqueSources)); + } + + byte[] sourceIndices = new byte[numSources]; + DatanodeInfo[] sourceDatanodes = new DatanodeInfo[numSources]; + + seenIndices.clear(); + int uniquePos = 0; + int duplicatePos = numUniqueSources; + for (int i = 0; i < blockIndices.length; i++) { + byte blockIndex = blockIndices[i]; + if (blockIndex == errIndex) { + continue; + } + + int sourcePos; + if (seenIndices.get(blockIndex)) { + sourcePos = duplicatePos++; + } else { + seenIndices.set(blockIndex); + sourcePos = uniquePos++; + } + + sourceIndices[sourcePos] = blockIndex; + sourceDatanodes[sourcePos] = datanodes[i]; + } StripedReconstructionInfo stripedReconInfo = new StripedReconstructionInfo( - blockGroup, ecPolicy, blockIndices, datanodes, errIndices); + blockGroup, ecPolicy, sourceIndices, sourceDatanodes, + new byte[] {errIndex}); BlockChecksumType groupChecksumType = getBlockChecksumOptions().getBlockChecksumType(); try (StripedBlockChecksumReconstructor checksumRecon = diff --git a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/TestFileChecksum.java b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/TestFileChecksum.java index 980c72958989c6..39d2e0e77f0dbc 100644 --- a/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/TestFileChecksum.java +++ b/hadoop-hdfs-project/hadoop-hdfs/src/test/java/org/apache/hadoop/hdfs/TestFileChecksum.java @@ -23,12 +23,24 @@ import org.apache.hadoop.fs.Options.ChecksumCombineMode; import org.apache.hadoop.fs.Path; import org.apache.hadoop.fs.permission.FsPermission; +import org.apache.hadoop.hdfs.protocol.BlockChecksumOptions; +import org.apache.hadoop.hdfs.protocol.BlockChecksumType; import org.apache.hadoop.hdfs.protocol.DatanodeInfo; import org.apache.hadoop.hdfs.protocol.ErasureCodingPolicy; import org.apache.hadoop.hdfs.protocol.LocatedBlock; import org.apache.hadoop.hdfs.protocol.LocatedBlocks; +import org.apache.hadoop.hdfs.protocol.LocatedStripedBlock; +import org.apache.hadoop.hdfs.protocol.StripedBlockInfo; +import org.apache.hadoop.hdfs.protocol.datatransfer.DataTransferProtoUtil; +import org.apache.hadoop.hdfs.protocol.datatransfer.IOStreamPair; +import org.apache.hadoop.hdfs.protocol.datatransfer.Sender; +import org.apache.hadoop.hdfs.protocol.proto.DataTransferProtos.BlockOpResponseProto; +import org.apache.hadoop.hdfs.protocol.proto.DataTransferProtos.OpBlockChecksumResponseProto; +import org.apache.hadoop.hdfs.protocolPB.PBHelperClient; +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.security.token.Token; import org.apache.hadoop.test.GenericTestUtils; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Timeout; @@ -39,6 +51,7 @@ import org.apache.hadoop.hdfs.client.HdfsClientConfigKeys; import org.slf4j.event.Level; +import java.io.DataOutputStream; import java.io.IOException; import java.util.Random; @@ -100,12 +113,19 @@ public static Object[] getParameters() { } public static void setup(String mode) throws IOException { + setup(mode, false); + } + + private static void setup(String mode, + boolean enableEcReconstructionValidation) throws IOException { checksumCombineMode = mode; int numDNs = dataBlocks + parityBlocks + 2; conf = new Configuration(); conf.setLong(DFSConfigKeys.DFS_BLOCK_SIZE_KEY, blockSize); conf.setInt(DFSConfigKeys.DFS_NAMENODE_REPLICATION_MAX_STREAMS_KEY, 0); conf.setBoolean(DFS_BLOCK_ACCESS_TOKEN_ENABLE_KEY, true); + conf.setBoolean(DFSConfigKeys.DFS_DN_EC_RECONSTRUCTION_VALIDATION_KEY, + enableEcReconstructionValidation); conf.set(HdfsClientConfigKeys.DFS_CHECKSUM_COMBINE_MODE_KEY, checksumCombineMode); cluster = new MiniDFSCluster.Builder(conf).numDataNodes(numDNs).build(); @@ -681,6 +701,84 @@ public void testStripedFileChecksumWithReconstructFail(String pMode) } } + @MethodSource("getParameters") + @ParameterizedTest + @Timeout(value = 90) + public void testStripedFileChecksumReconstructionExcludesFailedTarget( + String pMode) throws Exception { + setup(pMode, true); + BlockChecksumType checksumType = + ChecksumCombineMode.COMPOSITE_CRC.name().equals(pMode) + ? BlockChecksumType.COMPOSITE_CRC : BlockChecksumType.MD5CRC; + + assertBlockGroupChecksumWithFailedTarget( + ecDir + "/stripedFileChecksumFailedTarget", blockGroupSize, + checksumType); + assertBlockGroupChecksumWithFailedTarget( + ecDir + "/shortStripedFileChecksumFailedTarget", 100, + checksumType); + } + + private void assertBlockGroupChecksumWithFailedTarget(String stripedFile, + int fileLength, BlockChecksumType checksumType) throws Exception { + prepareTestFiles(fileLength, new String[] {stripedFile}); + + LocatedStripedBlock blockGroup = (LocatedStripedBlock) + client.getLocatedBlocks(stripedFile, 0).get(0); + OpBlockChecksumResponseProto expected = getBlockGroupChecksum(blockGroup, + blockGroup.getBlockTokens(), checksumType); + + // Keep the target live and fail only its child BLOCK_CHECKSUM request. + // Reconstruction obtains a fresh token for its READ_BLOCK requests. + Token[] invalidTokens = + blockGroup.getBlockTokens().clone(); + int targetPos = firstSelectedDataBlockPosition(blockGroup); + invalidTokens[targetPos] = new Token<>(); + + OpBlockChecksumResponseProto reconstructed = getBlockGroupChecksum( + blockGroup, invalidTokens, checksumType); + assertEquals(expected, reconstructed); + } + + private int firstSelectedDataBlockPosition(LocatedStripedBlock blockGroup) { + byte[] blockIndices = blockGroup.getBlockIndices(); + int cellsNum = (int) ((blockGroup.getBlock().getNumBytes() - 1) + / cellSize + 1); + int minRequiredSources = Math.min(cellsNum, dataBlocks); + for (int i = 0; i < minRequiredSources; i++) { + if (blockIndices[i] < dataBlocks) { + return i; + } + } + throw new AssertionError( + "No data block found in the selected reconstruction sources"); + } + + private OpBlockChecksumResponseProto getBlockGroupChecksum( + LocatedStripedBlock blockGroup, + Token[] blockTokens, + BlockChecksumType checksumType) throws IOException { + StripedBlockInfo stripedBlockInfo = new StripedBlockInfo( + blockGroup.getBlock(), blockGroup.getLocations(), blockTokens, + blockGroup.getBlockIndices(), ecPolicy); + int timeout = client.getConf().getChecksumEcSocketTimeout() + + client.getConf().getSocketTimeout(); + + try (IOStreamPair pair = client.connectToDN(blockGroup.getLocations()[0], + timeout, blockGroup.getBlockToken())) { + new Sender((DataOutputStream) pair.out).blockGroupChecksum( + stripedBlockInfo, blockGroup.getBlockToken(), + blockGroup.getBlock().getNumBytes(), + new BlockChecksumOptions(checksumType)); + + BlockOpResponseProto reply = BlockOpResponseProto.parseFrom( + PBHelperClient.vintPrefixed(pair.in)); + DataTransferProtoUtil.checkBlockOpStatus(reply, + "for block group " + blockGroup.getBlock()); + return reply.getChecksumResponse(); + } + } + @MethodSource("getParameters") @ParameterizedTest @Timeout(value = 90)