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 d39ca0af90b HDDS-12669. Race condition between entries of
ContainerSet#recoveringContainerMap (#10481)
d39ca0af90b is described below
commit d39ca0af90b0c3d52014a3d71246445087625bad
Author: Chung En Lee <[email protected]>
AuthorDate: Thu Jun 11 16:10:12 2026 +0800
HDDS-12669. Race condition between entries of
ContainerSet#recoveringContainerMap (#10481)
---
.../ozone/container/common/impl/ContainerSet.java | 74 ++++++++++++++++++----
.../StaleRecoveringContainerScrubbingService.java | 10 +--
2 files changed, 67 insertions(+), 17 deletions(-)
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 d122dfb587f..43a3fc0a00b 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<Container<?>> {
ConcurrentSkipListMap<>();
private final ConcurrentSkipListSet<Long> missingContainerSet =
new ConcurrentSkipListSet<>();
- private final ConcurrentSkipListMap<Long, Long> recoveringContainerMap =
- new ConcurrentSkipListMap<>();
+
+ private final ConcurrentSkipListSet<RecoveringContainer>
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<Map.Entry<Long, Long>> it = getRecoveringContainerIterator();
+ Iterator<RecoveringContainer> it = getRecoveringContainerIterator();
while (it.hasNext()) {
- Map.Entry<Long, Long> 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<Container<?>> iterator() {
/**
* Return an container Iterator over
- * {@link ContainerSet#recoveringContainerMap}.
- * @return {@literal Iterator<Container<?>>}
+ * {@link ContainerSet#recoveringContainerSet}.
+ * @return {@literal Iterator<RecoveringContainer>}
*/
- public Iterator<Map.Entry<Long, Long>> getRecoveringContainerIterator() {
- return recoveringContainerMap.entrySet().iterator();
+ public Iterator<RecoveringContainer> getRecoveringContainerIterator() {
+ return recoveringContainerSet.iterator();
}
/**
@@ -667,4 +669,52 @@ public <T> void buildMissingContainerSetAndValidate(Map<T,
Long> container2BCSID
}
});
}
+
+ /**
+ * A class that holds information about a recovering container.
+ */
+ public static class RecoveringContainer
+ implements Comparable<RecoveringContainer> {
+ 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 5535c5128cc..9c535e5f6e9 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<Map.Entry<Long, Long>> it =
+ Iterator<RecoveringContainer> it =
containerSet.getRecoveringContainerIterator();
while (it.hasNext()) {
- Map.Entry<Long, Long> 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;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]