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

zyk pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 755c5d55e24 PBTree: fix deed lock in ReleaseFlushMonitor and Scheduler 
(#11849)
755c5d55e24 is described below

commit 755c5d55e245a0d3e97d1293257a67212dab7d80
Author: Chen YZ <[email protected]>
AuthorDate: Fri Jan 5 09:28:23 2024 +0800

    PBTree: fix deed lock in ReleaseFlushMonitor and Scheduler (#11849)
---
 .../mtree/impl/pbtree/CachedMTreeStore.java        | 13 ++--
 .../mtree/impl/pbtree/flush/Scheduler.java         | 74 +++++++++++++--------
 .../impl/pbtree/memory/ReleaseFlushMonitor.java    | 77 ++++++++--------------
 .../db/metadata/mtree/schemafile/MonitorTest.java  | 13 ++--
 4 files changed, 82 insertions(+), 95 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/CachedMTreeStore.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/CachedMTreeStore.java
index 0d68da924d3..d4d2e73e46c 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/CachedMTreeStore.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/CachedMTreeStore.java
@@ -605,15 +605,10 @@ public class CachedMTreeStore implements 
IMTreeStore<ICachedMNode> {
    * @return should not continue releasing
    */
   public boolean executeMemoryRelease(AtomicLong releaseNodeNum, AtomicLong 
releaseMemorySize) {
-    lockManager.globalReadLock(true);
-    try {
-      if (regionStatistics.getUnpinnedMemorySize() != 0) {
-        return !memoryManager.evict(releaseNodeNum, releaseMemorySize);
-      } else {
-        return true;
-      }
-    } finally {
-      lockManager.globalReadUnlock();
+    if (regionStatistics.getUnpinnedMemorySize() != 0) {
+      return !memoryManager.evict(releaseNodeNum, releaseMemorySize);
+    } else {
+      return true;
     }
   }
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/flush/Scheduler.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/flush/Scheduler.java
index 6552cf1dff5..38282d12c2a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/flush/Scheduler.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/flush/Scheduler.java
@@ -28,7 +28,6 @@ import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.lock.Lock
 import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.memcontrol.IReleaseFlushStrategy;
 import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.memory.IMemoryManager;
 import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.schemafile.ISchemaFile;
-import org.apache.iotdb.tsfile.utils.Pair;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -96,8 +95,12 @@ public class Scheduler {
                     entry ->
                         CompletableFuture.runAsync(
                             () -> {
-                              CachedMTreeStore store = entry.getValue();
                               int regionId = entry.getKey();
+                              CachedMTreeStore store = entry.getValue();
+                              if (store == null) {
+                                // store has been closed
+                                return;
+                              }
                               IMemoryManager memoryManager = 
store.getMemoryManager();
                               ISchemaFile file = store.getSchemaFile();
                               LockManager lockManager = store.getLockManager();
@@ -106,8 +109,8 @@ public class Scheduler {
                               AtomicLong flushMemSize = new AtomicLong(0);
                               try {
                                 lockManager.globalReadLock();
-                                if (file == null) {
-                                  // store has been closed
+                                if (!regionToStore.containsKey(regionId)) {
+                                  // double check store have not been closed
                                   return;
                                 }
                                 PBTreeFlushExecutor flushExecutor =
@@ -151,25 +154,38 @@ public class Scheduler {
    */
   public synchronized void scheduleRelease(boolean force) {
     CompletableFuture.allOf(
-            regionToStore.values().stream()
+            regionToStore.entrySet().stream()
                 .map(
-                    store ->
+                    entry ->
                         CompletableFuture.runAsync(
                             () -> {
-                              AtomicLong releaseNodeNum = new AtomicLong(0);
-                              AtomicLong releaseMemorySize = new AtomicLong(0);
-                              long startTime = System.currentTimeMillis();
-                              while (force || 
releaseFlushStrategy.isExceedReleaseThreshold()) {
-                                // store try to release memory if not exceed 
release threshold
-                                if (store.executeMemoryRelease(releaseNodeNum, 
releaseMemorySize)) {
-                                  // if store can not release memory, break
-                                  break;
+                              int regionId = entry.getKey();
+                              CachedMTreeStore store = entry.getValue();
+                              if (!regionToStore.containsKey(regionId)) {
+                                // double check store have not been closed
+                                return;
+                              }
+                              LockManager lockManager = store.getLockManager();
+                              try {
+                                lockManager.globalReadLock(true);
+                                AtomicLong releaseNodeNum = new AtomicLong(0);
+                                AtomicLong releaseMemorySize = new 
AtomicLong(0);
+                                long startTime = System.currentTimeMillis();
+                                while (force || 
releaseFlushStrategy.isExceedReleaseThreshold()) {
+                                  // store try to release memory if not exceed 
release threshold
+                                  if (store.executeMemoryRelease(
+                                      releaseNodeNum, releaseMemorySize)) {
+                                    // if store can not release memory, break
+                                    break;
+                                  }
                                 }
+                                store.recordReleaseMetrics(
+                                    System.currentTimeMillis() - startTime,
+                                    releaseNodeNum.get(),
+                                    releaseMemorySize.get());
+                              } finally {
+                                lockManager.globalReadUnlock();
                               }
-                              store.recordReleaseMetrics(
-                                  System.currentTimeMillis() - startTime,
-                                  releaseNodeNum.get(),
-                                  releaseMemorySize.get());
                             },
                             workerPool))
                 .toArray(CompletableFuture[]::new))
@@ -185,19 +201,20 @@ public class Scheduler {
    * @param regionIds determine the MTreeStore to select subtrees, the head of 
the list is the first
    *     MTreeStore to select subtrees
    */
-  public synchronized void scheduleFlush(List<Pair<Integer, Long>> regionIds) {
+  public synchronized void scheduleFlush(List<Integer> regionIds) {
     AtomicInteger remainToFlush = new AtomicInteger(BATCH_FLUSH_SUBTREE);
-    boolean hasBreak = false;
-    for (Pair<Integer, Long> pair : regionIds) {
-      int regionId = pair.getLeft();
-      if (hasBreak || flushingRegionSet.contains(regionId)) {
-        
regionToStore.get(regionId).getLockManager().globalStampedReadUnlock(pair.getRight());
+    for (int regionId : regionIds) {
+      if (flushingRegionSet.contains(regionId)) {
         continue;
       }
       flushingRegionSet.add(regionId);
       workerPool.submit(
           () -> {
             CachedMTreeStore store = regionToStore.get(regionId);
+            if (store == null) {
+              // store has been closed
+              return;
+            }
             IMemoryManager memoryManager = store.getMemoryManager();
             ISchemaFile file = store.getSchemaFile();
             LockManager lockManager = store.getLockManager();
@@ -205,8 +222,9 @@ public class Scheduler {
             AtomicLong flushNodeNum = new AtomicLong(0);
             AtomicLong flushMemSize = new AtomicLong(0);
             try {
-              if (file == null) {
-                // store has been closed
+              lockManager.globalReadLock();
+              if (!regionToStore.containsKey(regionId)) {
+                // double check store have not been closed
                 return;
               }
               PBTreeFlushExecutor flushExecutor =
@@ -228,12 +246,12 @@ public class Scheduler {
                 LOGGER.debug("It takes {}ms to flush MTree in SchemaRegion 
{}", time, regionId);
               }
               store.recordFlushMetrics(time, flushNodeNum.get(), 
flushMemSize.get());
-              lockManager.globalStampedReadUnlock(pair.getRight());
+              lockManager.globalReadUnlock();
               flushingRegionSet.remove(regionId);
             }
           });
       if (remainToFlush.get() <= 0) {
-        hasBreak = true;
+        break;
       }
     }
   }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/memory/ReleaseFlushMonitor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/memory/ReleaseFlushMonitor.java
index a5ba3232e92..5f584a9061e 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/memory/ReleaseFlushMonitor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/schemaregion/mtree/impl/pbtree/memory/ReleaseFlushMonitor.java
@@ -28,7 +28,6 @@ import 
org.apache.iotdb.db.schemaengine.rescon.CachedSchemaEngineStatistics;
 import org.apache.iotdb.db.schemaengine.rescon.ISchemaEngineStatistics;
 import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.CachedMTreeStore;
 import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.flush.Scheduler;
-import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.lock.LockManager;
 import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.memcontrol.IReleaseFlushStrategy;
 import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.memcontrol.ReleaseFlushStrategyNumBasedImpl;
 import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.memcontrol.ReleaseFlushStrategySizeBasedImpl;
@@ -42,7 +41,6 @@ import javax.annotation.concurrent.NotThreadSafe;
 
 import java.util.ArrayList;
 import java.util.Comparator;
-import java.util.HashMap;
 import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
@@ -93,6 +91,7 @@ public class ReleaseFlushMonitor {
 
   public void clearCachedMTreeStore(CachedMTreeStore store) {
     regionToStoreMap.remove(store.getRegionStatistics().getSchemaRegionId());
+    
regionToTraverserTime.remove(store.getRegionStatistics().getSchemaRegionId());
   }
 
   public void init(ISchemaEngineStatistics engineStatistics) {
@@ -190,66 +189,42 @@ public class ReleaseFlushMonitor {
     regionToTraverserTime.computeIfAbsent(regionId, k -> new RecordList());
   }
 
-  public List<Pair<Integer, Long>> getRegionsToFlush(long windowsEndTime) {
+  public List<Integer> getRegionsToFlush(long windowsEndTime) {
     long windowsStartTime = windowsEndTime - MONITOR_INETRVAL_MILLISECONDS;
     List<Pair<Integer, Long>> regionAndFreeTimeList = new ArrayList<>();
-    Map<Integer, Long> storeToLockStamp = new HashMap<>();
     for (Map.Entry<Integer, RecordList> entry : 
regionToTraverserTime.entrySet()) {
       int regionId = entry.getKey();
-      CachedMTreeStore store = regionToStoreMap.get(regionId);
-      if (store == null) {
-        // already been removed
-        continue;
-      }
-
-      LockManager lockManager = store.getLockManager();
-      long lockStamp = lockManager.globalStampedReadLock();
-      boolean needReleaseLock = true;
-      try {
-        if (!regionToStoreMap.containsKey(regionId)) {
-          // already been removed
-          continue;
-        }
 
-        long traverserEndTime = windowsStartTime;
-        long traverserFreeTime = 0;
-        RecordList recordList = entry.getValue();
-        Iterator<RecordNode> iterator = recordList.iterator();
-        while (iterator.hasNext()) {
-          RecordNode recordNode = iterator.next();
-          if (recordNode.startTime > windowsEndTime) {
-            break;
-          }
-          if (recordNode.startTime > traverserEndTime) {
-            traverserFreeTime += (recordNode.startTime - traverserEndTime);
-            traverserEndTime = recordNode.endTime;
-          } else if (recordNode.endTime > traverserEndTime) {
-            traverserEndTime = recordNode.endTime;
-          }
-          if (recordNode.endTime < windowsStartTime) {
-            iterator.remove();
-          } else if (recordNode.endTime >= windowsEndTime) {
-            break;
-          }
-        }
-        if (traverserEndTime < windowsEndTime) {
-          traverserFreeTime += (windowsEndTime - traverserEndTime);
+      long traverserEndTime = windowsStartTime;
+      long traverserFreeTime = 0;
+      RecordList recordList = entry.getValue();
+      Iterator<RecordNode> iterator = recordList.iterator();
+      while (iterator.hasNext()) {
+        RecordNode recordNode = iterator.next();
+        if (recordNode.startTime > windowsEndTime) {
+          break;
         }
-        if (traverserFreeTime > FREE_FLUSH_PROPORTION * 
MONITOR_INETRVAL_MILLISECONDS) {
-          regionAndFreeTimeList.add(new Pair<>(regionId, traverserFreeTime));
-          storeToLockStamp.put(regionId, lockStamp);
-          needReleaseLock = false;
+        if (recordNode.startTime > traverserEndTime) {
+          traverserFreeTime += (recordNode.startTime - traverserEndTime);
+          traverserEndTime = recordNode.endTime;
+        } else if (recordNode.endTime > traverserEndTime) {
+          traverserEndTime = recordNode.endTime;
         }
-      } finally {
-        if (needReleaseLock) {
-          lockManager.globalStampedReadUnlock(lockStamp);
+        if (recordNode.endTime < windowsStartTime) {
+          iterator.remove();
+        } else if (recordNode.endTime >= windowsEndTime) {
+          break;
         }
       }
+      if (traverserEndTime < windowsEndTime) {
+        traverserFreeTime += (windowsEndTime - traverserEndTime);
+      }
+      if (traverserFreeTime > FREE_FLUSH_PROPORTION * 
MONITOR_INETRVAL_MILLISECONDS) {
+        regionAndFreeTimeList.add(new Pair<>(regionId, traverserFreeTime));
+      }
     }
     regionAndFreeTimeList.sort(Comparator.comparing((Pair<Integer, Long> o) -> 
o.right).reversed());
-    return regionAndFreeTimeList.stream()
-        .map(o -> new Pair<>(o.getLeft(), storeToLockStamp.get(o.getLeft())))
-        .collect(Collectors.toList());
+    return 
regionAndFreeTimeList.stream().map(Pair::getLeft).collect(Collectors.toList());
   }
 
   @TestOnly
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/metadata/mtree/schemafile/MonitorTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/metadata/mtree/schemafile/MonitorTest.java
index fd8ef2fd21f..c376355b8ad 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/metadata/mtree/schemafile/MonitorTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/metadata/mtree/schemafile/MonitorTest.java
@@ -23,7 +23,6 @@ import 
org.apache.iotdb.db.schemaengine.rescon.CachedSchemaRegionStatistics;
 import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.CachedMTreeStore;
 import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.lock.LockManager;
 import 
org.apache.iotdb.db.schemaengine.schemaregion.mtree.impl.pbtree.memory.ReleaseFlushMonitor;
-import org.apache.iotdb.tsfile.utils.Pair;
 
 import org.junit.After;
 import org.junit.Assert;
@@ -78,11 +77,11 @@ public class MonitorTest {
     setRecord(4, Arrays.asList(700L, 800L, 2500L), Arrays.asList(1000L, 1500L, 
5000L));
     // free =2100
     setRecord(5, Arrays.asList(0L, 2000L), Arrays.asList(1000L, 3900L));
-    List<Pair<Integer, Long>> regions = 
releaseFlushMonitor.getRegionsToFlush(5000);
+    List<Integer> regions = releaseFlushMonitor.getRegionsToFlush(5000);
     Assert.assertEquals(3, regions.size());
-    Assert.assertEquals(5, regions.get(0).left.intValue());
-    Assert.assertEquals(2, regions.get(1).left.intValue());
-    Assert.assertEquals(4, regions.get(2).left.intValue());
+    Assert.assertEquals(5, regions.get(0).intValue());
+    Assert.assertEquals(2, regions.get(1).intValue());
+    Assert.assertEquals(4, regions.get(2).intValue());
   }
 
   @Test
@@ -90,9 +89,9 @@ public class MonitorTest {
     mockCachedMTreeStore(2);
     setRecord(1, Arrays.asList(0L, 2000L), Arrays.asList(100L, 7000L));
     setRecord(2, Collections.singletonList(3000L), 
Collections.singletonList(3500L));
-    List<Pair<Integer, Long>> regions = 
releaseFlushMonitor.getRegionsToFlush(7000);
+    List<Integer> regions = releaseFlushMonitor.getRegionsToFlush(7000);
     Assert.assertEquals(1, regions.size());
-    Assert.assertEquals(2, regions.get(0).left.intValue());
+    Assert.assertEquals(2, regions.get(0).intValue());
   }
 
   private void setRecord(int regionId, List<Long> startTimes, List<Long> 
eneTimes) {

Reply via email to