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]

Reply via email to