This is an automated email from the ASF dual-hosted git repository.
ashishkumar50 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 0dd906c4f15 HDDS-15541. PlacementPolicy to use PendingContainerTracker
to check space availability. (#10497)
0dd906c4f15 is described below
commit 0dd906c4f1527aad4c53b7ca86727d724bb88111
Author: Ashish Kumar <[email protected]>
AuthorDate: Thu Jun 25 22:10:38 2026 +0530
HDDS-15541. PlacementPolicy to use PendingContainerTracker to check space
availability. (#10497)
---
.../hadoop/hdds/scm/SCMCommonPlacementPolicy.java | 54 +++++++++---------
.../apache/hadoop/hdds/scm/node/NodeManager.java | 9 +++
.../hdds/scm/node/PendingContainerTracker.java | 54 ++++++++++++++++--
.../hadoop/hdds/scm/node/SCMNodeManager.java | 43 ++++++++++----
.../hadoop/hdds/scm/node/SCMNodeMetrics.java | 10 ++++
.../hdds/scm/pipeline/PipelinePlacementPolicy.java | 2 +-
.../hadoop/hdds/scm/pipeline/PipelineProvider.java | 3 +-
.../hdds/scm/TestSCMCommonPlacementPolicy.java | 37 ++----------
.../hadoop/hdds/scm/container/MockNodeManager.java | 10 ++++
.../hdds/scm/container/SimpleMockNodeManager.java | 5 ++
.../algorithms/TestContainerPlacementFactory.java | 2 +
.../TestSCMContainerPlacementCapacity.java | 4 ++
.../TestSCMContainerPlacementRackAware.java | 2 +
.../TestSCMContainerPlacementRackScatter.java | 5 ++
.../TestSCMContainerPlacementRandom.java | 9 +++
.../hdds/scm/node/TestPendingContainerTracker.java | 65 ++++++++++++++++++++++
.../scm/pipeline/TestPipelinePlacementFactory.java | 3 +
.../scm/pipeline/TestPipelinePlacementPolicy.java | 1 +
.../scm/pipeline/TestRatisPipelineProvider.java | 5 ++
19 files changed, 245 insertions(+), 78 deletions(-)
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/SCMCommonPlacementPolicy.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/SCMCommonPlacementPolicy.java
index 93940b7770f..0deece79c6d 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/SCMCommonPlacementPolicy.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/SCMCommonPlacementPolicy.java
@@ -37,7 +37,6 @@
import org.apache.hadoop.hdds.conf.ConfigurationSource;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.MetadataStorageReportProto;
-import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.StorageReportProto;
import org.apache.hadoop.hdds.scm.container.ContainerReplica;
import
org.apache.hadoop.hdds.scm.container.placement.algorithms.ContainerPlacementStatusDefault;
import org.apache.hadoop.hdds.scm.exceptions.SCMException;
@@ -47,7 +46,6 @@
import org.apache.hadoop.hdds.scm.node.NodeManager;
import org.apache.hadoop.hdds.scm.node.NodeStatus;
import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException;
-import org.apache.hadoop.ozone.container.common.volume.VolumeUsage;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -252,7 +250,7 @@ protected List<DatanodeDetails> chooseDatanodesInternal(
}
return filterNodesWithSpace(healthyNodes, nodesRequired,
- metadataSizeRequired, dataSizeRequired);
+ metadataSizeRequired);
}
/**
@@ -270,18 +268,17 @@ protected boolean usedNodesPassed(List<DatanodeDetails>
list) {
}
public List<DatanodeDetails> filterNodesWithSpace(List<DatanodeDetails>
nodes,
- int nodesRequired, long metadataSizeRequired, long dataSizeRequired)
+ int nodesRequired, long metadataSizeRequired)
throws SCMException {
List<DatanodeDetails> nodesWithSpace = nodes.stream().filter(d ->
- hasEnoughSpace(d, metadataSizeRequired, dataSizeRequired))
+ hasEnoughSpace(d, metadataSizeRequired, nodeManager))
.collect(Collectors.toList());
if (nodesWithSpace.size() < nodesRequired) {
String msg = String.format("Unable to find enough nodes that meet the " +
- "space requirement of %d bytes for metadata and %d bytes for " +
- "data in healthy node set. Required %d. Found %d.",
- metadataSizeRequired, dataSizeRequired, nodesRequired,
- nodesWithSpace.size());
+ "space requirement of %d bytes for metadata and an available " +
+ "container slot in healthy node set. Required %d. Found %d.",
+ metadataSizeRequired, nodesRequired, nodesWithSpace.size());
LOG.warn(msg);
throw new SCMException(msg,
SCMException.ResultCodes.FAILED_TO_FIND_NODES_WITH_SPACE);
@@ -291,35 +288,33 @@ public List<DatanodeDetails>
filterNodesWithSpace(List<DatanodeDetails> nodes,
}
/**
- * Returns true if this node has enough space to meet our requirement.
+ * Returns true if this node has enough space to satisfy the placement
request.
+ *
+ * <p>Data-space is checked via {@link NodeManager#hasAvailableSpace}, which
+ * delegates to {@link
org.apache.hadoop.hdds.scm.node.PendingContainerTracker}
+ * and accounts for both current disk usage and in-flight allocations.
+ * The check always uses {@code maxContainerSize} as the unit of allocation,
+ * regardless of the actual container's used bytes.
*
- * @param datanodeDetails DatanodeDetails
- * @return true if we have enough space.
+ * @param datanodeDetails the datanode to evaluate
+ * @param metadataSizeRequired minimum metadata volume space required in
bytes
+ * @param nodeManager used to check slot availability via
PendingContainerTracker
+ * @return true if the datanode has both an available data slot and enough
metadata space
*/
public static boolean hasEnoughSpace(DatanodeDetails datanodeDetails,
long metadataSizeRequired,
- long dataSizeRequired) {
+ NodeManager nodeManager) {
Preconditions.checkArgument(datanodeDetails instanceof DatanodeInfo);
- boolean enoughForData = false;
boolean enoughForMeta = false;
DatanodeInfo datanodeInfo = (DatanodeInfo) datanodeDetails;
- if (dataSizeRequired > 0) {
- for (StorageReportProto reportProto : datanodeInfo.getStorageReports()) {
- if (VolumeUsage.getUsableSpace(reportProto) > dataSizeRequired) {
- enoughForData = true;
- break;
- }
- }
- } else {
- enoughForData = true;
- }
-
- if (!enoughForData) {
- LOG.debug("Datanode {} has no volumes with enough space to allocate {} "
+
- "bytes for data.", datanodeDetails, dataSizeRequired);
+ // Data-space check: use PendingContainerTracker slot availability.
+ // This accounts for both current disk usage and in-flight allocations.
+ // Always slot-based (maxContainerSize unit).
+ if (!nodeManager.hasAvailableSpace(datanodeInfo)) {
+ LOG.debug("Datanode {} has no available container slots.",
datanodeDetails);
return false;
}
@@ -539,7 +534,8 @@ public boolean isValidNode(DatanodeDetails datanodeDetails,
return false;
}
NodeStatus nodeStatus = datanodeInfo.getNodeStatus();
- if (nodeStatus.isNodeWritable() && (hasEnoughSpace(datanodeInfo,
metadataSizeRequired, dataSizeRequired))) {
+ if (nodeStatus.isNodeWritable() && (hasEnoughSpace(datanodeInfo,
metadataSizeRequired,
+ nodeManager))) {
LOG.debug("Datanode {} is chosen. Required metadata size is {} and " +
"required data size is {} and NodeStatus is {}",
datanodeDetails, metadataSizeRequired, dataSizeRequired, nodeStatus);
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
index 6a2d92b64a5..3da84aad3a7 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
@@ -195,6 +195,15 @@ default int getAllNodeCount() {
*/
boolean checkSpaceAndRecordAllocation(DatanodeInfo datanodeInfo, ContainerID
containerID);
+ /**
+ * Returns true if the datanode has at least one available container slot
considering
+ * in-flight allocations tracked by PendingContainerTracker.
+ *
+ * @param datanodeInfo the datanode to check
+ * @return true if at least one slot is free
+ */
+ boolean hasAvailableSpace(DatanodeInfo datanodeInfo);
+
/**
* Removes a pending container allocation from a datanode.
*
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/PendingContainerTracker.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/PendingContainerTracker.java
index 3821727ed97..229dc5fad92 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/PendingContainerTracker.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/PendingContainerTracker.java
@@ -105,7 +105,11 @@ synchronized void rollIfNeeded() {
previousWindow.clear();
currentWindow.clear();
lastRollTime = now;
- LOG.debug("Double roll interval elapsed ({}ms): dropped {} pending
containers", elapsed, dropped);
+ if (dropped > 0) {
+ LOG.warn("PendingContainerTracker: force-dropped {} unconfirmed
pending containers "
+ + "on DN {} after {}ms (2x rollInterval). "
+ + "Container reports may have been lost.", dropped, datanodeID,
elapsed);
+ }
} else if (elapsed >= rollIntervalMs) {
previousWindow.clear();
final Set<ContainerID> tmp = previousWindow;
@@ -156,7 +160,11 @@ synchronized boolean checkSpaceAndAdd(
final int pendingAllocationCount = getCount();
long allocatableCount = 0;
for (StorageReportProto report : storageReports) {
- final long allocatableCountOnThisDisk =
VolumeUsage.getUsableSpace(report) / maxContainerSize;
+ if (report.hasFailed() && report.getFailed()) {
+ continue;
+ }
+ final long allocatableCountOnThisDisk =
+ Math.max(0L, VolumeUsage.getUsableSpace(report)) /
maxContainerSize;
allocatableCount += allocatableCountOnThisDisk;
if (allocatableCount > pendingAllocationCount) {
final boolean added = currentWindow.add(containerID);
@@ -208,11 +216,45 @@ public boolean checkSpaceAndRecordAllocation(DatanodeInfo
datanodeInfo, Containe
}
/**
- * Remove a pending container allocation from a specific DataNode.
- * Removes from both current and previous windows.
- * Called when container is confirmed.
+ * Returns true if the given datanode has at least one allocatable container
slot
+ * available, accounting for pending in-flight allocations.
+ *
+ * <p>Slot availability is based on {@code maxContainerSize}: a slot exists
for each
+ * {@code maxContainerSize}-worth of usable space on any volume. This check
is intended for the placement policy.
+ * This rolls expired-window entries but does not consume a slot.
+ *
+ * @param datanodeInfo the datanode to check
+ * @return true if at least one container slot is available
+ */
+ public boolean hasAvailableSpace(DatanodeInfo datanodeInfo) {
+ Objects.requireNonNull(datanodeInfo, "datanodeInfo == null");
+ List<StorageReportProto> storageReports = datanodeInfo.getStorageReports();
+ if (storageReports.isEmpty()) {
+ return false;
+ }
+ TwoWindowBucket bucket = datanodeInfo.getPendingContainerAllocations();
+ bucket.rollIfNeeded();
+ final int pendingCount = bucket.getCount();
+ long allocatableCount = 0;
+ for (StorageReportProto report : storageReports) {
+ if (report.hasFailed() && report.getFailed()) {
+ continue;
+ }
+ allocatableCount += Math.max(0L, VolumeUsage.getUsableSpace(report)) /
maxContainerSize;
+ if (allocatableCount > pendingCount) {
+ return true;
+ }
+ }
+ LOG.debug("Datanode {} has no available container slots. Pending: {},
Allocatable: {}",
+ datanodeInfo.getID(), pendingCount, allocatableCount);
+ return false;
+ }
+
+ /**
+ * Remove pending allocation from the bucket for the given container.
*
- * @param containerID The container to remove from pending
+ * @param bucket TWO window bucket of the datanode
+ * @param containerID containerID
*/
public void removePendingAllocation(TwoWindowBucket bucket, ContainerID
containerID) {
Objects.requireNonNull(containerID, "containerID == null");
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
index 926884c5b61..66a41fe773d 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
@@ -22,7 +22,6 @@
import static
org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.IN_SERVICE;
import static
org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeState.HEALTHY;
import static
org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeState.HEALTHY_READONLY;
-import static
org.apache.hadoop.hdds.scm.SCMCommonPlacementPolicy.hasEnoughSpace;
import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Preconditions;
@@ -63,6 +62,7 @@
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.CommandQueueReportProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.LayoutVersionProto;
+import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.MetadataStorageReportProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.NodeReportProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.PipelineReportsProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto;
@@ -213,7 +213,7 @@ public SCMNodeManager(
ScmConfigKeys.OZONE_SCM_PIPELINE_OWNER_CONTAINER_COUNT_DEFAULT);
this.scmContext = scmContext;
this.sendCommandNotifyMap = new HashMap<>();
- this.nonWritableNodeFilter = new NonWritableNodeFilter(conf);
+ this.nonWritableNodeFilter = new NonWritableNodeFilter(conf,
pendingContainerTracker);
}
@Override
@@ -1084,6 +1084,11 @@ public boolean
checkSpaceAndRecordAllocation(DatanodeInfo datanodeInfo, Containe
return pendingContainerTracker.checkSpaceAndRecordAllocation(datanodeInfo,
containerID);
}
+ @Override
+ public boolean hasAvailableSpace(DatanodeInfo datanodeInfo) {
+ return pendingContainerTracker.hasAvailableSpace(datanodeInfo);
+ }
+
@Override
public void removePendingAllocationForDatanode(DatanodeInfo datanodeInfo,
ContainerID containerID) {
pendingContainerTracker.removePendingAllocation(
@@ -1430,7 +1435,9 @@ private void nodeSpaceStatistics(Map<String, String>
nodeStatics) {
long capacityByte = 0;
long scmUsedByte = 0;
long remainingByte = 0;
+ long totalPending = 0;
for (DatanodeInfo dni : nodeStateManager.getAllNodes()) {
+ totalPending += dni.getPendingContainerAllocations().getCount();
List<StorageReportProto> storageReports = dni.getStorageReports();
if (storageReports != null && !storageReports.isEmpty()) {
for (StorageReportProto storageReport : storageReports) {
@@ -1440,6 +1447,7 @@ private void nodeSpaceStatistics(Map<String, String>
nodeStatics) {
}
}
}
+ metrics.setTotalPendingContainerSlots(totalPending);
long nonScmUsedByte = capacityByte - scmUsedByte - remainingByte;
if (nonScmUsedByte < 0) {
@@ -1467,9 +1475,9 @@ static class NonWritableNodeFilter implements
Predicate<DatanodeInfo> {
private final long blockSize;
private final long minRatisVolumeSizeBytes;
- private final long containerSize;
+ private final PendingContainerTracker tracker;
- NonWritableNodeFilter(ConfigurationSource conf) {
+ NonWritableNodeFilter(ConfigurationSource conf, PendingContainerTracker
tracker) {
blockSize = (long) conf.getStorageSize(
OzoneConfigKeys.OZONE_SCM_BLOCK_SIZE,
OzoneConfigKeys.OZONE_SCM_BLOCK_SIZE_DEFAULT,
@@ -1478,17 +1486,32 @@ static class NonWritableNodeFilter implements
Predicate<DatanodeInfo> {
ScmConfigKeys.OZONE_DATANODE_RATIS_VOLUME_FREE_SPACE_MIN,
ScmConfigKeys.OZONE_DATANODE_RATIS_VOLUME_FREE_SPACE_MIN_DEFAULT,
StorageUnit.BYTES);
- containerSize = (long) conf.getStorageSize(
- ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE,
- ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE_DEFAULT,
- StorageUnit.BYTES);
+ this.tracker = tracker;
}
@Override
public boolean test(DatanodeInfo dn) {
return !dn.getNodeStatus().isNodeWritable()
- || (!hasEnoughSpace(dn, minRatisVolumeSizeBytes, containerSize)
- && !hasEnoughCommittedVolumeSpace(dn));
+ || (!hasEnoughSpaceForNode(dn) &&
!hasEnoughCommittedVolumeSpace(dn));
+ }
+
+ /**
+ * Returns true if the datanode has both an available data slot (via
+ * {@link PendingContainerTracker}) and sufficient Ratis metadata volume
space.
+ */
+ private boolean hasEnoughSpaceForNode(DatanodeInfo dn) {
+ if (!tracker.hasAvailableSpace(dn)) {
+ return false;
+ }
+ if (minRatisVolumeSizeBytes <= 0) {
+ return true;
+ }
+ for (MetadataStorageReportProto report : dn.getMetadataStorageReports())
{
+ if (report.getRemaining() > minRatisVolumeSizeBytes) {
+ return true;
+ }
+ }
+ return false;
}
/**
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeMetrics.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeMetrics.java
index 0b90247ae8a..ee78ea3dd2f 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeMetrics.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeMetrics.java
@@ -30,6 +30,7 @@
import org.apache.hadoop.metrics2.lib.Interns;
import org.apache.hadoop.metrics2.lib.MetricsRegistry;
import org.apache.hadoop.metrics2.lib.MutableCounterLong;
+import org.apache.hadoop.metrics2.lib.MutableGaugeLong;
import org.apache.hadoop.ozone.OzoneConsts;
import org.apache.hadoop.util.StringUtils;
@@ -54,6 +55,7 @@ public final class SCMNodeMetrics implements MetricsSource {
private @Metric MutableCounterLong numPendingContainersAdded;
private @Metric MutableCounterLong numPendingContainersRemoved;
private @Metric MutableCounterLong numSkippedFullNodeContainerAllocation;
+ private @Metric MutableGaugeLong totalPendingContainerSlots;
private final MetricsRegistry registry;
private final NodeManagerMXBean managerMXBean;
@@ -148,6 +150,14 @@ void incNumSkippedFullNodeContainerAllocation() {
numSkippedFullNodeContainerAllocation.incr();
}
+ void setTotalPendingContainerSlots(long value) {
+ totalPendingContainerSlots.set(value);
+ }
+
+ public long getTotalPendingContainerSlots() {
+ return totalPendingContainerSlots.value();
+ }
+
/**
* Get aggregated counter and gauge metrics.
*/
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelinePlacementPolicy.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelinePlacementPolicy.java
index 8e1437d474c..615d466f6e4 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelinePlacementPolicy.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelinePlacementPolicy.java
@@ -156,7 +156,7 @@ List<DatanodeDetails> filterViableNodes(
}
healthyNodes = filterNodesWithSpace(healthyNodes, nodesRequired,
- metadataSizeRequired, dataSizeRequired);
+ metadataSizeRequired);
boolean multipleRacks = multipleRacksAvailable(healthyNodes);
int excludedNodesSize = 0;
if (excludedNodes != null) {
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineProvider.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineProvider.java
index 0b8506f58c2..978ebbb0406 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineProvider.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineProvider.java
@@ -87,7 +87,8 @@ List<DatanodeDetails> pickNodesNotUsed(REPLICATION_CONFIG
replicationConfig,
int nodesRequired = replicationConfig.getRequiredNodes();
List<DatanodeDetails> healthyDNs = pickAllNodesNotUsed(replicationConfig);
List<DatanodeDetails> healthyDNsWithSpace = healthyDNs.stream()
- .filter(dn -> SCMCommonPlacementPolicy.hasEnoughSpace(dn,
metadataSizeRequired, dataSizeRequired))
+ .filter(dn -> SCMCommonPlacementPolicy.hasEnoughSpace(
+ dn, metadataSizeRequired, nodeManager))
.limit(nodesRequired)
.collect(Collectors.toList());
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/TestSCMCommonPlacementPolicy.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/TestSCMCommonPlacementPolicy.java
index 5ecd3fa1c74..9290e8d83be 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/TestSCMCommonPlacementPolicy.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/TestSCMCommonPlacementPolicy.java
@@ -17,7 +17,6 @@
package org.apache.hadoop.hdds.scm;
-import static
org.apache.hadoop.hdds.protocol.proto.HddsProtos.StorageTypeProto.DISK;
import static
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReplicaProto.State.CLOSED;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -474,7 +473,7 @@ protected List<DatanodeDetails> chooseDatanodesInternal(
}
@Test
- public void testDatanodeIsInvalidInCaseOfIncreasingCommittedBytes() {
+ public void testDatanodeIsInvalidWhenNoSlotsAvailable() {
NodeManager nodeMngr = mock(NodeManager.class);
final DatanodeID datanodeID = DatanodeID.of(UUID.randomUUID());
DummyPlacementPolicy placementPolicy =
@@ -488,43 +487,19 @@ public void
testDatanodeIsInvalidInCaseOfIncreasingCommittedBytes() {
when(datanodeInfo.getNodeStatus()).thenReturn(nodeStatus);
when(nodeMngr.getNode(eq(datanodeID))).thenReturn(datanodeInfo);
- // capacity = 200000, used = 90000, remaining = 101000, committed = 500
- StorageContainerDatanodeProtocolProtos.StorageReportProto storageReport1 =
- HddsTestUtils.createStorageReport(DatanodeID.randomID(), "/data/hdds",
- 200000, 90000, 101000, DISK).toBuilder()
- .setCommitted(500)
- .setFreeSpaceToSpare(10000)
- .build();
- // capacity = 200000, used = 90000, remaining = 101000, committed = 1000
- StorageContainerDatanodeProtocolProtos.StorageReportProto storageReport2 =
- HddsTestUtils.createStorageReport(DatanodeID.randomID(), "/data/hdds",
- 200000, 90000, 101000, DISK).toBuilder()
- .setCommitted(1000)
- .setFreeSpaceToSpare(100000)
- .build();
StorageContainerDatanodeProtocolProtos.MetadataStorageReportProto
metaReport =
HddsTestUtils.createMetadataStorageReport("/data/metadata",
200);
- when(datanodeInfo.getStorageReports())
- .thenReturn(Collections.singletonList(storageReport1))
- .thenReturn(Collections.singletonList(storageReport2));
when(datanodeInfo.getMetadataStorageReports())
.thenReturn(Collections.singletonList(metaReport));
-
- // 500 committed bytes:
- //
- // 101000 500
- // | |
- // (remaining - committed) > Math.max(4000, freeSpaceToSpare)
- // |
- // 100000
- //
- // Summary: 101000 - 500 > 100000 == true
+ // Space check now uses PendingContainerTracker.hasAvailableSpace:
+ // slot available → isValidNode returns true
+ when(nodeMngr.hasAvailableSpace(datanodeInfo)).thenReturn(true);
assertTrue(placementPolicy.isValidNode(datanodeDetails, 100, 4000));
- // 1000 committed bytes:
- // Summary: 101000 - 1000 > 100000 == false
+ // No slot available (all pending) → isValidNode returns false
+ when(nodeMngr.hasAvailableSpace(datanodeInfo)).thenReturn(false);
assertFalse(placementPolicy.isValidNode(datanodeDetails, 100, 4000));
}
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java
index 1b8c4bbf384..b9aa2930f01 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/MockNodeManager.java
@@ -461,6 +461,11 @@ public boolean checkSpaceAndRecordAllocation(DatanodeInfo
datanodeInfo, Containe
return pendingContainerTracker.checkSpaceAndRecordAllocation(datanodeInfo,
containerID);
}
+ @Override
+ public boolean hasAvailableSpace(DatanodeInfo datanodeInfo) {
+ return pendingContainerTracker.hasAvailableSpace(datanodeInfo);
+ }
+
@Override
public void removePendingAllocationForDatanode(DatanodeInfo datanodeInfo,
ContainerID containerID) {
if (datanodeInfo != null) {
@@ -955,6 +960,11 @@ public PendingContainerTracker
getPendingContainerTracker() {
return pendingContainerTracker;
}
+ public void setPendingContainerMaxSize(long maxContainerSize) {
+ this.pendingContainerTracker = new
PendingContainerTracker(maxContainerSize,
+ HddsTestUtils.ROLL_INTERVAL_MS_DEFAULT, null);
+ }
+
@Override
public long getLastHeartbeat(DatanodeDetails datanodeDetails) {
return -1;
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/SimpleMockNodeManager.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/SimpleMockNodeManager.java
index e47b979b475..ec481f07973 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/SimpleMockNodeManager.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/SimpleMockNodeManager.java
@@ -256,6 +256,11 @@ public boolean checkSpaceAndRecordAllocation(DatanodeInfo
datanodeInfo, Containe
return true;
}
+ @Override
+ public boolean hasAvailableSpace(DatanodeInfo datanodeInfo) {
+ return true;
+ }
+
@Override
public void removePendingAllocationForDatanode(DatanodeInfo datanodeInfo,
ContainerID containerID) {
}
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestContainerPlacementFactory.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestContainerPlacementFactory.java
index ed4e96b8de5..5e6f0ef3959 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestContainerPlacementFactory.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestContainerPlacementFactory.java
@@ -26,6 +26,7 @@
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertThrows;
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;
@@ -146,6 +147,7 @@ public void testRackAwarePolicy() throws IOException {
when(nodeManager.getNode(dn.getID()))
.thenReturn(dn);
}
+
when(nodeManager.hasAvailableSpace(any(DatanodeInfo.class))).thenReturn(true);
PlacementPolicy policy = ContainerPlacementPolicyFactory
.getPolicy(conf, nodeManager, cluster, true,
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementCapacity.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementCapacity.java
index 1fb3f53504d..b86bf3e58d0 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementCapacity.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementCapacity.java
@@ -118,6 +118,10 @@ public void chooseDatanodes() throws SCMException {
.filter(dn -> dn.getID().equals(invocation.getArgument(0)))
.findFirst()
.orElse(null));
+
when(mockNodeManager.hasAvailableSpace(any(DatanodeInfo.class))).thenAnswer(invocation
-> {
+ DatanodeInfo di = invocation.getArgument(0);
+ return di.getStorageReports().stream().anyMatch(r -> r.getRemaining() >=
15L);
+ });
SCMContainerPlacementCapacity scmContainerPlacementRandom =
new SCMContainerPlacementCapacity(mockNodeManager, conf, null, true,
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackAware.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackAware.java
index dd068f55cdf..5b1961df107 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackAware.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackAware.java
@@ -32,6 +32,7 @@
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assumptions.assumeTrue;
+import static org.mockito.Mockito.any;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
@@ -180,6 +181,7 @@ private void setup(int datanodeCount) {
}
when(nodeManager.getClusterNetworkTopologyMap())
.thenReturn(cluster);
+
when(nodeManager.hasAvailableSpace(any(DatanodeInfo.class))).thenReturn(true);
// create placement policy instances
policy = new SCMContainerPlacementRackAware(
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackScatter.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackScatter.java
index e015b93c1e3..cbbfae95298 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackScatter.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRackScatter.java
@@ -33,6 +33,7 @@
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assumptions.assumeTrue;
+import static org.mockito.Mockito.any;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
@@ -249,6 +250,10 @@ private void createMocksAndUpdateStorageReports(int
datanodeCount) {
}
when(nodeManager.getClusterNetworkTopologyMap())
.thenReturn(cluster);
+
when(nodeManager.hasAvailableSpace(any(DatanodeInfo.class))).thenAnswer(invocation
-> {
+ DatanodeInfo di = invocation.getArgument(0);
+ return di.getStorageReports().stream().anyMatch(r -> r.getRemaining() >
1L);
+ });
// create placement policy instances
policy = new SCMContainerPlacementRackScatter(
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRandom.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRandom.java
index 47602a385fd..b8c1a0f4c4f 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRandom.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/placement/algorithms/TestSCMContainerPlacementRandom.java
@@ -22,6 +22,7 @@
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
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;
@@ -89,6 +90,10 @@ public void chooseDatanodes() throws SCMException {
NodeManager mockNodeManager = mock(NodeManager.class);
when(mockNodeManager.getNodes(NodeStatus.inServiceHealthy()))
.thenReturn(new ArrayList<>(datanodes));
+
when(mockNodeManager.hasAvailableSpace(any(DatanodeInfo.class))).thenAnswer(invocation
-> {
+ DatanodeInfo di = invocation.getArgument(0);
+ return di.getStorageReports().stream().anyMatch(r -> r.getRemaining() >=
15L);
+ });
SCMContainerPlacementRandom scmContainerPlacementRandom =
new SCMContainerPlacementRandom(mockNodeManager, conf, null, true,
@@ -209,6 +214,10 @@ public void testIsValidNode() throws SCMException {
.thenReturn(datanodes.get(1));
when(mockNodeManager.getNode(datanodes.get(2).getID()))
.thenReturn(datanodes.get(2));
+
when(mockNodeManager.hasAvailableSpace(any(DatanodeInfo.class))).thenAnswer(invocation
-> {
+ DatanodeInfo di = invocation.getArgument(0);
+ return di.getStorageReports().stream().anyMatch(r -> r.getRemaining() >=
15L);
+ });
SCMContainerPlacementRandom scmContainerPlacementRandom =
new SCMContainerPlacementRandom(mockNodeManager, conf, null, true,
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestPendingContainerTracker.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestPendingContainerTracker.java
index b486715deb0..ff789aba141 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestPendingContainerTracker.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestPendingContainerTracker.java
@@ -23,6 +23,7 @@
import java.io.IOException;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.List;
import org.apache.hadoop.hdds.protocol.MockDatanodeDetails;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.StorageReportProto;
@@ -376,7 +377,71 @@ public void testMultiVolumeWithCommittedBytes() {
assertFalse(tracker.checkSpaceAndRecordAllocation(dnInfo,
containers.get(1)));
}
+ /**
+ * Pending in-flight replications recorded via checkSpaceAndRecordAllocation
count against
+ * slots, same as write-path containers. hasAvailableSpace reflects the
combined total.
+ */
+ @Test
+ public void testInFlightReplicationCountsAgainstAvailableSlots() {
+ long containerSize = MAX_CONTAINER_SIZE;
+ DatanodeInfo dnInfo = datanodes.get(0);
+
+ // Two slots of usable space
+ List<StorageReportProto> twoSlotReports = new ArrayList<>();
+ twoSlotReports.add(createStorageReport(dnInfo, 10 * containerSize, 2 *
containerSize, 0));
+ dnInfo.updateStorageReports(twoSlotReports);
+
+ assertTrue(tracker.hasAvailableSpace(dnInfo)); // 2 slots free
+ assertTrue(tracker.checkSpaceAndRecordAllocation(dnInfo,
containers.get(0))); // slot 1 used
+ assertTrue(tracker.hasAvailableSpace(dnInfo)); // 1 slot free
+ assertTrue(tracker.checkSpaceAndRecordAllocation(dnInfo,
containers.get(1))); // slot 2 used
+ assertFalse(tracker.hasAvailableSpace(dnInfo)); // 0 slots
free
+ assertFalse(tracker.checkSpaceAndRecordAllocation(dnInfo,
containers.get(2))); // rejected
+ }
+
+ /**
+ * hasAvailableSpace on a DN with no storage reports returns false.
+ */
+ @Test
+ public void testHasAvailableSpaceWithNoStorageReports() {
+ DatanodeInfo emptyDn = new DatanodeInfo(
+ MockDatanodeDetails.randomLocalDatanodeDetails(),
NodeStatus.inServiceHealthy(), null,
+ HddsTestUtils.ROLL_INTERVAL_MS_DEFAULT);
+ // No storage reports set
+ assertFalse(tracker.hasAvailableSpace(emptyDn));
+ }
+
+ /**
+ * A failed volume has remaining=0 by DN convention, but the tracker should
+ * explicitly skip it (report.getFailed() == true) so a stale non-zero
+ * remaining value on a failed volume can never grant spurious slots.
+ */
+ @Test
+ public void testFailedVolumeNotCountedAsAllocatableSlot() {
+ StorageReportProto failed = createFailedStorageReport(dn1);
+ StorageReportProto healthy = createStorageReport(dn1,
+ 10 * MAX_CONTAINER_SIZE, MAX_CONTAINER_SIZE, 0); // 1 real slot
+ dn1.updateStorageReports(new ArrayList<>(Arrays.asList(failed, healthy)));
+ assertTrue(tracker.hasAvailableSpace(dn1)); //
healthy vol → 1 slot
+ assertTrue(tracker.checkSpaceAndRecordAllocation(dn1, container1)); //
consumes it
+ assertFalse(tracker.hasAvailableSpace(dn1)); // 0
slots left
+ assertFalse(tracker.checkSpaceAndRecordAllocation(dn1, container2)); //
rejected
+ }
+
+ @Test
+ public void testAllVolumesFailedReturnsFalse() {
+ dn1.updateStorageReports(new ArrayList<>(Arrays.asList((
+ createFailedStorageReport(dn1)),
+ createFailedStorageReport(dn1))));
+ assertFalse(tracker.hasAvailableSpace(dn1));
+ assertFalse(tracker.checkSpaceAndRecordAllocation(dn1, container1));
+ }
+
private StorageReportProto createStorageReport(DatanodeInfo dn, long
capacity, long remaining, long committed) {
return HddsTestUtils.createStorageReports(dn.getID(), capacity, remaining,
committed).get(0);
}
+
+ private StorageReportProto createFailedStorageReport(DatanodeInfo dn) {
+ return HddsTestUtils.createStorageReport(dn.getID(), "", 0, 0, 0, null,
true);
+ }
}
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementFactory.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementFactory.java
index 84afb74f0a5..3ddb0361e08 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementFactory.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementFactory.java
@@ -27,6 +27,8 @@
import static org.junit.jupiter.api.Assertions.assertNotSame;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.any;
+import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.when;
@@ -128,6 +130,7 @@ private void setupRacks(int datanodeCount, int nodesPerRack,
when(nodeManager.getNode(dn.getID()))
.thenReturn(dn);
}
+
doReturn(true).when(nodeManager).hasAvailableSpace(any(DatanodeInfo.class));
DBStore dbStore = DBStoreBuilder.createDBStore(conf,
SCMDBDefinition.get());
SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true);
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java
index 90b7747350b..d76bf22b3fe 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java
@@ -258,6 +258,7 @@ public void testChooseNodeNotEnoughSpace() throws
IOException {
"the space requirement";
// A huge container size
+ localNodeManager.setPendingContainerMaxSize(200 * OzoneConsts.TB);
SCMException ex =
assertThrows(SCMException.class,
() -> localPlacementPolicy.chooseDatanodes(new
ArrayList<>(datanodes.size()),
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java
index ca9d1f5a6c3..bf352f2051e 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java
@@ -41,6 +41,7 @@
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
import org.apache.hadoop.hdds.client.ReplicationConfig;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.conf.StorageUnit;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.DatanodeID;
import org.apache.hadoop.hdds.protocol.MockDatanodeDetails;
@@ -97,6 +98,10 @@ public void init(int maxPipelinePerNode, OzoneConfiguration
conf, File dir) thro
dbStore = DBStoreBuilder.createDBStore(conf, SCMDBDefinition.get());
nodeManager = new MockNodeManager(true, nodeCount);
nodeManager.setNumPipelinePerDatanode(maxPipelinePerNode);
+ long containerSize = (long) conf.getStorageSize(
+ ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE,
+ ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE_DEFAULT, StorageUnit.BYTES);
+ nodeManager.setPendingContainerMaxSize(containerSize);
SCMHAManager scmhaManager = SCMHAManagerStub.getInstance(true);
conf.setInt(ScmConfigKeys.OZONE_DATANODE_PIPELINE_LIMIT,
maxPipelinePerNode);
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]