This is an automated email from the ASF dual-hosted git repository.
szetszwo 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 930f8fa4aa6 HDDS-15909. Refactor OzoneManagerLock. (#10852)
930f8fa4aa6 is described below
commit 930f8fa4aa62b7e6ceb256f586b836a5fe2ff090
Author: Tsz-Wo Nicholas Sze <[email protected]>
AuthorDate: Fri Jul 24 11:16:44 2026 -0700
HDDS-15909. Refactor OzoneManagerLock. (#10852)
---
.../apache/hadoop/hdds/utils/SimpleStriped.java | 4 +-
.../hadoop/hdds/utils/TestSimpleStriped.java | 3 +-
.../ozone/om/lock/DAGResourceLockTracker.java | 5 +
.../ozone/om/lock/LeveledResourceLockTracker.java | 6 +
.../hadoop/ozone/om/lock/OzoneManagerLock.java | 316 ++++++++++-----------
.../hadoop/ozone/om/lock/ResourceLockTracker.java | 2 +
.../hadoop/ozone/om/lock/TestKeyPathLock.java | 2 +-
7 files changed, 172 insertions(+), 166 deletions(-)
diff --git
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/SimpleStriped.java
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/SimpleStriped.java
index ec83553473e..390b11e7a18 100644
---
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/SimpleStriped.java
+++
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/SimpleStriped.java
@@ -18,7 +18,6 @@
package org.apache.hadoop.hdds.utils;
import com.google.common.util.concurrent.Striped;
-import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
/**
@@ -45,8 +44,7 @@ private SimpleStriped() {
* @param fair whether to use a fair ordering policy
* @return a new {@code Striped<ReadWriteLock>}
*/
- public static Striped<ReadWriteLock> readWriteLock(int stripes,
- boolean fair) {
+ public static Striped<ReentrantReadWriteLock> readWriteLock(int stripes,
boolean fair) {
return Striped.custom(stripes, () -> new ReentrantReadWriteLock(fair));
}
diff --git
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/TestSimpleStriped.java
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/TestSimpleStriped.java
index ccd80b9fd24..d1ce0476529 100644
---
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/TestSimpleStriped.java
+++
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/TestSimpleStriped.java
@@ -36,8 +36,7 @@ void testReadWriteLocks() {
}
private void testReadWriteLocks(boolean fair) {
- Striped<ReadWriteLock> striped = SimpleStriped.readWriteLock(128,
- fair);
+ Striped<ReentrantReadWriteLock> striped = SimpleStriped.readWriteLock(128,
fair);
assertEquals(128, striped.size());
ReadWriteLock lock = striped.get("key1");
assertEquals(fair, ((ReentrantReadWriteLock) lock).isFair());
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/DAGResourceLockTracker.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/DAGResourceLockTracker.java
index 7fd44059dd6..a669eec517d 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/DAGResourceLockTracker.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/DAGResourceLockTracker.java
@@ -72,6 +72,11 @@ public static DAGResourceLockTracker get() {
return instance;
}
+ @Override
+ Class<DAGLeveledResource> getResourceClass() {
+ return DAGLeveledResource.class;
+ }
+
/**
* Performs a Depth-First Search (DFS) traversal on a directed acyclic graph
(DAG)
* composed of {@code DAGLeveledResource} objects. This method populates a
mapping
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/LeveledResourceLockTracker.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/LeveledResourceLockTracker.java
index bbe9cd9076c..783652a56a0 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/LeveledResourceLockTracker.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/LeveledResourceLockTracker.java
@@ -19,6 +19,7 @@
import java.util.Arrays;
import java.util.stream.Stream;
+import org.apache.hadoop.ozone.om.lock.OzoneManagerLock.LeveledResource;
/**
* The LeveledResourceLockTracker class is a singleton that extends the
@@ -57,6 +58,11 @@ final class LeveledResourceLockTracker extends
ResourceLockTracker<OzoneManagerL
private LeveledResourceLockTracker() {
}
+ @Override
+ Class<LeveledResource> getResourceClass() {
+ return LeveledResource.class;
+ }
+
public static LeveledResourceLockTracker get() {
if (instance == null) {
synchronized (LeveledResourceLockTracker.class) {
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/OzoneManagerLock.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/OzoneManagerLock.java
index f567f17766b..b0abd85f944 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/OzoneManagerLock.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/OzoneManagerLock.java
@@ -17,14 +17,12 @@
package org.apache.hadoop.ozone.om.lock;
-import static org.apache.hadoop.hdds.utils.CompositeKey.combineKeys;
import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_MANAGER_FAIR_LOCK;
import static
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_MANAGER_FAIR_LOCK_DEFAULT;
import static
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_MANAGER_STRIPED_LOCK_SIZE_DEFAULT;
import static
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_MANAGER_STRIPED_LOCK_SIZE_PREFIX;
import com.google.common.annotations.VisibleForTesting;
-import com.google.common.collect.ImmutableMap;
import com.google.common.util.concurrent.Striped;
import java.util.ArrayList;
import java.util.Arrays;
@@ -35,19 +33,18 @@
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.TimeUnit;
-import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.function.Function;
import java.util.stream.Collectors;
import java.util.stream.IntStream;
import java.util.stream.StreamSupport;
-import org.apache.commons.lang3.tuple.Pair;
import org.apache.hadoop.hdds.conf.ConfigurationSource;
import org.apache.hadoop.hdds.utils.CompositeKey;
import org.apache.hadoop.hdds.utils.SimpleStriped;
import org.apache.hadoop.ipc_.ProcessingDetails.Timing;
import org.apache.hadoop.ipc_.Server;
import org.apache.hadoop.util.Time;
+import org.apache.ratis.util.Preconditions;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -97,41 +94,140 @@ public class OzoneManagerLock implements IOzoneManagerLock
{
private static final Logger LOG =
LoggerFactory.getLogger(OzoneManagerLock.class);
- private final Map<Class<? extends Resource>,
- Pair<Map<Resource, Striped<ReadWriteLock>>, ResourceLockTracker>>
resourcelockMap;
+ private final ResourceLocks<LeveledResource> leveledResourceLocks;
+ private final ResourceLocks<DAGLeveledResource> dagLeveledResourceLocks;
- private OMLockMetrics omLockMetrics;
+ private final OMLockMetrics omLockMetrics = OMLockMetrics.create();
+
+ class ResourceLocks<R extends Resource> {
+ private final Map<R, Striped<ReentrantReadWriteLock>> lockMap;
+ private final ResourceLockTracker<R> tracker;
+
+ ResourceLocks(Map<R, Striped<ReentrantReadWriteLock>> lockMap,
ResourceLockTracker<R> tracker) {
+ this.lockMap = lockMap;
+ this.tracker = tracker;
+ }
+
+ R assertAcquire(Resource resource) {
+ final R r = Preconditions.assertInstanceOf(resource,
tracker.getResourceClass());
+ tracker.clearLockDetails();
+ if (!tracker.canLockResource(r)) {
+ final String errorMessage = "Thread '" +
Thread.currentThread().getName() + "' cannot acquire "
+ + r.getName() + " lock while holding " + getCurrentLocks() + "
lock(s).";
+ LOG.error(errorMessage);
+ // TODO: change it to IllegalStateException
+ throw new RuntimeException(errorMessage);
+ }
+ return r;
+ }
+
+ private ReentrantReadWriteLock getLockForTesting(Resource resource,
String... keys) {
+ final R r = Preconditions.assertInstanceOf(resource,
tracker.getResourceClass());
+ return getLock(r, keys);
+ }
+
+ private ReentrantReadWriteLock getLock(R r, String... keys) {
+ return lockMap.get(r).get(CompositeKey.combineKeys(keys));
+ }
+
+ private void acquireLock(R resource, boolean isRead,
ReentrantReadWriteLock lock, long startWaitingTimeNanos) {
+ if (isRead) {
+ lock.readLock().lock();
+ updateReadLockMetrics(resource, tracker, lock, startWaitingTimeNanos);
+ } else {
+ lock.writeLock().lock();
+ updateWriteLockMetrics(resource, tracker, lock, startWaitingTimeNanos);
+ }
+ }
+
+ OMLockDetails acquire(Resource resource, boolean isRead,
+ Function<Striped<ReentrantReadWriteLock>,
Iterable<ReentrantReadWriteLock>> getLocks) {
+ final R r = assertAcquire(resource);
+ final long startWaitingTimeNanos = Time.monotonicNowNanos();
+ for (ReentrantReadWriteLock lock : getLocks.apply(lockMap.get(r))) {
+ acquireLock(r, isRead, lock, startWaitingTimeNanos);
+ }
+ return tracker.lockResource(r);
+ }
+
+ OMLockDetails acquire(Resource resource, boolean isRead, String... keys) {
+ final R r = assertAcquire(resource);
+ final long startWaitingTimeNanos = Time.monotonicNowNanos();
+ acquireLock(r, isRead, getLock(r, keys), startWaitingTimeNanos);
+ return tracker.lockResource(r);
+ }
+
+ void releaseLock(R resource, boolean isRead, ReentrantReadWriteLock lock) {
+ if (isRead) {
+ lock.readLock().unlock();
+ updateReadUnlockMetrics(resource, tracker, lock);
+ } else {
+ boolean isWriteLocked = lock.isWriteLockedByCurrentThread();
+ lock.writeLock().unlock();
+ updateWriteUnlockMetrics(resource, tracker, lock, isWriteLocked);
+ }
+ }
+
+ OMLockDetails release(Resource resource, boolean isRead, String... keys) {
+ final R r = Preconditions.assertInstanceOf(resource,
tracker.getResourceClass());
+ tracker.clearLockDetails();
+ final ReentrantReadWriteLock lock = getLock(r, keys);
+ releaseLock(r, isRead, lock);
+ return tracker.unlockResource(r);
+ }
+
+ private OMLockDetails release(Resource resource, boolean isRead,
+ Function<Striped<ReentrantReadWriteLock>,
Iterable<ReentrantReadWriteLock>> getLock) {
+ final R r = Preconditions.assertInstanceOf(resource,
tracker.getResourceClass());
+ tracker.clearLockDetails();
+ final Iterable<ReentrantReadWriteLock> i = getLock.apply(lockMap.get(r));
+ final List<ReentrantReadWriteLock> locks =
StreamSupport.stream(i.spliterator(), false)
+ .collect(Collectors.toList());
+ // Release locks in reverse order.
+ Collections.reverse(locks);
+ for (ReentrantReadWriteLock lock : locks) {
+ releaseLock(r, isRead, lock);
+ }
+ return tracker.unlockResource(r);
+ }
+
+ List<String> getCurrentLocks() {
+ return tracker.getCurrentLockedResources()
+ .map(Resource::getName)
+ .collect(Collectors.toList());
+ }
+ }
/**
* Creates new OzoneManagerLock instance.
* @param conf Configuration object
*/
public OzoneManagerLock(ConfigurationSource conf) {
- omLockMetrics = OMLockMetrics.create();
- this.resourcelockMap = ImmutableMap.of(LeveledResource.class,
getLeveledLocks(conf), DAGLeveledResource.class,
- getFlatLocks(conf));
+ this.leveledResourceLocks =
newResourceLocks(LeveledResourceLockTracker.get(), conf);
+ this.dagLeveledResourceLocks =
newResourceLocks(DAGResourceLockTracker.get(), conf);
}
- private Pair<Map<Resource, Striped<ReadWriteLock>>, ResourceLockTracker>
getLeveledLocks(
- ConfigurationSource conf) {
- Map<LeveledResource, Striped<ReadWriteLock>> stripedLockMap = new
EnumMap<>(LeveledResource.class);
- for (LeveledResource r : LeveledResource.values()) {
+ private <T extends Enum<T> & Resource> ResourceLocks<T> newResourceLocks(
+ ResourceLockTracker<T> tracker, ConfigurationSource conf) {
+ final Class<T> clazz = tracker.getResourceClass();
+ final EnumMap<T, Striped<ReentrantReadWriteLock>> stripedLockMap = new
EnumMap<>(clazz);
+ for (T r : clazz.getEnumConstants()) {
stripedLockMap.put(r, createStripeLock(r, conf));
}
- return Pair.of(Collections.unmodifiableMap(stripedLockMap),
LeveledResourceLockTracker.get());
+ return new ResourceLocks<>(Collections.unmodifiableMap(stripedLockMap),
tracker);
}
- private Pair<Map<Resource, Striped<ReadWriteLock>>, ResourceLockTracker>
getFlatLocks(
- ConfigurationSource conf) {
- Map<DAGLeveledResource, Striped<ReadWriteLock>> stripedLockMap = new
EnumMap<>(DAGLeveledResource.class);
- for (DAGLeveledResource r : DAGLeveledResource.values()) {
- stripedLockMap.put(r, createStripeLock(r, conf));
+ private ResourceLocks<?> getResourceLocks(Resource instance) {
+ final Class<?> clazz = instance.getClass();
+ if (clazz == LeveledResource.class) {
+ return leveledResourceLocks;
+ } else if (clazz == DAGLeveledResource.class) {
+ return dagLeveledResourceLocks;
}
- return Pair.of(Collections.unmodifiableMap(stripedLockMap),
DAGResourceLockTracker.get());
+ throw new IllegalArgumentException("Unsupported resource class: " + clazz);
}
- private Striped<ReadWriteLock> createStripeLock(Resource r,
- ConfigurationSource conf) {
+ private static Striped<ReentrantReadWriteLock> createStripeLock(Resource r,
ConfigurationSource conf) {
boolean fair = conf.getBoolean(OZONE_MANAGER_FAIR_LOCK,
OZONE_MANAGER_FAIR_LOCK_DEFAULT);
String stripeSizeKey = OZONE_MANAGER_STRIPED_LOCK_SIZE_PREFIX +
@@ -141,11 +237,12 @@ private Striped<ReadWriteLock> createStripeLock(Resource
r,
return SimpleStriped.readWriteLock(size, fair);
}
- private Iterable<ReadWriteLock> getAllLocks(Striped<ReadWriteLock> striped) {
+ private Iterable<ReentrantReadWriteLock>
getAllLocks(Striped<ReentrantReadWriteLock> striped) {
return IntStream.range(0,
striped.size()).mapToObj(striped::getAt).collect(Collectors.toList());
}
- private Iterable<ReadWriteLock> bulkGetLock(Striped<ReadWriteLock> striped,
Collection<String[]> keys) {
+ private Iterable<ReentrantReadWriteLock>
bulkGetLock(Striped<ReentrantReadWriteLock> striped,
+ Collection<String[]> keys) {
List<Object> lockKeys = new ArrayList<>(keys.size());
for (String[] key : keys) {
if (Objects.nonNull(key)) {
@@ -155,13 +252,6 @@ private Iterable<ReadWriteLock>
bulkGetLock(Striped<ReadWriteLock> striped, Coll
return striped.bulkGet(lockKeys);
}
- private ReentrantReadWriteLock getLock(Map<Resource, Striped<ReadWriteLock>>
lockMap, Resource resource,
- String... keys) {
- Striped<ReadWriteLock> striped = lockMap.get(resource);
- Object key = combineKeys(keys);
- return (ReentrantReadWriteLock) striped.get(key);
- }
-
/**
* Acquire read lock on resource.
*
@@ -181,7 +271,8 @@ private ReentrantReadWriteLock getLock(Map<Resource,
Striped<ReadWriteLock>> loc
*/
@Override
public OMLockDetails acquireReadLock(Resource resource, String... keys) {
- return acquireLock(resource, true, keys);
+ return getResourceLocks(resource)
+ .acquire(resource, true, keys);
}
/**
@@ -203,7 +294,8 @@ public OMLockDetails acquireReadLock(Resource resource,
String... keys) {
*/
@Override
public OMLockDetails acquireReadLocks(Resource resource,
Collection<String[]> keys) {
- return acquireLocks(resource, true, striped -> bulkGetLock(striped, keys));
+ return getResourceLocks(resource)
+ .acquire(resource, true, striped -> bulkGetLock(striped, keys));
}
/**
@@ -225,7 +317,8 @@ public OMLockDetails acquireReadLocks(Resource resource,
Collection<String[]> ke
*/
@Override
public OMLockDetails acquireWriteLock(Resource resource, String... keys) {
- return acquireLock(resource, false, keys);
+ return getResourceLocks(resource)
+ .acquire(resource, false, keys);
}
/**
@@ -247,7 +340,8 @@ public OMLockDetails acquireWriteLock(Resource resource,
String... keys) {
*/
@Override
public OMLockDetails acquireWriteLocks(Resource resource,
Collection<String[]> keys) {
- return acquireLocks(resource, false, striped -> bulkGetLock(striped,
keys));
+ return getResourceLocks(resource)
+ .acquire(resource, false, striped -> bulkGetLock(striped, keys));
}
/**
@@ -257,59 +351,11 @@ public OMLockDetails acquireWriteLocks(Resource resource,
Collection<String[]> k
*/
@Override
public OMLockDetails acquireResourceWriteLock(Resource resource) {
- return acquireLocks(resource, false, this::getAllLocks);
- }
-
- private void acquireLock(Resource resource, boolean isReadLock,
ReadWriteLock lock,
- long startWaitingTimeNanos) {
- if (isReadLock) {
- lock.readLock().lock();
- updateReadLockMetrics(resource, (ReentrantReadWriteLock) lock,
startWaitingTimeNanos);
- } else {
- lock.writeLock().lock();
- updateWriteLockMetrics(resource, (ReentrantReadWriteLock) lock,
startWaitingTimeNanos);
- }
- }
-
- private OMLockDetails acquireLocks(Resource resource, boolean isReadLock,
- Function<Striped<ReadWriteLock>, Iterable<ReadWriteLock>>
lockListProvider) {
- Pair<Map<Resource, Striped<ReadWriteLock>>, ResourceLockTracker>
resourceLockPair =
- resourcelockMap.get(resource.getClass());
- ResourceLockTracker<Resource> resourceLockTracker =
resourceLockPair.getRight();
- resourceLockTracker.clearLockDetails();
- if (!resourceLockTracker.canLockResource(resource)) {
- String errorMessage = getErrorMessage(resource);
- LOG.error(errorMessage);
- throw new RuntimeException(errorMessage);
- }
-
- long startWaitingTimeNanos = Time.monotonicNowNanos();
-
- for (ReadWriteLock lock :
lockListProvider.apply(resourceLockPair.getKey().get(resource))) {
- acquireLock(resource, isReadLock, lock, startWaitingTimeNanos);
- }
- return resourceLockTracker.lockResource(resource);
- }
-
- private OMLockDetails acquireLock(Resource resource, boolean isReadLock,
String... keys) {
- Pair<Map<Resource, Striped<ReadWriteLock>>, ResourceLockTracker>
resourceLockPair =
- resourcelockMap.get(resource.getClass());
- ResourceLockTracker<Resource> resourceLockTracker =
resourceLockPair.getRight();
- resourceLockTracker.clearLockDetails();
- if (!resourceLockTracker.canLockResource(resource)) {
- String errorMessage = getErrorMessage(resource);
- LOG.error(errorMessage);
- throw new RuntimeException(errorMessage);
- }
-
- long startWaitingTimeNanos = Time.monotonicNowNanos();
-
- ReentrantReadWriteLock lock = getLock(resourceLockPair.getKey(), resource,
keys);
- acquireLock(resource, isReadLock, lock, startWaitingTimeNanos);
- return resourceLockTracker.lockResource(resource);
+ return getResourceLocks(resource)
+ .acquire(resource, false, this::getAllLocks);
}
- private void updateReadLockMetrics(Resource resource,
+ private void updateReadLockMetrics(Resource resource, ResourceLockTracker<?
extends Resource> tracker,
ReentrantReadWriteLock lock, long startWaitingTimeNanos) {
/*
@@ -323,14 +369,13 @@ private void updateReadLockMetrics(Resource resource,
// Adds a snapshot to the metric readLockWaitingTimeMsStat.
omLockMetrics.setReadLockWaitingTimeMsStat(
TimeUnit.NANOSECONDS.toMillis(readLockWaitingTimeNanos));
-
updateProcessingDetails(resourcelockMap.get(resource.getClass()).getValue(),
- Timing.LOCKWAIT, readLockWaitingTimeNanos);
+ updateProcessingDetails(tracker, Timing.LOCKWAIT,
readLockWaitingTimeNanos);
resource.getResourceManager().setStartReadHeldTimeNanos(Time.monotonicNowNanos());
}
}
- private void updateWriteLockMetrics(Resource resource,
+ private void updateWriteLockMetrics(Resource resource, ResourceLockTracker<?
extends Resource> tracker,
ReentrantReadWriteLock lock, long startWaitingTimeNanos) {
/*
* writeHoldCount helps in metrics updation only once in case
@@ -345,25 +390,15 @@ private void updateWriteLockMetrics(Resource resource,
// Adds a snapshot to the metric writeLockWaitingTimeMsStat.
omLockMetrics.setWriteLockWaitingTimeMsStat(
TimeUnit.NANOSECONDS.toMillis(writeLockWaitingTimeNanos));
-
updateProcessingDetails(resourcelockMap.get(resource.getClass()).getValue(),
Timing.LOCKWAIT,
- writeLockWaitingTimeNanos);
+ updateProcessingDetails(tracker, Timing.LOCKWAIT,
writeLockWaitingTimeNanos);
resource.getResourceManager().setStartWriteHeldTimeNanos(Time.monotonicNowNanos());
}
}
- private String getErrorMessage(Resource resource) {
- return "Thread '" + Thread.currentThread().getName() + "' cannot " +
- "acquire " + resource.getName() + " lock while holding " +
- getCurrentLocks().toString() + " lock(s).";
- }
-
@VisibleForTesting
- List<String> getCurrentLocks() {
- return resourcelockMap.values().stream().map(Pair::getValue)
- .flatMap(rlm -> ((ResourceLockTracker<? extends
Resource>)rlm).getCurrentLockedResources())
- .map(Resource::getName)
- .collect(Collectors.toList());
+ int getCurrentLockSizeForTesting() {
+ return leveledResourceLocks.getCurrentLocks().size() +
dagLeveledResourceLocks.getCurrentLocks().size();
}
/**
@@ -397,7 +432,8 @@ public void releaseMultiUserLock(String firstUser, String
secondUser) {
*/
@Override
public OMLockDetails releaseWriteLock(Resource resource, String... keys) {
- return releaseLock(resource, false, keys);
+ return getResourceLocks(resource)
+ .release(resource, false, keys);
}
/**
@@ -410,7 +446,8 @@ public OMLockDetails releaseWriteLock(Resource resource,
String... keys) {
*/
@Override
public OMLockDetails releaseWriteLocks(Resource resource,
Collection<String[]> keys) {
- return releaseLocks(resource, false, striped -> bulkGetLock(striped,
keys));
+ return getResourceLocks(resource)
+ .release(resource, false, striped -> bulkGetLock(striped, keys));
}
/**
@@ -420,7 +457,8 @@ public OMLockDetails releaseWriteLocks(Resource resource,
Collection<String[]> k
*/
@Override
public OMLockDetails releaseResourceWriteLock(Resource resource) {
- return releaseLocks(resource, false, this::getAllLocks);
+ return getResourceLocks(resource)
+ .release(resource, false, this::getAllLocks);
}
/**
@@ -433,7 +471,8 @@ public OMLockDetails releaseResourceWriteLock(Resource
resource) {
*/
@Override
public OMLockDetails releaseReadLock(Resource resource, String... keys) {
- return releaseLock(resource, true, keys);
+ return getResourceLocks(resource)
+ .release(resource, true, keys);
}
/**
@@ -446,51 +485,11 @@ public OMLockDetails releaseReadLock(Resource resource,
String... keys) {
*/
@Override
public OMLockDetails releaseReadLocks(Resource resource,
Collection<String[]> keys) {
- return releaseLocks(resource, true, striped -> bulkGetLock(striped, keys));
- }
-
- private OMLockDetails releaseLock(Resource resource, boolean isReadLock,
- String... keys) {
- Pair<Map<Resource, Striped<ReadWriteLock>>, ResourceLockTracker>
resourceLockPair =
- resourcelockMap.get(resource.getClass());
- ResourceLockTracker<Resource> resourceLockTracker =
resourceLockPair.getRight();
- resourceLockTracker.clearLockDetails();
- ReentrantReadWriteLock lock = getLock(resourceLockPair.getKey(), resource,
keys);
- if (isReadLock) {
- lock.readLock().unlock();
- updateReadUnlockMetrics(resource, lock);
- } else {
- boolean isWriteLocked = lock.isWriteLockedByCurrentThread();
- lock.writeLock().unlock();
- updateWriteUnlockMetrics(resource, lock, isWriteLocked);
- }
- return resourceLockTracker.unlockResource(resource);
- }
-
- private OMLockDetails releaseLocks(Resource resource, boolean isReadLock,
- Function<Striped<ReadWriteLock>, Iterable<ReadWriteLock>>
lockListProvider) {
- Pair<Map<Resource, Striped<ReadWriteLock>>, ResourceLockTracker>
resourceLockPair =
- resourcelockMap.get(resource.getClass());
- ResourceLockTracker<Resource> resourceLockTracker =
resourceLockPair.getRight();
- resourceLockTracker.clearLockDetails();
- List<ReadWriteLock> locks =
StreamSupport.stream(lockListProvider.apply(resourceLockPair.getKey().get(resource))
- .spliterator(), false).collect(Collectors.toList());
- // Release locks in reverse order.
- Collections.reverse(locks);
- for (ReadWriteLock lock : locks) {
- if (isReadLock) {
- lock.readLock().unlock();
- updateReadUnlockMetrics(resource, (ReentrantReadWriteLock) lock);
- } else {
- boolean isWriteLocked =
((ReentrantReadWriteLock)lock).isWriteLockedByCurrentThread();
- lock.writeLock().unlock();
- updateWriteUnlockMetrics(resource, (ReentrantReadWriteLock) lock,
isWriteLocked);
- }
- }
- return resourceLockTracker.unlockResource(resource);
+ return getResourceLocks(resource)
+ .release(resource, true, striped -> bulkGetLock(striped, keys));
}
- private void updateReadUnlockMetrics(Resource resource,
+ private void updateReadUnlockMetrics(Resource resource,
ResourceLockTracker<? extends Resource> tracker,
ReentrantReadWriteLock lock) {
/*
* readHoldCount helps in metrics updation only once in case
@@ -503,12 +502,11 @@ private void updateReadUnlockMetrics(Resource resource,
// Adds a snapshot to the metric readLockHeldTimeMsStat.
omLockMetrics.setReadLockHeldTimeMsStat(
TimeUnit.NANOSECONDS.toMillis(readLockHeldTimeNanos));
-
updateProcessingDetails(resourcelockMap.get(resource.getClass()).getValue(),
Timing.LOCKSHARED,
- readLockHeldTimeNanos);
+ updateProcessingDetails(tracker, Timing.LOCKSHARED,
readLockHeldTimeNanos);
}
}
- private void updateWriteUnlockMetrics(Resource resource,
+ private void updateWriteUnlockMetrics(Resource resource,
ResourceLockTracker<? extends Resource> tracker,
ReentrantReadWriteLock lock, boolean isWriteLocked) {
/*
* writeHoldCount helps in metrics updation only once in case
@@ -522,8 +520,7 @@ private void updateWriteUnlockMetrics(Resource resource,
// Adds a snapshot to the metric writeLockHeldTimeMsStat.
omLockMetrics.setWriteLockHeldTimeMsStat(
TimeUnit.NANOSECONDS.toMillis(writeLockHeldTimeNanos));
-
updateProcessingDetails(resourcelockMap.get(resource.getClass()).getValue(),
Timing.LOCKEXCLUSIVE,
- writeLockHeldTimeNanos);
+ updateProcessingDetails(tracker, Timing.LOCKEXCLUSIVE,
writeLockHeldTimeNanos);
}
}
@@ -535,7 +532,7 @@ private void updateWriteUnlockMetrics(Resource resource,
@Override
@VisibleForTesting
public int getReadHoldCount(Resource resource, String... keys) {
- return getLock(resourcelockMap.get(resource.getClass()).getKey(),
resource, keys).getReadHoldCount();
+ return getResourceLocks(resource).getLockForTesting(resource,
keys).getReadHoldCount();
}
@@ -547,7 +544,7 @@ public int getReadHoldCount(Resource resource, String...
keys) {
@Override
@VisibleForTesting
public int getWriteHoldCount(Resource resource, String... keys) {
- return getLock(resourcelockMap.get(resource.getClass()).getKey(),
resource, keys).getWriteHoldCount();
+ return getResourceLocks(resource).getLockForTesting(resource,
keys).getWriteHoldCount();
}
/**
@@ -559,9 +556,8 @@ public int getWriteHoldCount(Resource resource, String...
keys) {
*/
@Override
@VisibleForTesting
- public boolean isWriteLockedByCurrentThread(Resource resource,
- String... keys) {
- return getLock(resourcelockMap.get(resource.getClass()).getKey(),
resource, keys).isWriteLockedByCurrentThread();
+ public boolean isWriteLockedByCurrentThread(Resource resource, String...
keys) {
+ return getResourceLocks(resource).getLockForTesting(resource,
keys).isWriteLockedByCurrentThread();
}
/**
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/ResourceLockTracker.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/ResourceLockTracker.java
index 8a551e5f06d..80e40711383 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/ResourceLockTracker.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/lock/ResourceLockTracker.java
@@ -31,6 +31,8 @@ abstract class ResourceLockTracker<T extends
IOzoneManagerLock.Resource> {
private final ThreadLocal<OMLockDetails> omLockDetails =
ThreadLocal.withInitial(OMLockDetails::new);
+ abstract Class<T> getResourceClass();
+
abstract boolean canLockResource(T resource);
abstract Stream<T> getCurrentLockedResources();
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/lock/TestKeyPathLock.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/lock/TestKeyPathLock.java
index 3122f65a0d4..77b7999d616 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/lock/TestKeyPathLock.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/lock/TestKeyPathLock.java
@@ -216,7 +216,7 @@ private void testDiffKeyPathWriteLockMultiThreadingUtil(
// Waiting for all the threads to be instantiated/to reach
// acquireWriteLock.
countDown.countDown();
- assertEquals(1, lock.getCurrentLocks().size());
+ assertEquals(1, lock.getCurrentLockSizeForTesting());
lock.releaseWriteLock(resource, sampleResourceName);
LOG.info("Write Lock Released by " + Thread.currentThread().getName());
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]