This is an automated email from the ASF dual-hosted git repository.
aswinshakil pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/master by this push:
new 458c0d14498 HDDS-15791. EC reconstructed RECOVERING container while
still idle can be deleted by SCM and recreated as new partial OPEN container
(#10702)
458c0d14498 is described below
commit 458c0d14498231b558f131fe0de26a7754a4cea3
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)
---
.../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 05ddf4b9e4f..5329428a348 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
@@ -318,8 +318,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);
}
@@ -356,10 +366,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)
@@ -467,6 +486,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()
@@ -485,6 +515,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 1b3f399b46f..7b8f5c8a7c8 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 fe4204b159f..e2881bd6f7e 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 e33b07b50aa..0277638f11b 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]