This is an automated email from the ASF dual-hosted git repository.

jianyun pushed a commit to branch rocksdb/dev
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 9640cf17566e9301b6f6d7a18bb09006ad2f9564
Author: lisijia <[email protected]>
AuthorDate: Wed Mar 9 15:25:30 2022 +0800

    fix timeout issue when obtaining lock
---
 .../iotdb/db/metadata/rocksdb/MRocksDBManager.java | 39 +++++++++++-----------
 1 file changed, 20 insertions(+), 19 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBManager.java
 
b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBManager.java
index a1f5aa9..7b0b5f1 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBManager.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/metadata/rocksdb/MRocksDBManager.java
@@ -85,7 +85,7 @@ import 
org.apache.iotdb.tsfile.write.schema.IMeasurementSchema;
 import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
 import org.apache.iotdb.tsfile.write.schema.TimeseriesSchema;
 
-import com.google.common.util.concurrent.Striped;
+import com.google.common.collect.MapMaker;
 import org.apache.commons.lang3.ArrayUtils;
 import org.apache.commons.lang3.StringUtils;
 import org.rocksdb.Holder;
@@ -114,6 +114,7 @@ import java.util.concurrent.ConcurrentLinkedQueue;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
 import java.util.function.BiFunction;
 import java.util.function.Function;
 import java.util.stream.Collectors;
