diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/ContainerSet.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/ContainerSet.java index d122dfb587f9..43a3fc0a00bd 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/ContainerSet.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/ContainerSet.java @@ -73,8 +73,9 @@ public class ContainerSet implements Iterable> { ConcurrentSkipListMap<>(); private final ConcurrentSkipListSet missingContainerSet = new ConcurrentSkipListSet<>(); - private final ConcurrentSkipListMap recoveringContainerMap = - new ConcurrentSkipListMap<>(); + + private final ConcurrentSkipListSet recoveringContainerSet = + new ConcurrentSkipListSet<>(); private final Clock clock; private long recoveringTimeout; @Nullable @@ -208,8 +209,9 @@ private boolean addContainer(Container container, boolean overwrite) throws updateContainerIdTable(containerId, container.getContainerData()); missingContainerSet.remove(containerId); if (container.getContainerData().getState() == RECOVERING) { - recoveringContainerMap.put( - clock.millis() + recoveringTimeout, containerId); + recoveringContainerSet.add( + new RecoveringContainer(clock.millis() + recoveringTimeout, + containerId)); } HddsVolume volume = container.getContainerData().getVolume(); if (volume != null) { @@ -421,16 +423,16 @@ public boolean removeRecoveringContainer(long containerId) { Preconditions.checkState(containerId >= 0, "Container Id cannot be negative."); //it might take a little long time to iterate all the entries - // in recoveringContainerMap, but it seems ok here since: + // in recoveringContainerSet, but it seems ok here since: // 1 In the vast majority of cases,there will not be too // many recovering containers. // 2 closing container is not a sort of urgent action // // we can revisit here if any performance problem happens - Iterator> it = getRecoveringContainerIterator(); + Iterator it = getRecoveringContainerIterator(); while (it.hasNext()) { - Map.Entry entry = it.next(); - if (entry.getValue() == containerId) { + RecoveringContainer entry = it.next(); + if (entry.getContainerId() == containerId) { it.remove(); return true; } @@ -489,11 +491,11 @@ public Iterator> iterator() { /** * Return an container Iterator over - * {@link ContainerSet#recoveringContainerMap}. - * @return {@literal Iterator>} + * {@link ContainerSet#recoveringContainerSet}. + * @return {@literal Iterator} */ - public Iterator> getRecoveringContainerIterator() { - return recoveringContainerMap.entrySet().iterator(); + public Iterator getRecoveringContainerIterator() { + return recoveringContainerSet.iterator(); } /** @@ -667,4 +669,52 @@ public void buildMissingContainerSetAndValidate(Map container2BCSID } }); } + + /** + * A class that holds information about a recovering container. + */ + public static class RecoveringContainer + implements Comparable { + private final long timeout; + private final long containerId; + + public RecoveringContainer(long timeout, long containerId) { + this.timeout = timeout; + this.containerId = containerId; + } + + public long getTimeout() { + return timeout; + } + + public long getContainerId() { + return containerId; + } + + @Override + public int compareTo(RecoveringContainer other) { + int timeoutCompare = Long.compare(this.timeout, other.timeout); + if (timeoutCompare != 0) { + return timeoutCompare; + } + return Long.compare(this.containerId, other.containerId); + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (o == null || getClass() != o.getClass()) { + return false; + } + RecoveringContainer that = (RecoveringContainer) o; + return timeout == that.timeout && containerId == that.containerId; + } + + @Override + public int hashCode() { + return Objects.hash(timeout, containerId); + } + } } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/statemachine/background/StaleRecoveringContainerScrubbingService.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/statemachine/background/StaleRecoveringContainerScrubbingService.java index 5535c5128ccc..9c535e5f6e94 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/statemachine/background/StaleRecoveringContainerScrubbingService.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/statemachine/background/StaleRecoveringContainerScrubbingService.java @@ -18,13 +18,13 @@ package org.apache.hadoop.ozone.container.keyvalue.statemachine.background; import java.util.Iterator; -import java.util.Map; import java.util.concurrent.TimeUnit; import org.apache.hadoop.hdds.utils.BackgroundService; import org.apache.hadoop.hdds.utils.BackgroundTask; import org.apache.hadoop.hdds.utils.BackgroundTaskQueue; import org.apache.hadoop.hdds.utils.BackgroundTaskResult; import org.apache.hadoop.ozone.container.common.impl.ContainerSet; +import org.apache.hadoop.ozone.container.common.impl.ContainerSet.RecoveringContainer; import org.apache.hadoop.ozone.container.common.interfaces.Container; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -55,13 +55,13 @@ public BackgroundTaskQueue getTasks() { BackgroundTaskQueue backgroundTaskQueue = new BackgroundTaskQueue(); long currentTime = containerSet.getCurrentTime(); - Iterator> it = + Iterator it = containerSet.getRecoveringContainerIterator(); while (it.hasNext()) { - Map.Entry entry = it.next(); - if (currentTime >= entry.getKey()) { + RecoveringContainer entry = it.next(); + if (currentTime >= entry.getTimeout()) { backgroundTaskQueue.add(new RecoveringContainerScrubbingTask( - containerSet, entry.getValue())); + containerSet, entry.getContainerId())); it.remove(); } else { break;