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 29b2ad5fcd4 HDDS-14924. Account replication in PendingContainerTracker 
for space usage. (#10440)
29b2ad5fcd4 is described below

commit 29b2ad5fcd40642c77f771128e7f0f7c4fe88e42
Author: Ashish Kumar <[email protected]>
AuthorDate: Fri Jun 26 09:52:48 2026 +0530

    HDDS-14924. Account replication in PendingContainerTracker for space usage. 
(#10440)
---
 .../replication/ContainerReplicaPendingOps.java    | 22 ++++++--
 .../ContainerReplicaPendingOpsSubscriber.java      | 10 ++++
 .../apache/hadoop/hdds/scm/node/NodeManager.java   |  8 +++
 .../hdds/scm/node/PendingContainerTracker.java     | 27 ++++++++++
 .../hadoop/hdds/scm/node/SCMNodeManager.java       | 29 ++++++++++-
 .../hdds/scm/server/StorageContainerManager.java   |  5 ++
 .../hadoop/hdds/scm/container/MockNodeManager.java |  7 +++
 .../hdds/scm/container/SimpleMockNodeManager.java  |  4 ++
 .../TestContainerReplicaPendingOps.java            | 59 ++++++++++++++++++++++
 9 files changed, 167 insertions(+), 4 deletions(-)

diff --git 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ContainerReplicaPendingOps.java
 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ContainerReplicaPendingOps.java
index 2905ae4d4a3..1405c6e85f6 100644
--- 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ContainerReplicaPendingOps.java
+++ 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ContainerReplicaPendingOps.java
@@ -329,16 +329,18 @@ private void updateTimeoutMetrics(ContainerReplicaOp op) {
   private void addReplica(ContainerReplicaOp.PendingOpType opType,
       ContainerID containerID, DatanodeDetails target, int replicaIndex, 
SCMCommand<?> command,
       long deadlineEpochMillis, long containerSize, long scheduledEpochMillis) 
{
+    ContainerReplicaOp op = new ContainerReplicaOp(opType,
+        target, replicaIndex, command, deadlineEpochMillis, containerSize);
     Lock lock = writeLock(containerID);
     lock(lock);
+    boolean found;
     try {
       // Remove any existing duplicate op for the same target and replicaIndex 
before adding
       // the new one. Especially for delete ops, they could be getting resent 
after expiry.
-      completeOp(opType, containerID, target, replicaIndex, false);
+      found = completeOp(opType, containerID, target, replicaIndex, false);
       List<ContainerReplicaOp> ops = pendingOps.computeIfAbsent(
           containerID, s -> new ArrayList<>());
-      ops.add(new ContainerReplicaOp(opType,
-          target, replicaIndex, command, deadlineEpochMillis, containerSize));
+      ops.add(op);
       DatanodeID id = target.getID();
       if (opType == ADD) {
         containerSizeScheduled.compute(id, (k, v) -> {
@@ -353,6 +355,10 @@ private void addReplica(ContainerReplicaOp.PendingOpType 
opType,
     } finally {
       unlock(lock);
     }
+    // Notify for ADD ops to record container slot.
+    if (opType == ADD && !found) {
+      notifySubscribersOpAdded(op, containerID);
+    }
   }
 
   private boolean completeOp(ContainerReplicaOp.PendingOpType opType,
@@ -417,6 +423,16 @@ private void notifySubscribers(List<ContainerReplicaOp> 
ops,
     }
   }
 
+  /**
+   * Notifies subscribers that an ADD op was added for the given containerID.
+   */
+  private void notifySubscribersOpAdded(ContainerReplicaOp op,
+      ContainerID containerID) {
+    for (ContainerReplicaPendingOpsSubscriber subscriber : subscribers) {
+      subscriber.opAdded(op, containerID);
+    }
+  }
+
   /**
    * Registers a subscriber that will be notified about completed ops.
    *
diff --git 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ContainerReplicaPendingOpsSubscriber.java
 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ContainerReplicaPendingOpsSubscriber.java
index c0c9085679b..3a9ec2c4c25 100644
--- 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ContainerReplicaPendingOpsSubscriber.java
+++ 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ContainerReplicaPendingOpsSubscriber.java
@@ -25,6 +25,16 @@
  */
 public interface ContainerReplicaPendingOpsSubscriber {
 
+  /**
+   * Notifies that the specified op has been added for the specified
+   * containerID.
+   *
+   * @param op Add or Delete op
+   * @param containerID container on which the operation is being performed
+   */
+  default void opAdded(ContainerReplicaOp op, ContainerID containerID) {
+  }
+
   /**
    * Notifies that the specified op has been completed for the specified
    * containerID. Might have completed normally or timed out.
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 3da84aad3a7..51eb0e4e41f 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,14 @@ default int getAllNodeCount() {
    */
   boolean checkSpaceAndRecordAllocation(DatanodeInfo datanodeInfo, ContainerID 
containerID);
 
+  /**
+   * Records a container allocation on the given datanode.
+   * Unlike {@link #checkSpaceAndRecordAllocation}, this does not check for
+   * available space — it is called after the placement policy has already
+   * validated space and a replication command has been committed.
+   */
+  void recordAllocationForDatanode(DatanodeInfo datanodeInfo, ContainerID 
containerID);
+
   /**
    * Returns true if the datanode has at least one available container slot 
considering
    * in-flight allocations tracked by PendingContainerTracker.
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 229dc5fad92..eb17b3ccf86 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
@@ -144,6 +144,15 @@ synchronized int getCount() {
       return currentWindow.size() + previousWindow.size();
     }
 
+    /**
+     * Records a container allocation in the current window,
+     * without checking available space. Use this when the space check has
+     * already been performed by the placement policy.
+     */
+    synchronized void add(ContainerID containerID) {
+      currentWindow.add(containerID);
+    }
+
     /**
      * Atomically checks whether there is allocatable space for one more 
container of
      * {@code maxContainerSize} given the current pending count, and adds 
{@code containerID}
@@ -215,6 +224,24 @@ public boolean checkSpaceAndRecordAllocation(DatanodeInfo 
datanodeInfo, Containe
     return added;
   }
 
+  /**
+   * Records a container allocation on the given datanode in the
+   * current window, without performing a space check. This is used when the
+   * space check was already done by the placement policy (e.g. from
+   * {@link 
org.apache.hadoop.hdds.scm.container.replication.ContainerReplicaPendingOps}).
+   *
+   * @param datanodeInfo the datanode receiving the container
+   * @param containerID  the container being allocated
+   */
+  public void recordAllocation(DatanodeInfo datanodeInfo, ContainerID 
containerID) {
+    Objects.requireNonNull(datanodeInfo, "datanodeInfo == null");
+    Objects.requireNonNull(containerID, "containerID == null");
+    datanodeInfo.getPendingContainerAllocations().add(containerID);
+    if (metrics != null) {
+      metrics.incNumPendingContainersAdded();
+    }
+  }
+
   /**
    * Returns true if the given datanode has at least one allocatable container 
slot
    * available, accounting for pending in-flight allocations.
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 66a41fe773d..394793bb9a3 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
@@ -74,6 +74,8 @@
 import org.apache.hadoop.hdds.scm.container.ContainerID;
 import org.apache.hadoop.hdds.scm.container.placement.metrics.SCMNodeMetric;
 import org.apache.hadoop.hdds.scm.container.placement.metrics.SCMNodeStat;
+import org.apache.hadoop.hdds.scm.container.replication.ContainerReplicaOp;
+import 
org.apache.hadoop.hdds.scm.container.replication.ContainerReplicaPendingOpsSubscriber;
 import org.apache.hadoop.hdds.scm.events.SCMEvents;
 import org.apache.hadoop.hdds.scm.ha.SCMContext;
 import org.apache.hadoop.hdds.scm.net.NetworkTopology;
@@ -117,7 +119,7 @@
  * get functions in this file as a snap-shot of information that is 
inconsistent
  * as soon as you read it.
  */
-public class SCMNodeManager implements NodeManager {
+public class SCMNodeManager implements NodeManager, 
ContainerReplicaPendingOpsSubscriber {
 
   private static final Logger LOG =
       LoggerFactory.getLogger(SCMNodeManager.class);
@@ -1084,6 +1086,11 @@ public boolean 
checkSpaceAndRecordAllocation(DatanodeInfo datanodeInfo, Containe
     return pendingContainerTracker.checkSpaceAndRecordAllocation(datanodeInfo, 
containerID);
   }
 
+  @Override
+  public void recordAllocationForDatanode(DatanodeInfo datanodeInfo, 
ContainerID containerID) {
+    pendingContainerTracker.recordAllocation(datanodeInfo, containerID);
+  }
+  
   @Override
   public boolean hasAvailableSpace(DatanodeInfo datanodeInfo) {
     return pendingContainerTracker.hasAvailableSpace(datanodeInfo);
@@ -1095,6 +1102,26 @@ public void 
removePendingAllocationForDatanode(DatanodeInfo datanodeInfo, Contai
         datanodeInfo.getPendingContainerAllocations(), containerID);
   }
 
+  @Override
+  public void opAdded(ContainerReplicaOp op, ContainerID containerID) {
+    if (op.getOpType() == ContainerReplicaOp.PendingOpType.ADD) {
+      DatanodeInfo dnInfo = getDatanodeInfo(op.getTarget());
+      if (dnInfo != null) {
+        recordAllocationForDatanode(dnInfo, containerID);
+      }
+    }
+  }
+
+  @Override
+  public void opCompleted(ContainerReplicaOp op, ContainerID containerID, 
boolean timedOut) {
+    if (op.getOpType() == ContainerReplicaOp.PendingOpType.ADD && !timedOut) {
+      DatanodeInfo dnInfo = getDatanodeInfo(op.getTarget());
+      if (dnInfo != null) {
+        removePendingAllocationForDatanode(dnInfo, containerID);
+      }
+    }
+  }
+
   /**
    * Return the node stat of the specified datanode.
    *
diff --git 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java
 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java
index 69f7973ff1b..59a04d6a9a5 100644
--- 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java
+++ 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java
@@ -96,6 +96,7 @@
 import 
org.apache.hadoop.hdds.scm.container.placement.metrics.SCMPerformanceMetrics;
 import 
org.apache.hadoop.hdds.scm.container.reconciliation.ReconcileContainerEventHandler;
 import 
org.apache.hadoop.hdds.scm.container.replication.ContainerReplicaPendingOps;
+import 
org.apache.hadoop.hdds.scm.container.replication.ContainerReplicaPendingOpsSubscriber;
 import 
org.apache.hadoop.hdds.scm.container.replication.DatanodeCommandCountUpdatedHandler;
 import org.apache.hadoop.hdds.scm.container.replication.ReplicationManager;
 import 
org.apache.hadoop.hdds.scm.container.replication.ReplicationManagerEventHandler;
@@ -469,6 +470,10 @@ private StorageContainerManager(OzoneConfiguration conf,
 
     moveManager = new MoveManager(replicationManager, containerManager);
     containerReplicaPendingOps.registerSubscriber(moveManager);
+    if (scmNodeManager instanceof ContainerReplicaPendingOpsSubscriber) {
+      containerReplicaPendingOps.registerSubscriber(
+          (ContainerReplicaPendingOpsSubscriber) scmNodeManager);
+    }
     containerBalancer = new ContainerBalancer(this);
 
     // Emit initial safe mode status, as now handlers are registered.
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 b9aa2930f01..3da501228cb 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
@@ -462,6 +462,13 @@ public boolean checkSpaceAndRecordAllocation(DatanodeInfo 
datanodeInfo, Containe
   }
 
   @Override
+  public void recordAllocationForDatanode(DatanodeInfo datanodeInfo, 
ContainerID containerID) {
+    if (datanodeInfo != null) {
+      pendingContainerTracker.recordAllocation(datanodeInfo, containerID);
+    }
+  }
+    
+  @Override  
   public boolean hasAvailableSpace(DatanodeInfo datanodeInfo) {
     return pendingContainerTracker.hasAvailableSpace(datanodeInfo);
   }
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 ec481f07973..5b10f9d2d75 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,10 @@ public boolean checkSpaceAndRecordAllocation(DatanodeInfo 
datanodeInfo, Containe
     return true;
   }
 
+  @Override
+  public void recordAllocationForDatanode(DatanodeInfo datanodeInfo, 
ContainerID containerID) {
+  }
+  
   @Override
   public boolean hasAvailableSpace(DatanodeInfo datanodeInfo) {
     return true;
diff --git 
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestContainerReplicaPendingOps.java
 
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestContainerReplicaPendingOps.java
index ee813f0942c..2130e009dd5 100644
--- 
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestContainerReplicaPendingOps.java
+++ 
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestContainerReplicaPendingOps.java
@@ -582,4 +582,63 @@ public void 
testOnlyExpiredOpSizeIsRemovedFromSizeScheduledMap() {
     assertNull(scheduled.get(dn2.getID()));
     assertEquals(THREE_GB_CONTAINER_SIZE, 
scheduled.get(dn1.getID()).getSize());
   }
+
+  /**
+   * scheduleAddReplica must notify subscribers via opAdded() so that
+   * a subscribed NodeManager can record the PendingContainerTracker slot.
+   */
+  @Test
+  public void testScheduleAddReplicaNotifiesSubscriberOpAdded() {
+    ContainerReplicaPendingOpsSubscriber subscriber = 
mock(ContainerReplicaPendingOpsSubscriber.class);
+    ContainerID containerID = ContainerID.valueOf(1);
+
+    pendingOps.registerSubscriber(subscriber);
+    pendingOps.scheduleAddReplica(containerID, dn1, 0, addCmd, deadline,
+        FIVE_GB_CONTAINER_SIZE, clock.millis());
+
+    verify(subscriber, times(1)).opAdded(
+        org.mockito.ArgumentMatchers.argThat(op ->
+            op.getOpType() == ADD && op.getTarget().equals(dn1)),
+        org.mockito.ArgumentMatchers.eq(containerID));
+  }
+
+  /**
+   * scheduleDeleteReplica must NOT invoke opAdded() on subscribers — only ADD 
ops
+   * reserve container slots in the PendingContainerTracker.
+   */
+  @Test
+  public void testScheduleDeleteReplicaDoesNotNotifyOpAdded() {
+    ContainerReplicaPendingOpsSubscriber subscriber = 
mock(ContainerReplicaPendingOpsSubscriber.class);
+    ContainerID containerID = ContainerID.valueOf(1);
+
+    pendingOps.registerSubscriber(subscriber);
+    pendingOps.scheduleDeleteReplica(containerID, dn1, 0, deleteCmd, deadline);
+
+    verifyNoMoreInteractions(subscriber);
+  }
+
+  /**
+   * completeAddReplica must notify subscribers via 
opCompleted(timedOut=false) so that
+   * a subscribed NodeManager can release the PendingContainerTracker slot.
+   */
+  @Test
+  public void testCompleteAddReplicaNotifiesSubscriberOpCompleted() {
+    ContainerReplicaPendingOpsSubscriber subscriber = 
mock(ContainerReplicaPendingOpsSubscriber.class);
+    ContainerID containerID = ContainerID.valueOf(1);
+
+    pendingOps.registerSubscriber(subscriber);
+    pendingOps.scheduleAddReplica(containerID, dn1, 0, addCmd, deadline,
+        FIVE_GB_CONTAINER_SIZE, clock.millis());
+    pendingOps.completeAddReplica(containerID, dn1, 0);
+
+    verify(subscriber, times(1)).opAdded(
+        org.mockito.ArgumentMatchers.argThat(op ->
+            op.getOpType() == ADD && op.getTarget().equals(dn1)),
+        org.mockito.ArgumentMatchers.eq(containerID));
+    verify(subscriber, times(1)).opCompleted(
+        org.mockito.ArgumentMatchers.argThat(op ->
+            op.getOpType() == ADD && op.getTarget().equals(dn1)),
+        org.mockito.ArgumentMatchers.eq(containerID),
+        org.mockito.ArgumentMatchers.eq(false));
+  }
 }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to