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) {