This is an automated email from the ASF dual-hosted git repository. jojochuang pushed a commit to branch ozone-2.1 in repository https://gitbox.apache.org/repos/asf/ozone.git
commit 4901d9a23409c2541aa60fa7dabe6bca3b3f2c9c Author: Aswin Shakil Balasubramanian <[email protected]> AuthorDate: Sat Jul 18 00:12:27 2026 +0530 HDDS-15791. EC reconstructed RECOVERING container while still idle can be deleted by SCM and recreated as new partial OPEN container (#10702) (cherry picked from commit 458c0d14498231b558f131fe0de26a7754a4cea3) Change-Id: I9cb787e5f3c72d32f2099a17fc2e80770331bcb6 --- .../hadoop/hdds/scm/storage/BlockOutputStream.java | 16 +- .../hdds/scm/storage/ECBlockOutputStream.java | 25 +- .../storage/TestBlockOutputStreamCorrectness.java | 52 +++- .../hdds/scm/storage/ContainerProtocolCalls.java | 35 ++- .../container/common/helpers/ContainerUtils.java | 26 ++ .../ozone/container/common/impl/ContainerSet.java | 95 ++----- .../container/common/impl/HddsDispatcher.java | 8 + .../ECReconstructionCoordinator.java | 3 +- .../ozone/container/keyvalue/KeyValueHandler.java | 9 + .../StaleRecoveringContainerScrubbingService.java | 21 +- ...stStaleRecoveringContainerScrubbingService.java | 50 ++++ .../container/common/impl/TestHddsDispatcher.java | 272 ++++++++++++++++++++- .../src/main/proto/DatanodeClientProtocol.proto | 2 + 13 files changed, 519 insertions(+), 95 deletions(-) diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java index a03fd037bda..ca2c60f30fe 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java @@ -596,7 +596,7 @@ CompletableFuture<PutBlockResult> executePutBlock(boolean close, // if block is full, send the eof boolean isBlockFull = (blockSize != -1 && flushPos == blockSize); - asyncReply = putBlockAsync(xceiverClient, blockData, close || isBlockFull, tokenString); + asyncReply = putBlockAsync(xceiverClient, blockData, close || isBlockFull, tokenString, containerAutoCreate()); CompletableFuture<ContainerCommandResponseProto> future = asyncReply.getResponse(); flushFuture = future.thenApplyAsync(e -> { try { @@ -963,7 +963,7 @@ private CompletableFuture<PutBlockResult> writeChunkToContainer( } asyncReply = writeChunkAsync(xceiverClient, chunkInfo, - blockID.get(), data, tokenString, replicationIndex, blockData, close); + blockID.get(), data, tokenString, replicationIndex, blockData, close, containerAutoCreate()); CompletableFuture<ContainerCommandResponseProto> respFuture = asyncReply.getResponse(); validateFuture = respFuture.thenApplyAsync(e -> { @@ -1159,6 +1159,18 @@ private ChunkInfo createChunkInfo(long lastPartialChunkOffset) return revisedChunkInfo.build(); } + /** + * @return true when the DataNode may auto-create a missing container for this write. + */ + protected boolean containerAutoCreate() { + return true; + } + + @VisibleForTesting + public boolean isContainerAutoCreate() { + return containerAutoCreate(); + } + private boolean isFullChunk(ChunkInfo chunkInfo) { Preconditions.checkState( chunkInfo.getLen() <= config.getStreamBufferSize()); diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/ECBlockOutputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/ECBlockOutputStream.java index d798b3a9385..5f33e02f4d5 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/ECBlockOutputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/ECBlockOutputStream.java @@ -59,6 +59,7 @@ public class ECBlockOutputStream extends BlockOutputStream { private final DatanodeDetails datanodeDetails; + private final boolean containerAutoCreate; private CompletableFuture<ContainerProtos.ContainerCommandResponseProto> currentChunkRspFuture = null; @@ -83,11 +84,33 @@ public ECBlockOutputStream( Token<? extends TokenIdentifier> token, ContainerClientMetrics clientMetrics, StreamBufferArgs streamBufferArgs, Supplier<ExecutorService> executorServiceSupplier + ) throws IOException { + this(blockID, xceiverClientManager, pipeline, bufferPool, config, token, clientMetrics, + streamBufferArgs, executorServiceSupplier, true); + } + + @SuppressWarnings("checkstyle:ParameterNumber") + public ECBlockOutputStream( + BlockID blockID, + XceiverClientFactory xceiverClientManager, + Pipeline pipeline, + BufferPool bufferPool, + OzoneClientConfig config, + Token<? extends TokenIdentifier> token, + ContainerClientMetrics clientMetrics, StreamBufferArgs streamBufferArgs, + Supplier<ExecutorService> executorServiceSupplier, + boolean containerAutoCreate ) throws IOException { super(blockID, -1, xceiverClientManager, pipeline, bufferPool, config, token, clientMetrics, streamBufferArgs, executorServiceSupplier); // In EC stream, there will be only one node in pipeline. this.datanodeDetails = pipeline.getClosestNode(); + this.containerAutoCreate = containerAutoCreate; + } + + @Override + protected boolean containerAutoCreate() { + return containerAutoCreate; } @Override @@ -272,7 +295,7 @@ public CompletableFuture<PutBlockResult> executePutBlock(boolean close, try { ContainerProtos.BlockData blockData = getContainerBlockData().build(); XceiverClientReply asyncReply = - putBlockAsync(getXceiverClient(), blockData, close, getTokenString()); + putBlockAsync(getXceiverClient(), blockData, close, getTokenString(), containerAutoCreate()); CompletableFuture<ContainerProtos.ContainerCommandResponseProto> future = asyncReply.getResponse(); flushFuture = future.thenApplyAsync(e -> { diff --git a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockOutputStreamCorrectness.java b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockOutputStreamCorrectness.java index 440b5b3d4d5..81deaaf36bb 100644 --- a/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockOutputStreamCorrectness.java +++ b/hadoop-hdds/client/src/test/java/org/apache/hadoop/hdds/scm/storage/TestBlockOutputStreamCorrectness.java @@ -19,6 +19,8 @@ import static java.util.concurrent.Executors.newFixedThreadPool; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.any; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; @@ -132,6 +134,46 @@ public void testMissingStripeChecksumDoesNotMakeExecutePutBlockFailDuringECRecon } } + @Test + public void testEcReconstructionStreamDisablesContainerAutoCreate() throws IOException { + OzoneClientConfig config = new OzoneClientConfig(); + ECReplicationConfig replicationConfig = new ECReplicationConfig(3, 2); + BlockID blockID = new BlockID(1, 1); + DatanodeDetails datanodeDetails = MockDatanodeDetails.randomDatanodeDetails(); + Pipeline pipeline = Pipeline.newBuilder() + .setId(datanodeDetails.getID()) + .setReplicationConfig(replicationConfig) + .setNodes(ImmutableList.of(datanodeDetails)) + .setState(Pipeline.PipelineState.CLOSED) + .setReplicaIndexes(ImmutableMap.of(datanodeDetails, 2)) + .build(); + + try (ECBlockOutputStream ecBlockOutputStream = createECBlockOutputStream(config, replicationConfig, + blockID, pipeline, false)) { + assertFalse(ecBlockOutputStream.isContainerAutoCreate()); + } + } + + @Test + public void testEcClientStreamAllowsContainerAutoCreate() throws IOException { + OzoneClientConfig config = new OzoneClientConfig(); + ECReplicationConfig replicationConfig = new ECReplicationConfig(3, 2); + BlockID blockID = new BlockID(1, 1); + DatanodeDetails datanodeDetails = MockDatanodeDetails.randomDatanodeDetails(); + Pipeline pipeline = Pipeline.newBuilder() + .setId(datanodeDetails.getID()) + .setReplicationConfig(replicationConfig) + .setNodes(ImmutableList.of(datanodeDetails)) + .setState(Pipeline.PipelineState.CLOSED) + .setReplicaIndexes(ImmutableMap.of(datanodeDetails, 2)) + .build(); + + try (ECBlockOutputStream ecBlockOutputStream = createECBlockOutputStream(config, replicationConfig, + blockID, pipeline)) { + assertTrue(ecBlockOutputStream.isContainerAutoCreate()); + } + } + /** * Creates a BlockData array with {@link ECReplicationConfig#getRequiredNodes()} number of elements. */ @@ -183,7 +225,8 @@ private BlockOutputStream createBlockOutputStream(BufferPool bufferPool) } private ECBlockOutputStream createECBlockOutputStream(OzoneClientConfig clientConfig, - ECReplicationConfig repConfig, BlockID blockID, Pipeline pipeline) throws IOException { + ECReplicationConfig repConfig, BlockID blockID, Pipeline pipeline, + boolean containerAutoCreate) throws IOException { final XceiverClientManager xcm = mock(XceiverClientManager.class); when(xcm.acquireClient(any())) .thenReturn(new MockXceiverClientSpi(pipeline)); @@ -193,7 +236,12 @@ private ECBlockOutputStream createECBlockOutputStream(OzoneClientConfig clientCo StreamBufferArgs.getDefaultStreamBufferArgs(repConfig, clientConfig); return new ECBlockOutputStream(blockID, xcm, pipeline, BufferPool.empty(), clientConfig, null, - clientMetrics, streamBufferArgs, () -> newFixedThreadPool(2)); + clientMetrics, streamBufferArgs, () -> newFixedThreadPool(2), containerAutoCreate); + } + + private ECBlockOutputStream createECBlockOutputStream(OzoneClientConfig clientConfig, + ECReplicationConfig repConfig, BlockID blockID, Pipeline pipeline) throws IOException { + return createECBlockOutputStream(clientConfig, repConfig, blockID, pipeline, true); } /** diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java index a934fc51372..1bcbc64248b 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java @@ -288,8 +288,18 @@ public static XceiverClientReply putBlockAsync(XceiverClientSpi xceiverClient, boolean eof, String tokenString) throws IOException, InterruptedException, ExecutionException { + return putBlockAsync(xceiverClient, containerBlockData, eof, tokenString, true); + } + + public static XceiverClientReply putBlockAsync(XceiverClientSpi xceiverClient, + BlockData containerBlockData, + boolean eof, + String tokenString, + boolean containerAutoCreate) + throws IOException, InterruptedException, ExecutionException { final ContainerCommandRequestProto request = getPutBlockRequest( - xceiverClient.getPipeline(), containerBlockData, eof, tokenString); + xceiverClient.getPipeline(), containerBlockData, eof, tokenString, + containerAutoCreate); return xceiverClient.sendCommandAsync(request); } @@ -326,10 +336,19 @@ public static ContainerProtos.FinalizeBlockResponseProto finalizeBlock( public static ContainerCommandRequestProto getPutBlockRequest( Pipeline pipeline, BlockData containerBlockData, boolean eof, String tokenString) throws IOException { + return getPutBlockRequest(pipeline, containerBlockData, eof, tokenString, true); + } + + public static ContainerCommandRequestProto getPutBlockRequest( + Pipeline pipeline, BlockData containerBlockData, boolean eof, + String tokenString, boolean containerAutoCreate) throws IOException { PutBlockRequestProto.Builder createBlockRequest = PutBlockRequestProto.newBuilder() .setBlockData(containerBlockData) .setEof(eof); + if (!containerAutoCreate) { + createBlockRequest.setContainerAutoCreate(false); + } final String id = pipeline.getFirstNode().getUuidString(); ContainerCommandRequestProto.Builder builder = ContainerCommandRequestProto.newBuilder().setCmdType(Type.PutBlock) @@ -437,6 +456,17 @@ public static XceiverClientReply writeChunkAsync( ByteString data, String tokenString, int replicationIndex, BlockData blockData, boolean close) throws IOException, ExecutionException, InterruptedException { + return writeChunkAsync(xceiverClient, chunk, blockID, data, tokenString, + replicationIndex, blockData, close, true); + } + + @SuppressWarnings("parameternumber") + public static XceiverClientReply writeChunkAsync( + XceiverClientSpi xceiverClient, ChunkInfo chunk, BlockID blockID, + ByteString data, String tokenString, + int replicationIndex, BlockData blockData, boolean close, + boolean containerAutoCreate) + throws IOException, ExecutionException, InterruptedException { WriteChunkRequestProto.Builder writeChunkRequest = WriteChunkRequestProto.newBuilder() @@ -455,6 +485,9 @@ public static XceiverClientReply writeChunkAsync( .setEof(close); writeChunkRequest.setBlock(createBlockRequest); } + if (!containerAutoCreate) { + writeChunkRequest.setContainerAutoCreate(false); + } String id = xceiverClient.getPipeline().getFirstNode().getUuidString(); ContainerCommandRequestProto.Builder builder = ContainerCommandRequestProto.newBuilder() diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/helpers/ContainerUtils.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/helpers/ContainerUtils.java index 07f4e09aa9e..253433824cc 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/helpers/ContainerUtils.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/helpers/ContainerUtils.java @@ -329,4 +329,30 @@ public static void assertSpaceAvailability(long containerId, HddsVolume volume, + currentUsage + ", minimum free space spared=" + spared, DISK_OUT_OF_SPACE); } } + + /** + * @return true if the DataNode may auto-create a missing container for this request + */ + public static boolean isContainerCreatable(ContainerCommandRequestProto request) { + switch (request.getCmdType()) { + case PutBlock: + return isContainerAutoCreateAllowed(request.getPutBlock()); + case WriteChunk: + return isContainerAutoCreateAllowed(request.getWriteChunk()); + case PutSmallFile: + return isContainerAutoCreateAllowed(request.getPutSmallFile().getBlock()); + default: + return true; + } + } + + private static boolean isContainerAutoCreateAllowed( + ContainerProtos.PutBlockRequestProto putBlock) { + return !putBlock.hasContainerAutoCreate() || putBlock.getContainerAutoCreate(); + } + + private static boolean isContainerAutoCreateAllowed( + ContainerProtos.WriteChunkRequestProto writeChunk) { + return !writeChunk.hasContainerAutoCreate() || writeChunk.getContainerAutoCreate(); + } } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/ContainerSet.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/ContainerSet.java index e351b9360a7..5953b4e5eaf 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/ContainerSet.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/ContainerSet.java @@ -34,6 +34,7 @@ import java.util.Map; import java.util.Objects; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentNavigableMap; import java.util.concurrent.ConcurrentSkipListMap; import java.util.concurrent.ConcurrentSkipListSet; @@ -68,8 +69,7 @@ public class ContainerSet implements Iterable<Container<?>> { private final ConcurrentSkipListSet<Long> missingContainerSet = new ConcurrentSkipListSet<>(); - private final ConcurrentSkipListSet<RecoveringContainer> recoveringContainerSet = - new ConcurrentSkipListSet<>(); + private final ConcurrentHashMap<Long, Long> recoveringContainerMap = new ConcurrentHashMap<>(); private final Clock clock; private long recoveringTimeout; @Nullable @@ -203,9 +203,7 @@ private boolean addContainer(Container<?> container, boolean overwrite) throws updateContainerIdTable(containerId, container.getContainerData()); missingContainerSet.remove(containerId); if (container.getContainerData().getState() == RECOVERING) { - recoveringContainerSet.add( - new RecoveringContainer(clock.millis() + recoveringTimeout, - containerId)); + recoveringContainerMap.put(containerId, getCurrentTime() + recoveringTimeout); } return true; } else { @@ -326,22 +324,20 @@ private void deleteFromContainerTable(long containerId) throws StorageContainerE public boolean removeRecoveringContainer(long containerId) { Preconditions.checkState(containerId >= 0, "Container Id cannot be negative."); - //it might take a little long time to iterate all the entries - // in recoveringContainerSet, but it seems ok here since: - // 1 In the vast majority of cases,there will not be too - // many recovering containers. - // 2 closing container is not a sort of urgent action - // - // we can revisit here if any performance problem happens - Iterator<RecoveringContainer> it = getRecoveringContainerIterator(); - while (it.hasNext()) { - RecoveringContainer entry = it.next(); - if (entry.getContainerId() == containerId) { - it.remove(); - return true; - } - } - return false; + return recoveringContainerMap.remove(containerId) != null; + } + + /** + * Reset the stale recovering scrub deadline for an active RECOVERING container. + */ + public void updateRecoveringContainerTimeout(long containerId) { + Preconditions.checkState(containerId >= 0, "Container Id cannot be negative."); + recoveringContainerMap.put(containerId, getCurrentTime() + recoveringTimeout); + } + + @VisibleForTesting + public Map<Long, Long> getRecoveringContainerMap() { + return recoveringContainerMap; } /** @@ -393,15 +389,6 @@ public Iterator<Container<?>> iterator() { return containerMap.values().iterator(); } - /** - * Return an container Iterator over - * {@link ContainerSet#recoveringContainerSet}. - * @return {@literal Iterator<RecoveringContainer>} - */ - public Iterator<RecoveringContainer> getRecoveringContainerIterator() { - return recoveringContainerSet.iterator(); - } - /** * Return an iterator of containers associated with the specified volume. * The iterator is sorted by last data scan timestamp in increasing order. @@ -572,52 +559,4 @@ public <T> void buildMissingContainerSetAndValidate(Map<T, Long> container2BCSID } }); } - - /** - * A class that holds information about a recovering container. - */ - public static class RecoveringContainer - implements Comparable<RecoveringContainer> { - private final long timeout; - private final long containerId; - - public RecoveringContainer(long timeout, long containerId) { - this.timeout = timeout; - this.containerId = containerId; - } - - public long getTimeout() { - return timeout; - } - - public long getContainerId() { - return containerId; - } - - @Override - public int compareTo(RecoveringContainer other) { - int timeoutCompare = Long.compare(this.timeout, other.timeout); - if (timeoutCompare != 0) { - return timeoutCompare; - } - return Long.compare(this.containerId, other.containerId); - } - - @Override - public boolean equals(Object o) { - if (this == o) { - return true; - } - if (o == null || getClass() != o.getClass()) { - return false; - } - RecoveringContainer that = (RecoveringContainer) o; - return timeout == that.timeout && containerId == that.containerId; - } - - @Override - public int hashCode() { - return Objects.hash(timeout, containerId); - } - } } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java index ea47c4945b8..357402a8e05 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java @@ -291,6 +291,14 @@ && getMissingContainerSet().contains(containerID)) { if (container == null && ((isWriteStage || isCombinedStage) || cmdType == Type.PutSmallFile || cmdType == Type.PutBlock)) { + + if (!ContainerUtils.isContainerCreatable(msg)) { + StorageContainerException sce = new StorageContainerException( + "ContainerID " + containerID + " does not exist", + ContainerProtos.Result.CONTAINER_NOT_FOUND); + audit(action, eventType, msg, dispatcherContext, AuditEventStatus.FAILURE, sce); + return ContainerUtils.logAndReturnError(LOG, sce, msg); + } // If container does not exist, create one for WriteChunk and // PutSmallFile request responseProto = createContainer(msg); diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ec/reconstruction/ECReconstructionCoordinator.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ec/reconstruction/ECReconstructionCoordinator.java index 9b869f4a81f..426b1df6706 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ec/reconstruction/ECReconstructionCoordinator.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ec/reconstruction/ECReconstructionCoordinator.java @@ -232,7 +232,8 @@ private ECBlockOutputStream getECBlockOutputStream( containerOperationClient.singleNodePipeline(datanodeDetails, repConfig, replicaIndex), BufferPool.empty(), ozoneClientConfig, - blockLocationInfo.getToken(), clientMetrics, streamBufferArgs, ecReconstructWriteExecutor); + blockLocationInfo.getToken(), clientMetrics, streamBufferArgs, ecReconstructWriteExecutor, + false); } @VisibleForTesting diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java index 584cb98b367..2bb3738e9a7 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java @@ -681,6 +681,7 @@ ContainerCommandResponseProto handlePutBlock( request); } + updateRecoveringContainerTimeout(kvContainer); return putBlockResponseSuccess(request, blockDataProto); } @@ -1090,9 +1091,17 @@ ContainerCommandResponseProto handleWriteChunk( request); } + updateRecoveringContainerTimeout(kvContainer); return getWriteChunkResponseSuccess(request, blockDataProto); } + private void updateRecoveringContainerTimeout(KeyValueContainer kvContainer) { + if (kvContainer.getContainerState() != RECOVERING) { + return; + } + containerSet.updateRecoveringContainerTimeout(kvContainer.getContainerData().getContainerID()); + } + /** * Handle Write Chunk operation for closed container. Calls ChunkManager to process the request. */ diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/statemachine/background/StaleRecoveringContainerScrubbingService.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/statemachine/background/StaleRecoveringContainerScrubbingService.java index 9c535e5f6e9..6fd0bd59f8e 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/statemachine/background/StaleRecoveringContainerScrubbingService.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/statemachine/background/StaleRecoveringContainerScrubbingService.java @@ -18,13 +18,13 @@ package org.apache.hadoop.ozone.container.keyvalue.statemachine.background; import java.util.Iterator; +import java.util.Map; import java.util.concurrent.TimeUnit; import org.apache.hadoop.hdds.utils.BackgroundService; import org.apache.hadoop.hdds.utils.BackgroundTask; import org.apache.hadoop.hdds.utils.BackgroundTaskQueue; import org.apache.hadoop.hdds.utils.BackgroundTaskResult; import org.apache.hadoop.ozone.container.common.impl.ContainerSet; -import org.apache.hadoop.ozone.container.common.impl.ContainerSet.RecoveringContainer; import org.apache.hadoop.ozone.container.common.interfaces.Container; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -55,16 +55,15 @@ public BackgroundTaskQueue getTasks() { BackgroundTaskQueue backgroundTaskQueue = new BackgroundTaskQueue(); long currentTime = containerSet.getCurrentTime(); - Iterator<RecoveringContainer> it = - containerSet.getRecoveringContainerIterator(); + Iterator<Map.Entry<Long, Long>> it = containerSet.getRecoveringContainerMap().entrySet().iterator(); while (it.hasNext()) { - RecoveringContainer entry = it.next(); - if (currentTime >= entry.getTimeout()) { + Map.Entry<Long, Long> entry = it.next(); + long containerId = entry.getKey(); + long deadline = entry.getValue(); + if (currentTime >= deadline) { backgroundTaskQueue.add(new RecoveringContainerScrubbingTask( - containerSet, entry.getContainerId())); + containerSet, containerId)); it.remove(); - } else { - break; } } return backgroundTaskQueue; @@ -82,6 +81,12 @@ static class RecoveringContainerScrubbingTask implements BackgroundTask { @Override public BackgroundTaskResult call() throws Exception { + Long deadline = containerSet.getRecoveringContainerMap().get(containerID); + if (deadline != null && containerSet.getCurrentTime() < deadline) { + LOG.debug("Skipping stale recovering scrub for container {} - deadline extended", + containerID); + return new BackgroundTaskResult.EmptyTaskResult(); + } Container con = containerSet.getContainer(containerID); if (null != con) { con.markContainerUnhealthy(); diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/TestStaleRecoveringContainerScrubbingService.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/TestStaleRecoveringContainerScrubbingService.java index a21813956f5..3460c705c55 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/TestStaleRecoveringContainerScrubbingService.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/TestStaleRecoveringContainerScrubbingService.java @@ -22,6 +22,8 @@ import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerDataProto.State.UNHEALTHY; import static org.apache.hadoop.ozone.container.common.impl.ContainerImplTestUtils.newContainerSet; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.anyList; import static org.mockito.Mockito.anyLong; import static org.mockito.Mockito.mock; @@ -47,6 +49,8 @@ import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.scm.ScmConfigKeys; import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; +import org.apache.hadoop.hdds.utils.BackgroundTask; +import org.apache.hadoop.hdds.utils.BackgroundTaskQueue; import org.apache.hadoop.ozone.container.common.impl.ContainerLayoutVersion; import org.apache.hadoop.ozone.container.common.impl.ContainerSet; import org.apache.hadoop.ozone.container.common.interfaces.Container; @@ -187,4 +191,50 @@ public void testScrubbingStaleRecoveringContainers( containerStateMap.get(entry.getContainerData().getContainerID())); } } + + @ContainerTestVersionInfo.ContainerTest + public void testUpdateRecoveringContainerTimeoutExtendsScrubDeadline( + ContainerTestVersionInfo versionInfo) throws Exception { + initVersionInfo(versionInfo); + ContainerSet containerSet = newContainerSet(1000, testClock); + StaleRecoveringContainerScrubbingService srcss = + new StaleRecoveringContainerScrubbingService( + 50, TimeUnit.MILLISECONDS, 10, + Duration.ofSeconds(300).toMillis(), + containerSet); + List<Long> ids = createTestContainers(containerSet, 1, RECOVERING); + long containerId = ids.get(0); + testClock.fastForward(800L); + containerSet.updateRecoveringContainerTimeout(containerId); + testClock.fastForward(800L); + srcss.runPeriodicalTaskNow(); + assertEquals(RECOVERING, containerSet.getContainer(containerId).getContainerState()); + testClock.fastForward(500L); + srcss.runPeriodicalTaskNow(); + assertEquals(UNHEALTHY, containerSet.getContainer(containerId).getContainerState()); + } + + @ContainerTestVersionInfo.ContainerTest + public void testScrubSkippedWhenDeadlineExtendedBeforeTaskRuns( + ContainerTestVersionInfo versionInfo) throws Exception { + initVersionInfo(versionInfo); + ContainerSet containerSet = newContainerSet(1000, testClock); + StaleRecoveringContainerScrubbingService srcss = + new StaleRecoveringContainerScrubbingService( + 50, TimeUnit.MILLISECONDS, 10, + Duration.ofSeconds(300).toMillis(), + containerSet); + List<Long> ids = createTestContainers(containerSet, 1, RECOVERING); + long containerId = ids.get(0); + testClock.fastForward(1000L); + BackgroundTaskQueue tasks = srcss.getTasks(); + assertFalse(containerSet.getRecoveringContainerMap().containsKey(containerId)); + containerSet.updateRecoveringContainerTimeout(containerId); + while (!tasks.isEmpty()) { + BackgroundTask task = tasks.poll(); + task.call(); + } + assertEquals(RECOVERING, containerSet.getContainer(containerId).getContainerState()); + assertTrue(containerSet.getRecoveringContainerMap().containsKey(containerId)); + } } diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java index 8c664321c46..d60ca220ece 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java @@ -26,6 +26,8 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.any; @@ -339,8 +341,7 @@ public void testContainerCloseActionWhenVolumeFull( HddsDispatcher hddsDispatcher = new HddsDispatcher( conf, containerSet, volumeSet, handlers, context, metrics, null); hddsDispatcher.setClusterId(scmId.toString()); - containerData.getVolume().getVolumeUsage() - .ifPresent(usage -> usage.incrementUsedSpace(60)); + containerData.getVolume().incrementUsedSpace(60); usedSpace.addAndGet(60); ContainerCommandResponseProto response = hddsDispatcher .dispatch(getWriteChunkRequest(dd.getUuidString(), 1L, 1L), null); @@ -496,6 +497,90 @@ public void testWriteChunkWithCreateContainerFailure() throws IOException { } } + @Test + public void testCreateContainerWhenAlreadyExistsDoesNotMarkUnhealthy() throws IOException { + String testDirPath = testDir.getPath(); + try { + UUID scmId = UUID.randomUUID(); + OzoneConfiguration conf = new OzoneConfiguration(); + conf.set(HDDS_DATANODE_DIR_KEY, testDirPath); + conf.set(OzoneConfigKeys.OZONE_METADATA_DIRS, testDirPath); + DatanodeDetails dd = randomDatanodeDetails(); + HddsDispatcher hddsDispatcher = createDispatcher(dd, scmId, conf); + + // Create container via WriteChunk + ContainerCommandRequestProto writeChunkRequest = + getWriteChunkRequest(dd.getUuidString(), 1L, 1L); + ContainerCommandResponseProto initialResponse = + hddsDispatcher.dispatch(writeChunkRequest, null); + assertEquals(ContainerProtos.Result.SUCCESS, initialResponse.getResult()); + + // Send direct CreateContainer for existing container + ContainerCommandRequestProto createRequest = + ContainerCommandRequestProto.newBuilder() + .setCmdType(ContainerProtos.Type.CreateContainer) + .setContainerID(1L) + .setCreateContainer(ContainerProtos.CreateContainerRequestProto.newBuilder() + .setContainerType(ContainerProtos.ContainerType.KeyValueContainer) + .build()) + .setDatanodeUuid(dd.getUuidString()) + .build(); + + ContainerCommandResponseProto response = + hddsDispatcher.dispatch(createRequest, null); + assertEquals(ContainerProtos.Result.CONTAINER_ALREADY_EXISTS, response.getResult()); + + Container container = hddsDispatcher.getContainer(1L); + assertNotNull(container); + assertTrue(container.getContainerData().isOpen()); + assertFalse(container.getContainerData().isUnhealthy()); + } finally { + ContainerMetrics.remove(); + } + } + + @Test + public void testMalformedPutBlockDoesNotMarkContainerUnhealthy() throws IOException { + String testDirPath = testDir.getPath(); + try { + UUID scmId = UUID.randomUUID(); + OzoneConfiguration conf = new OzoneConfiguration(); + conf.set(HDDS_DATANODE_DIR_KEY, testDirPath); + conf.set(OzoneConfigKeys.OZONE_METADATA_DIRS, testDirPath); + DatanodeDetails dd = randomDatanodeDetails(); + HddsDispatcher hddsDispatcher = createDispatcher(dd, scmId, conf); + + ContainerCommandRequestProto writeChunkRequest = + getWriteChunkRequest(dd.getUuidString(), 1L, 1L); + ContainerCommandResponseProto writeChunkResponse = + hddsDispatcher.dispatch(writeChunkRequest, null); + assertEquals(ContainerProtos.Result.SUCCESS, writeChunkResponse.getResult()); + + ContainerCommandRequestProto putBlockRequest = + ContainerTestHelper.getPutBlockRequest(writeChunkRequest); + ContainerProtos.BlockData malformedBlockData = + putBlockRequest.getPutBlock().getBlockData().toBuilder() + .setSize(putBlockRequest.getPutBlock().getBlockData().getSize() + 1) + .build(); + ContainerCommandRequestProto malformedPutBlockRequest = + putBlockRequest.toBuilder() + .setPutBlock(putBlockRequest.getPutBlock().toBuilder() + .setBlockData(malformedBlockData)) + .build(); + + ContainerCommandResponseProto response = + hddsDispatcher.dispatch(malformedPutBlockRequest, null); + assertEquals(ContainerProtos.Result.MALFORMED_REQUEST, response.getResult()); + + Container container = hddsDispatcher.getContainer(1L); + assertNotNull(container); + assertTrue(container.getContainerData().isOpen()); + assertFalse(container.getContainerData().isUnhealthy()); + } finally { + ContainerMetrics.remove(); + } + } + @Test public void testDuplicateWriteChunkAndPutBlockRequest() throws IOException { String testDirPath = testDir.getPath(); @@ -645,6 +730,43 @@ private ContainerCommandRequestProto getWriteChunkRequest( .build(); } + private static ContainerCommandRequestProto withCreatableFalse( + ContainerCommandRequestProto writeChunk) { + return ContainerCommandRequestProto.newBuilder(writeChunk) + .setWriteChunk(writeChunk.getWriteChunk().toBuilder() + .setContainerAutoCreate(false) + .build()) + .build(); + } + + private static ContainerCommandRequestProto getEmptyPutBlockRequest( + String datanodeId, Long containerId, Long localId) { + BlockID blockID = new BlockID(containerId, localId); + ContainerProtos.BlockData blockData = ContainerProtos.BlockData.newBuilder() + .setBlockID(blockID.getDatanodeBlockIDProtobuf()) + .build(); + ContainerProtos.PutBlockRequestProto putBlockRequest = + ContainerProtos.PutBlockRequestProto.newBuilder() + .setBlockData(blockData) + .setEof(true) + .build(); + return ContainerCommandRequestProto.newBuilder() + .setContainerID(containerId) + .setCmdType(ContainerProtos.Type.PutBlock) + .setDatanodeUuid(datanodeId) + .setPutBlock(putBlockRequest) + .build(); + } + + private static ContainerCommandRequestProto withCreatableFalsePutBlock( + ContainerCommandRequestProto putBlock) { + return ContainerCommandRequestProto.newBuilder(putBlock) + .setPutBlock(putBlock.getPutBlock().toBuilder() + .setContainerAutoCreate(false) + .build()) + .build(); + } + static ChecksumData checksum(ByteString data) { try { return new Checksum(ContainerProtos.ChecksumType.CRC32, 256) @@ -814,6 +936,152 @@ public void verify(Token<?> token, } } + /** + * Verifies the soft/hard min-free-space split on the write path: + * + * <p>Setup (capacity=500 bytes): + * <pre> + * minFreeSpace bytes floor = 1 (ratio always dominates) + * softRatio = 10% → softSpare = 50 bytes (reported to SCM) + * hardRatio = 6% → hardSpare = 30 bytes (local write enforcement) + * softBand = 20 bytes + * writeChunk size ≈ 36 bytes (UUID string) + * </pre> + * + * <p>Three scenarios exercised in sequence using the same volume by calling + * {@code hddsVolume.incrementUsedSpace(delta)} to update the CachingSpaceUsageSource cache: + * <ol> + * <li>Well above both limits (usedSpace=400, available=100): write passes, no metric fires.</li> + * <li>Inside the soft band (usedSpace=425, available=75): write passes (75-30=45 > 36), + * {@code numWriteRequestsInSoftBandMinFreeSpace} incremented (75-50=25 < 36).</li> + * <li>Below hard limit (usedSpace=465, available=35): write rejected with DISK_OUT_OF_SPACE + * (35-30=5 < 36), {@code numWriteRequestsRejectedHardMinFreeSpace} incremented.</li> + * </ol> + */ + @ContainerLayoutTestInfo.ContainerTest + public void testWriteChunkEnforcesSoftHardMinFreeSpace( + ContainerLayoutVersion layoutVersion) throws Exception { + String testDirPath = testDir.getPath(); + OzoneConfiguration conf = new OzoneConfiguration(); + // 1-byte floor so the percentage ratios always dominate + conf.setStorageSize(DatanodeConfiguration.HDDS_DATANODE_VOLUME_MIN_FREE_SPACE, + 1.0, StorageUnit.BYTES); + // soft spare = 10% of 500 = 50 bytes; hard spare = 6% of 500 = 30 bytes; band = 20 bytes + conf.setFloat(DatanodeConfiguration.HDDS_DATANODE_VOLUME_MIN_FREE_SPACE_PERCENT, 0.1f); + conf.setFloat(DatanodeConfiguration.HDDS_DATANODE_VOLUME_MIN_FREE_SPACE_HARD_LIMIT_PERCENT, 0.06f); + conf.set(HDDS_DATANODE_DIR_KEY, testDirPath); + conf.set(OzoneConfigKeys.OZONE_METADATA_DIRS, testDirPath); + DatanodeDetails dd = randomDatanodeDetails(); + UUID scmId = UUID.randomUUID(); + AtomicLong usedSpace = new AtomicLong(400); // available = 100, well above both limits + SpaceUsageSource spaceUsage = MockSpaceUsageSource.of(500, usedSpace); + SpaceUsageCheckFactory factory = MockSpaceUsageCheckFactory.of( + spaceUsage, Duration.ZERO, inMemory(new AtomicLong(0))); + HddsVolume.Builder volumeBuilder = + new HddsVolume.Builder(testDirPath).datanodeUuid(dd.getUuidString()) + .conf(conf).usageCheckFactory(MockSpaceUsageCheckFactory.NONE).clusterID("test"); + volumeBuilder.usageCheckFactory(factory); + MutableVolumeSet volumeSet = mock(MutableVolumeSet.class); + when(volumeSet.getVolumesList()) + .thenReturn(Collections.singletonList(volumeBuilder.build())); + volumeSet.getVolumesList().get(0).setState(StorageVolume.VolumeState.NORMAL); + volumeSet.getVolumesList().get(0).start(); + HddsVolume hddsVolume = StorageVolumeUtil + .getHddsVolumesList(volumeSet.getVolumesList()).get(0); + try { + KeyValueContainerData containerData = new KeyValueContainerData(1L, + layoutVersion, 50, UUID.randomUUID().toString(), dd.getUuidString()); + Container container = new KeyValueContainer(containerData, conf); + StorageVolumeUtil.getHddsVolumesList(volumeSet.getVolumesList()) + .forEach(v -> v.setDbParentDir(tempDir.toFile())); + container.create(volumeSet, new RoundRobinVolumeChoosingPolicy(), scmId.toString()); + ContainerSet containerSet = newContainerSet(); + containerSet.addContainer(container); + StateContext context = ContainerTestUtils.getMockContext(dd, conf); + ContainerMetrics metrics = ContainerMetrics.create(conf); + Map<ContainerType, Handler> handlers = Maps.newHashMap(); + for (ContainerType containerType : ContainerType.values()) { + handlers.put(containerType, + Handler.getHandlerForContainerType(containerType, conf, + dd.getUuidString(), containerSet, volumeSet, volumeChoosingPolicy, + metrics, NO_OP_ICR_SENDER, new ContainerChecksumTreeManager(conf))); + } + HddsDispatcher hddsDispatcher = new HddsDispatcher( + conf, containerSet, volumeSet, handlers, context, metrics, null); + hddsDispatcher.setClusterId(scmId.toString()); + // --- Scenario 1: well above both limits (available=100) --- + // available(100) - hardSpare(30) = 70 > writeSize(~36): passes + // available(100) - softSpare(50) = 50 > writeSize(~36): not in soft band + ContainerCommandResponseProto response = + hddsDispatcher.dispatch(getWriteChunkRequest(dd.getUuidString(), 1L, 1L), null); + assertEquals(ContainerProtos.Result.SUCCESS, response.getResult()); + assertEquals(0, + hddsVolume.getVolumeInfoStats().getNumWriteRequestsInSoftBandMinFreeSpace()); + assertEquals(0, + hddsVolume.getVolumeInfoStats().getNumWriteRequestsRejectedHardMinFreeSpace()); + // --- Scenario 2: inside the soft band (usedSpace → 425, available=75) --- + // available(75) - hardSpare(30) = 45 > writeSize(~36): passes hard check + // available(75) - softSpare(50) = 25 < writeSize(~36): soft-band metric fires + // Use incrementUsedSpace so the CachingSpaceUsageSource internal cache is updated; + hddsVolume.incrementUsedSpace(25); // 400 → 425 + response = hddsDispatcher.dispatch(getWriteChunkRequest(dd.getUuidString(), 1L, 2L), null); + assertEquals(ContainerProtos.Result.SUCCESS, response.getResult()); + assertEquals(1, + hddsVolume.getVolumeInfoStats().getNumWriteRequestsInSoftBandMinFreeSpace()); + assertEquals(0, + hddsVolume.getVolumeInfoStats().getNumWriteRequestsRejectedHardMinFreeSpace()); + // --- Scenario 3: below hard limit (usedSpace → 465, available=35) --- + // available(35) - hardSpare(30) = 5 < writeSize(~36): DISK_OUT_OF_SPACE + hddsVolume.incrementUsedSpace(40); // 425 → 465 + response = hddsDispatcher.dispatch(getWriteChunkRequest(dd.getUuidString(), 1L, 3L), null); + assertEquals(ContainerProtos.Result.DISK_OUT_OF_SPACE, response.getResult()); + assertEquals(1, + hddsVolume.getVolumeInfoStats().getNumWriteRequestsInSoftBandMinFreeSpace()); + assertEquals(1, + hddsVolume.getVolumeInfoStats().getNumWriteRequestsRejectedHardMinFreeSpace()); + } finally { + volumeSet.shutdown(); + ContainerMetrics.remove(); + } + } + + @Test + public void testEcReconstructionWriteChunkDeniedWhenContainerCreatableFalse() + throws IOException { + String testDirPath = testDir.getPath(); + UUID scmId = UUID.randomUUID(); + OzoneConfiguration conf = new OzoneConfiguration(); + conf.set(HDDS_DATANODE_DIR_KEY, testDirPath); + conf.set(OzoneConfigKeys.OZONE_METADATA_DIRS, testDirPath); + DatanodeDetails dd = randomDatanodeDetails(); + HddsDispatcher dispatcher = createDispatcher(dd, scmId, conf); + long containerId = 99L; + + ContainerCommandResponseProto response = dispatcher.dispatch( + withCreatableFalse(getWriteChunkRequest(dd.getUuidString(), containerId, 1L)), null); + assertEquals(ContainerProtos.Result.CONTAINER_NOT_FOUND, response.getResult()); + assertNull(dispatcher.getContainer(containerId)); + } + + @Test + public void testEcReconstructionPutBlockDeniedWhenContainerCreatableFalse() + throws IOException { + String testDirPath = testDir.getPath(); + UUID scmId = UUID.randomUUID(); + OzoneConfiguration conf = new OzoneConfiguration(); + conf.set(HDDS_DATANODE_DIR_KEY, testDirPath); + conf.set(OzoneConfigKeys.OZONE_METADATA_DIRS, testDirPath); + DatanodeDetails dd = randomDatanodeDetails(); + HddsDispatcher dispatcher = createDispatcher(dd, scmId, conf); + long containerId = 100L; + + ContainerCommandResponseProto response = dispatcher.dispatch( + withCreatableFalsePutBlock(getEmptyPutBlockRequest(dd.getUuidString(), containerId, 1L)), + null); + assertEquals(ContainerProtos.Result.CONTAINER_NOT_FOUND, response.getResult()); + assertNull(dispatcher.getContainer(containerId)); + } + static DispatcherContext newContext(Op op) { return newContext(op, WriteChunkStage.COMBINED); } diff --git a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto index bd890eae64a..5144865e96b 100644 --- a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto +++ b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto @@ -344,6 +344,7 @@ message BlockData { message PutBlockRequestProto { required BlockData blockData = 1; optional bool eof = 2; + optional bool containerAutoCreate = 3; } message PutBlockResponseProto { @@ -446,6 +447,7 @@ message WriteChunkRequestProto { optional ChunkInfo chunkData = 2; optional bytes data = 3; optional PutBlockRequestProto block = 4; + optional bool containerAutoCreate = 5; } message WriteChunkResponseProto { --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
