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]