This is an automated email from the ASF dual-hosted git repository.

adoroszlai 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 4766aa8609b HDDS-15440. Combine pipelineMap and pipeline2container to 
a single map in PipelineStateMap (#10454)
4766aa8609b is described below

commit 4766aa8609b32ea279c3ea8249e1b835dfa57486
Author: Om Kenge <[email protected]>
AuthorDate: Sun Jul 19 02:02:46 2026 +0530

    HDDS-15440. Combine pipelineMap and pipeline2container to a single map in 
PipelineStateMap (#10454)
    
    Co-authored-by: Tsz-Wo Nicholas Sze <[email protected]>
---
 .../scm/pipeline/PipelineStateManagerImpl.java     |   2 +-
 .../hadoop/hdds/scm/pipeline/PipelineStateMap.java | 263 ++++++++++++---------
 2 files changed, 152 insertions(+), 113 deletions(-)

diff --git 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManagerImpl.java
 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManagerImpl.java
index 88cc37b0b5e..c23c424949c 100644
--- 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManagerImpl.java
+++ 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManagerImpl.java
@@ -128,7 +128,7 @@ public Pipeline getPipeline(PipelineID pipelineID)
       throws PipelineNotFoundException {
     lock.readLock().lock();
     try {
-      return pipelineStateMap.getPipeline(pipelineID);
+      return pipelineStateMap.getPipeline(pipelineID).getPipeline();
     } finally {
       lock.readLock().unlock();
     }
diff --git 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java
 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java
index 14fc9239149..2aab03b00b3 100644
--- 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java
+++ 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java
@@ -25,13 +25,12 @@
 import java.util.Collection;
 import java.util.Collections;
 import java.util.HashMap;
-import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
 import java.util.NavigableSet;
 import java.util.Objects;
-import java.util.Set;
 import java.util.TreeSet;
+import java.util.function.Predicate;
 import org.apache.hadoop.hdds.client.ReplicationConfig;
 import org.apache.hadoop.hdds.protocol.DatanodeDetails;
 import org.apache.hadoop.hdds.scm.container.ContainerID;
@@ -42,8 +41,11 @@
 /**
  * Holds the data structures which maintain the information about pipeline and
  * its state.
- * Invariant: If a pipeline exists in PipelineStateMap, both pipelineMap and
- * pipeline2container would have a non-null mapping for it.
+ *
+ * Invariant:
+ * If a pipeline exists in PipelineStateMap, pipelineMap contains a
+ * corresponding PipelineInfo, which stores both the Pipeline and its
+ * associated containers.
  *
  * Concurrency consideration:
  *   - thread-unsafe
@@ -52,8 +54,7 @@ class PipelineStateMap {
   private static final Logger LOG = 
LoggerFactory.getLogger(PipelineStateMap.class);
 
   // TODO: Use TreeMap for range operations?
-  private final Map<PipelineID, Pipeline> pipelineMap = new HashMap<>();
-  private final Map<PipelineID, NavigableSet<ContainerID>> pipeline2container 
= new HashMap<>();
+  private final Map<PipelineID, PipelineInfo> pipelineMap = new HashMap<>();
   private final Map<ReplicationConfig, List<Pipeline>> query2OpenPipelines = 
new HashMap<>();
 
   PipelineStateMap() { }
@@ -73,12 +74,12 @@ void addPipeline(Pipeline pipeline) throws 
DuplicatedPipelineIdException {
             pipeline.getNodes().size(), pipeline.getReplicationConfig()
                 .getRequiredNodes());
 
-    if (pipelineMap.putIfAbsent(pipeline.getId(), pipeline) != null) {
+    final PipelineInfo info = new PipelineInfo(pipeline);
+    if (pipelineMap.putIfAbsent(pipeline.getId(), info) != null) {
       LOG.warn("Duplicate pipeline ID detected. {}", pipeline.getId());
       throw new DuplicatedPipelineIdException(
           format("Duplicate pipeline ID %s detected.", pipeline.getId()));
     }
-    pipeline2container.put(pipeline.getId(), new TreeSet<>());
     if (pipeline.getPipelineState() == PipelineState.OPEN) {
       query2OpenPipelines.computeIfAbsent(pipeline.getReplicationConfig(), any 
-> new ArrayList<>())
           .add(pipeline);
@@ -93,17 +94,7 @@ void addPipeline(Pipeline pipeline) throws 
DuplicatedPipelineIdException {
    */
   void addContainerToPipeline(PipelineID pipelineID, ContainerID containerID)
       throws InvalidPipelineStateException, PipelineNotFoundException {
-    Objects.requireNonNull(pipelineID,
-        "Pipeline Id cannot be null");
-    Objects.requireNonNull(containerID,
-        "Container Id cannot be null");
-
-    Pipeline pipeline = getPipeline(pipelineID);
-    if (pipeline.isClosed()) {
-      throw new InvalidPipelineStateException(format(
-          "Cannot add container to pipeline=%s in closed state", pipelineID));
-    }
-    pipeline2container.get(pipelineID).add(containerID);
+    getPipeline(pipelineID).addContainerToOpenPipeline(containerID);
   }
 
   /**
@@ -114,23 +105,7 @@ void addContainerToPipeline(PipelineID pipelineID, 
ContainerID containerID)
    */
   void addContainerToPipelineSCMStart(PipelineID pipelineID, ContainerID 
containerID)
       throws PipelineNotFoundException {
-    Objects.requireNonNull(pipelineID,
-            "Pipeline Id cannot be null");
-    Objects.requireNonNull(containerID,
-            "Container Id cannot be null");
-
-    Pipeline pipeline = getPipeline(pipelineID);
-    if (pipeline.isClosed()) {
-      /*
-      When SCM restarts,the SCM DB may not be upto date where some
-      containers are in an OPEN state for a CLOSED pipeline. This happens when
-      close pipeline transaction in flushed before SCM goes down and close
-      container is not flushed into DB.
-      */
-      LOG.info("Container {} in open state for pipeline={} in closed state",
-              containerID, pipelineID);
-    }
-    pipeline2container.get(pipelineID).add(containerID);
+    getPipeline(pipelineID).addContainer(containerID);
   }
 
   /**
@@ -140,24 +115,27 @@ void addContainerToPipelineSCMStart(PipelineID 
pipelineID, ContainerID container
    * @return Pipeline
    * @throws PipelineNotFoundException if pipeline is not found
    */
-  Pipeline getPipeline(PipelineID pipelineID) throws PipelineNotFoundException 
{
-    Objects.requireNonNull(pipelineID,
-        "Pipeline Id cannot be null");
+  PipelineInfo getPipeline(PipelineID pipelineID) throws 
PipelineNotFoundException {
+    Objects.requireNonNull(pipelineID, "pipelineID == null");
+    final PipelineInfo info = pipelineMap.get(pipelineID);
 
-    Pipeline pipeline = pipelineMap.get(pipelineID);
-    if (pipeline == null) {
+    if (info == null) {
       throw new PipelineNotFoundException(
-          format("%s not found", pipelineID));
+          "Pipeline not found: " + pipelineID);
     }
-    return pipeline;
+    return info;
   }
 
   /**
    * Get list of pipelines in SCM.
    * @return List of pipelines
    */
-  public List<Pipeline> getPipelines() {
-    return new ArrayList<>(pipelineMap.values());
+  List<Pipeline> getPipelines() {
+    final List<Pipeline> pipelines = new ArrayList<>(pipelineMap.size());
+    for (PipelineInfo info : pipelineMap.values()) {
+      pipelines.add(info.getPipeline());
+    }
+    return pipelines;
   }
 
   /**
@@ -170,7 +148,8 @@ List<Pipeline> getPipelines(ReplicationConfig 
replicationConfig) {
     Objects.requireNonNull(replicationConfig, "ReplicationConfig cannot be 
null");
 
     List<Pipeline> pipelines = new ArrayList<>();
-    for (Pipeline pipeline : pipelineMap.values()) {
+    for (PipelineInfo info: pipelineMap.values()) {
+      final Pipeline pipeline = info.getPipeline();
       if (pipeline.getReplicationConfig().equals(replicationConfig)) {
         pipelines.add(pipeline);
       }
@@ -194,13 +173,12 @@ List<Pipeline> getPipelines(ReplicationConfig 
replicationConfig,
     Objects.requireNonNull(state, "Pipeline state cannot be null");
 
     if (state == PipelineState.OPEN) {
-      return new ArrayList<>(
-          query2OpenPipelines.getOrDefault(
-              replicationConfig, Collections.emptyList()));
+      return getOpenPipelines(replicationConfig);
     }
 
     List<Pipeline> pipelines = new ArrayList<>();
-    for (Pipeline pipeline : pipelineMap.values()) {
+    for (PipelineInfo info : pipelineMap.values()) {
+      final Pipeline pipeline = info.getPipeline();
       if (pipeline.getReplicationConfig().equals(replicationConfig)
           && pipeline.getPipelineState() == state) {
         pipelines.add(pipeline);
@@ -210,6 +188,11 @@ List<Pipeline> getPipelines(ReplicationConfig 
replicationConfig,
     return pipelines;
   }
 
+  private List<Pipeline> getOpenPipelines(ReplicationConfig replicationConfig) 
{
+    final List<Pipeline> pipelines = 
query2OpenPipelines.get(replicationConfig);
+    return pipelines != null && !pipelines.isEmpty() ? new 
ArrayList<>(pipelines) : Collections.emptyList();
+  }
+
   /**
    * Get a count of pipelines with the given replicationConfig and state.
    * This method is most efficient when getting a count for OPEN pipeline
@@ -225,12 +208,13 @@ int getPipelineCount(ReplicationConfig replicationConfig,
     Objects.requireNonNull(state, "Pipeline state cannot be null");
 
     if (state == PipelineState.OPEN) {
-      return query2OpenPipelines.getOrDefault(
-              replicationConfig, Collections.emptyList()).size();
+      final List<Pipeline> pipelines = 
query2OpenPipelines.get(replicationConfig);
+      return pipelines != null && !pipelines.isEmpty() ? pipelines.size() : 0;
     }
 
     int count = 0;
-    for (Pipeline pipeline : pipelineMap.values()) {
+    for (PipelineInfo info : pipelineMap.values()) {
+      final Pipeline pipeline = info.getPipeline();
       if (pipeline.getReplicationConfig().equals(replicationConfig)
           && pipeline.getPipelineState() == state) {
         count++;
@@ -239,6 +223,36 @@ int getPipelineCount(ReplicationConfig replicationConfig,
     return count;
   }
 
+  static Predicate<Pipeline> notInExcludeDatanodes(Collection<DatanodeDetails> 
excludeDatanodes) {
+    return p -> p.getNodeSet().stream().noneMatch(excludeDatanodes::contains);
+  }
+
+  static Predicate<Pipeline> notInExcludePipelines(Collection<PipelineID> 
excludePipelines) {
+    return p -> !excludePipelines.contains(p.getId());
+  }
+
+  static Predicate<Pipeline> getPredicate(
+      Collection<DatanodeDetails> excludeDatanodes,
+      Collection<PipelineID> excludePipelines) {
+    if (excludeDatanodes.isEmpty()) {
+      return excludePipelines.isEmpty() ? p -> true : 
notInExcludePipelines(excludePipelines);
+    } else {
+      final Predicate<Pipeline> n = notInExcludeDatanodes(excludeDatanodes);
+      return excludePipelines.isEmpty() ? n : p -> 
notInExcludePipelines(excludePipelines).test(p) && n.test(p);
+    }
+  }
+
+  static Predicate<Pipeline> getPredicate(
+      Collection<DatanodeDetails> excludeDatanodes,
+      Collection<PipelineID> excludePipelines,
+      ReplicationConfig replicationConfig,
+      PipelineState state) {
+    final Predicate<Pipeline> include = getPredicate(excludeDatanodes, 
excludePipelines);
+    return p -> p.getPipelineState() == state
+        && p.getReplicationConfig().equals(replicationConfig)
+        && include.test(p);
+  }
+
   /**
    * Get list of pipeline corresponding to specified replication type,
    * replication factor and pipeline state.
@@ -258,34 +272,25 @@ List<Pipeline> getPipelines(ReplicationConfig 
replicationConfig,
     Objects.requireNonNull(excludeDns, "Datanode exclude list cannot be null");
     Objects.requireNonNull(excludePipelines, "Pipeline exclude list cannot be 
null");
 
-    List<Pipeline> pipelines = null;
     if (state == PipelineState.OPEN) {
-      pipelines = new ArrayList<>(query2OpenPipelines.getOrDefault(
-          replicationConfig, Collections.emptyList()));
+      final List<Pipeline> pipelines = getOpenPipelines(replicationConfig);
       if (excludeDns.isEmpty() && excludePipelines.isEmpty()) {
         return pipelines;
       }
-    } else {
-      pipelines = new ArrayList<>(pipelineMap.values());
-    }
 
-    Iterator<Pipeline> iter = pipelines.iterator();
-    while (iter.hasNext()) {
-      Pipeline pipeline = iter.next();
-      if (!pipeline.getReplicationConfig().equals(replicationConfig) ||
-          pipeline.getPipelineState() != state ||
-          excludePipelines.contains(pipeline.getId())) {
-        iter.remove();
-      } else {
-        for (DatanodeDetails dn : pipeline.getNodes()) {
-          if (excludeDns.contains(dn)) {
-            iter.remove();
-            break;
-          }
-        }
-      }
+      final Predicate<Pipeline> include = getPredicate(excludeDns, 
excludePipelines);
+      pipelines.removeIf(pipeline -> !include.test(pipeline));
+      return pipelines;
     }
 
+    final Predicate<Pipeline> include = getPredicate(excludeDns, 
excludePipelines, replicationConfig, state);
+    final List<Pipeline> pipelines = new ArrayList<>(pipelineMap.size() / 2 + 
1); // only resize once
+    for (PipelineInfo info : pipelineMap.values()) {
+      final Pipeline pipeline = info.getPipeline();
+      if (include.test(pipeline)) {
+        pipelines.add(pipeline);
+      }
+    } 
     return pipelines;
   }
 
@@ -298,15 +303,7 @@ List<Pipeline> getPipelines(ReplicationConfig 
replicationConfig,
    */
   NavigableSet<ContainerID> getContainers(PipelineID pipelineID)
       throws PipelineNotFoundException {
-    Objects.requireNonNull(pipelineID,
-        "Pipeline Id cannot be null");
-
-    NavigableSet<ContainerID> containerIDs = 
pipeline2container.get(pipelineID);
-    if (containerIDs == null) {
-      throw new PipelineNotFoundException(
-          format("%s not found", pipelineID));
-    }
-    return new TreeSet<>(containerIDs);
+    return getPipeline(pipelineID).copyContainers();
   }
 
   /**
@@ -318,15 +315,7 @@ NavigableSet<ContainerID> getContainers(PipelineID 
pipelineID)
    */
   int getNumberOfContainers(PipelineID pipelineID)
       throws PipelineNotFoundException {
-    Objects.requireNonNull(pipelineID,
-        "Pipeline Id cannot be null");
-
-    Set<ContainerID> containerIDs = pipeline2container.get(pipelineID);
-    if (containerIDs == null) {
-      throw new PipelineNotFoundException(
-          format("%s not found", pipelineID));
-    }
-    return containerIDs.size();
+    return getPipeline(pipelineID).getContainers().size();
   }
 
   /**
@@ -337,14 +326,22 @@ int getNumberOfContainers(PipelineID pipelineID)
   Pipeline removePipeline(PipelineID pipelineID) throws 
PipelineNotFoundException, InvalidPipelineStateException {
     Objects.requireNonNull(pipelineID, "Pipeline Id cannot be null");
 
-    Pipeline pipeline = getPipeline(pipelineID);
+    // Check existence first, before removing
+    final PipelineInfo info = pipelineMap.get(pipelineID);
+    if (info == null) {
+      throw new PipelineNotFoundException("Pipeline not found: " + pipelineID);
+    }
+    final Pipeline pipeline = info.getPipeline();
     if (!pipeline.isClosed()) {
       throw new InvalidPipelineStateException(
           format("Pipeline with %s is not yet closed", pipelineID));
     }
+    List<Pipeline> pipelineList = 
query2OpenPipelines.get(pipeline.getReplicationConfig());
 
+    if (pipelineList != null) { 
+      pipelineList.remove(pipeline);
+    }
     pipelineMap.remove(pipelineID);
-    pipeline2container.remove(pipelineID);
     return pipeline;
   }
 
@@ -356,17 +353,7 @@ Pipeline removePipeline(PipelineID pipelineID) throws 
PipelineNotFoundException,
    * @param containerID - ContainerID of the container to remove
    */
   void removeContainerFromPipeline(PipelineID pipelineID, ContainerID 
containerID) throws PipelineNotFoundException {
-    Objects.requireNonNull(pipelineID,
-        "Pipeline Id cannot be null");
-    Objects.requireNonNull(containerID,
-        "container Id cannot be null");
-
-    Set<ContainerID> containerIDs = pipeline2container.get(pipelineID);
-    if (containerIDs == null) {
-      throw new PipelineNotFoundException(
-          format("%s not found", pipelineID));
-    }
-    containerIDs.remove(containerID);
+    getPipeline(pipelineID).removeContainer(containerID);
   }
 
   /**
@@ -380,29 +367,35 @@ void removeContainerFromPipeline(PipelineID pipelineID, 
ContainerID containerID)
    */
   Pipeline updatePipelineState(PipelineID pipelineID, PipelineState state)
       throws PipelineNotFoundException {
-    Objects.requireNonNull(pipelineID, "Pipeline Id cannot be null");
     Objects.requireNonNull(state, "Pipeline LifeCycleState cannot be null");
 
-    final Pipeline pipeline = getPipeline(pipelineID);
+    final PipelineInfo info = getPipeline(pipelineID);
+    final Pipeline pipeline = info.getPipeline();
     // Return the old pipeline if updating same state
     if (pipeline.getPipelineState() == state) {
       LOG.debug("CurrentState and NewState are the same, return from " +
           "updatePipelineState directly.");
       return pipeline;
     }
-    Pipeline updatedPipeline = pipelineMap.compute(pipelineID,
-        (id, p) -> pipeline.toBuilder().setState(state).build());
+    final Pipeline updated = pipeline.toBuilder().setState(state).build();
+    PipelineInfo newInfo = new PipelineInfo(updated);
+
+    for (ContainerID cid : info.getContainers()) {
+      newInfo.addContainer(cid);
+    }
+
+    pipelineMap.put(pipelineID, newInfo);
 
     List<Pipeline> pipelineList =
         query2OpenPipelines.get(pipeline.getReplicationConfig());
 
-    if (updatedPipeline.getPipelineState() == PipelineState.OPEN) {
+    if (updated.getPipelineState() == PipelineState.OPEN) {
       // for transition to OPEN state add pipeline to query2OpenPipelines
       if (pipelineList == null) {
         pipelineList = new ArrayList<>();
         query2OpenPipelines.put(pipeline.getReplicationConfig(), pipelineList);
       }
-      pipelineList.add(updatedPipeline);
+      pipelineList.add(updated);
     } else {
       // for transition from OPEN to CLOSED state remove pipeline from
       // query2OpenPipelines
@@ -410,7 +403,53 @@ Pipeline updatePipelineState(PipelineID pipelineID, 
PipelineState state)
         pipelineList.remove(pipeline);
       }
     }
-    return updatedPipeline;
+    return updated;
   }
 
+  static class PipelineInfo {
+    private final Pipeline pipeline;
+    private final NavigableSet<ContainerID> containers = new TreeSet<>();
+
+    PipelineInfo(Pipeline pipeline) {
+      this.pipeline = pipeline;
+    }
+
+    Pipeline getPipeline() {
+      return pipeline;
+    }
+
+    NavigableSet<ContainerID> getContainers() {
+      return containers;
+    }
+    
+    NavigableSet<ContainerID> copyContainers() {
+      return new TreeSet<>(containers);
+    }
+
+    void addContainerToOpenPipeline(ContainerID containerID) throws 
InvalidPipelineStateException {
+      Objects.requireNonNull(containerID, "Container Id == null");
+      if (pipeline.isClosed()) {
+        throw new InvalidPipelineStateException(
+            "Pipeline closed: Failed add container " + containerID + " to 
pipeline " + pipeline.getId());
+      }
+      containers.add(containerID);
+    }
+
+    void addContainer(ContainerID containerID) {
+      Objects.requireNonNull(containerID, "Container Id == null");
+      if (pipeline.isClosed()) {
+        // When SCM restarts, the SCM DB may not be up-to-dated,
+        // where some containers are in an OPEN state for a CLOSED pipeline.
+        // This happens when close pipeline transaction in flushed
+        // before SCM goes down and close container is not flushed into DB.
+        LOG.info("Container {} in open state for pipeline={} in closed state", 
containerID, pipeline.getId());
+      }
+      containers.add(containerID);
+    }
+
+    void removeContainer(ContainerID containerID) {
+      Objects.requireNonNull(containerID, "Container Id == null");
+      containers.remove(containerID);
+    }
+  }
 }


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

Reply via email to