@@ -182,8 +183,8 @@ public class MRocksDBManager implements IMetaManager {
 
   private RocksDBReadWriteHandler readWriteHandler;
 
-  // TODO: check how Stripped Lock consume memory
-  private Striped<Lock> locksPool = Striped.lazyWeakLock(10000);
+  private final Map<String, ReentrantLock> locksPool =
+      new MapMaker().weakValues().initialCapacity(10000).makeMap();
 
   private volatile Map<String, Boolean> storageGroupDeletingFlagMap = new 
ConcurrentHashMap<>();
 
@@ -385,7 +386,7 @@ public class MRocksDBManager implements IMetaManager {
       return;
     }
     String levelPath = RocksDBUtils.getLevelPath(nodes, start - 1);
-    Lock lock = locksPool.get(levelPath);
+    Lock lock = locksPool.computeIfAbsent(levelPath, x -> new ReentrantLock());
     if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
       lockedLocks.push(lock);
       try {
@@ -468,7 +469,7 @@ public class MRocksDBManager implements IMetaManager {
       aliasNodes[nodes.length - 1] = alias;
       String aliasLevelPath = RocksDBUtils.getLevelPath(aliasNodes, 
aliasNodes.length - 1);
       byte[] aliasNodeKey = RocksDBUtils.toAliasNodeKey(aliasLevelPath);
-      Lock lock = locksPool.get(aliasLevelPath);
+      Lock lock = locksPool.computeIfAbsent(aliasLevelPath, x -> new 
ReentrantLock());
       if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
         try {
           if 
(!readWriteHandler.keyExistByAllTypes(aliasLevelPath).existAnyKey()) {
@@ -552,7 +553,7 @@ public class MRocksDBManager implements IMetaManager {
       Stack<Lock> acquiredLock = new Stack<>();
       try {
         for (String lockKey : locks) {
-          Lock lock = locksPool.get(lockKey);
+          Lock lock = locksPool.computeIfAbsent(lockKey, x -> new 
ReentrantLock());
           if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
             acquiredLock.push(lock);
             if (readWriteHandler.keyExistByAllTypes(lockKey).existAnyKey()) {
@@ -591,7 +592,7 @@ public class MRocksDBManager implements IMetaManager {
       return;
     }
     String levelPath = RocksDBUtils.getLevelPath(nodes, start - 1);
-    Lock lock = locksPool.get(levelPath);
+    Lock lock = locksPool.computeIfAbsent(levelPath, x -> new ReentrantLock());
     if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
       try {
         CheckKeyResult checkResult = 
readWriteHandler.keyExistByAllTypes(levelPath);
@@ -657,7 +658,7 @@ public class MRocksDBManager implements IMetaManager {
 
         // Delete measurement node
         String mLevelPath = RocksDBUtils.getLevelPath(p.getNodes(), 
p.getNodeLength() - 1);
-        Lock lock = locksPool.get(mLevelPath);
+        Lock lock = locksPool.computeIfAbsent(mLevelPath, x -> new 
ReentrantLock());
         RMeasurementMNode deletedNode;
         if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
           try {
@@ -687,7 +688,7 @@ public class MRocksDBManager implements IMetaManager {
           PartialPath curPath = curNode.getPartialPath();
           String curLevelPath =
               RocksDBUtils.getLevelPath(curPath.getNodes(), 
curPath.getNodeLength() - 1);
-          Lock curLock = locksPool.get(curLevelPath);
+          Lock curLock = locksPool.computeIfAbsent(curLevelPath, x -> new 
ReentrantLock());
           if (curLock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
             try {
               IMNode toDelete = curNode.getParent();
@@ -799,7 +800,7 @@ public class MRocksDBManager implements IMetaManager {
       int len = nodes.length;
       for (int i = 1; i < nodes.length; i++) {
         String levelKey = RocksDBUtils.getLevelPath(nodes, i);
-        Lock lock = locksPool.get(levelKey);
+        Lock lock = locksPool.computeIfAbsent(levelKey, x -> new 
ReentrantLock());
         if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
           try {
             CheckKeyResult keyCheckResult = 
readWriteHandler.keyExistByAllTypes(levelKey);
@@ -886,7 +887,7 @@ public class MRocksDBManager implements IMetaManager {
     byte[] pathKey = RocksDBUtils.toStorageNodeKey(levelPath);
     Holder<byte[]> holder = new Holder<>();
     try {
-      Lock lock = locksPool.get(levelPath);
+      Lock lock = locksPool.computeIfAbsent(levelPath, x -> new 
ReentrantLock());
       if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
         try {
           if (readWriteHandler.keyExist(pathKey, holder)) {
@@ -1862,7 +1863,7 @@ public class MRocksDBManager implements IMetaManager {
     String levelPath = RocksDBUtils.getLevelPath(path.getNodes(), 
path.getNodeLength() - 1);
     byte[] originKey = RocksDBUtils.toMeasurementNodeKey(levelPath);
     try {
-      Lock rawKeyLock = locksPool.get(levelPath);
+      Lock rawKeyLock = locksPool.computeIfAbsent(levelPath, x -> new 
ReentrantLock());
       if (rawKeyLock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
         try {
           String[] nodes = path.getNodes();
@@ -1875,7 +1876,7 @@ public class MRocksDBManager implements IMetaManager {
             String newAliasLevel = RocksDBUtils.getLevelPath(newAlias, 
newAlias.length - 1);
             byte[] newAliasKey = RocksDBUtils.toAliasNodeKey(newAliasLevel);
 
-            Lock newAliasLock = locksPool.get(newAliasLevel);
+            Lock newAliasLock = locksPool.computeIfAbsent(newAliasLevel, x -> 
new ReentrantLock());
             Lock oldAliasLock = null;
             boolean lockedOldAlias = false;
             try {
@@ -1893,7 +1894,7 @@ public class MRocksDBManager implements IMetaManager {
                 oldAlias[nodes.length - 1] = mNode.getAlias();
                 String oldAliasLevel = RocksDBUtils.getLevelPath(oldAlias, 
oldAlias.length - 1);
                 byte[] oldAliasKey = 
RocksDBUtils.toAliasNodeKey(oldAliasLevel);
-                oldAliasLock = locksPool.get(oldAliasLevel);
+                oldAliasLock = locksPool.computeIfAbsent(oldAliasLevel, x -> 
new ReentrantLock());
                 if (oldAliasLock.tryLock(MAX_LOCK_WAIT_TIME, 
TimeUnit.MILLISECONDS)) {
                   lockedOldAlias = true;
                   if (!readWriteHandler.keyExist(oldAliasKey)) {
@@ -1962,7 +1963,7 @@ public class MRocksDBManager implements IMetaManager {
     byte[] key = RocksDBUtils.toMeasurementNodeKey(levelPath);
     Holder<byte[]> holder = new Holder<>();
     try {
-      Lock lock = locksPool.get(levelPath);
+      Lock lock = locksPool.computeIfAbsent(levelPath, x -> new 
ReentrantLock());
       if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
         try {
           if (!readWriteHandler.keyExist(key, holder)) {
@@ -2004,7 +2005,7 @@ public class MRocksDBManager implements IMetaManager {
     byte[] key = RocksDBUtils.toMeasurementNodeKey(levelPath);
     Holder<byte[]> holder = new Holder<>();
     try {
-      Lock lock = locksPool.get(levelPath);
+      Lock lock = locksPool.computeIfAbsent(levelPath, x -> new 
ReentrantLock());
       if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
         try {
           if (!readWriteHandler.keyExist(key, holder)) {
@@ -2052,7 +2053,7 @@ public class MRocksDBManager implements IMetaManager {
     byte[] key = RocksDBUtils.toMeasurementNodeKey(levelPath);
     Holder<byte[]> holder = new Holder<>();
     try {
-      Lock lock = locksPool.get(levelPath);
+      Lock lock = locksPool.computeIfAbsent(levelPath, x -> new 
ReentrantLock());
       if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
         try {
           if (!readWriteHandler.keyExist(key, holder)) {
@@ -2108,7 +2109,7 @@ public class MRocksDBManager implements IMetaManager {
     byte[] key = RocksDBUtils.toMeasurementNodeKey(levelPath);
     Holder<byte[]> holder = new Holder<>();
     try {
-      Lock lock = locksPool.get(levelPath);
+      Lock lock = locksPool.computeIfAbsent(levelPath, x -> new 
ReentrantLock());
       if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
         try {
           byte[] originValue = holder.getValue();
@@ -2159,7 +2160,7 @@ public class MRocksDBManager implements IMetaManager {
     byte[] nodeKey = RocksDBUtils.toMeasurementNodeKey(levelPath);
     Holder<byte[]> holder = new Holder<>();
     try {
-      Lock lock = locksPool.get(levelPath);
+      Lock lock = locksPool.computeIfAbsent(levelPath, x -> new 
ReentrantLock());
       if (lock.tryLock(MAX_LOCK_WAIT_TIME, TimeUnit.MILLISECONDS)) {
         try {
           if (!readWriteHandler.keyExist(nodeKey, holder)) {

Reply via email to