This is an automated email from the ASF dual-hosted git repository. jojochuang pushed a commit to branch ozone-2.2 in repository https://gitbox.apache.org/repos/asf/ozone.git
commit 778197970ec1641970b6d496c63b5383a9c7a080 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) --- .../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 | 75 +++++++++++++++++ .../src/main/proto/DatanodeClientProtocol.proto | 2 + 13 files changed, 324 insertions(+), 93 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 d7076df3ba0..b960e753744 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 @@ -597,7 +597,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 { @@ -964,7 +964,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 -> { @@ -1160,6 +1160,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 15879fb4764..70a362f0dda 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 @@ -289,8 +289,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); } @@ -327,10 +337,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) @@ -438,6 +457,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() @@ -456,6 +486,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 f2817e37b51..4cd68a2b034 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 @@ -403,4 +403,30 @@ public static long getPendingDeletionBytes(ContainerData containerData) { " not support."); } } + + /** + * @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 43a3fc0a00b..752fdaf9fd1 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; @@ -74,8 +75,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 @@ -209,9 +209,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); } HddsVolume volume = container.getContainerData().getVolume(); if (volume != null) { @@ -422,22 +420,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; } /** @@ -489,15 +485,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. @@ -669,52 +656,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 249e58df94d..a363d99b469 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 @@ -298,6 +298,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 7e73fdd76ee..3d5369a939c 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 @@ -231,7 +231,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 af7242541ab..3ccadc0736e 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 @@ -697,6 +697,7 @@ ContainerCommandResponseProto handlePutBlock( request); } + updateRecoveringContainerTimeout(kvContainer); return putBlockResponseSuccess(request, blockDataProto); } @@ -1105,9 +1106,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 0f10eaf1327..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 @@ -27,6 +27,7 @@ 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; @@ -729,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) @@ -1007,6 +1045,43 @@ public void testWriteChunkEnforcesSoftHardMinFreeSpace( } } + @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 05c94624c99..595a17e3022 100644 --- a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto +++ b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto @@ -326,6 +326,7 @@ message BlockData { message PutBlockRequestProto { required BlockData blockData = 1; optional bool eof = 2; + optional bool containerAutoCreate = 3; } message PutBlockResponseProto { @@ -441,6 +442,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]
