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

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

commit 42eb8fb808d5f1cae0004c89d9492cf0230a5a83
Author: HTHou <[email protected]>
AuthorDate: Thu Apr 15 14:53:33 2021 +0800

    [IOTDB-1306] DeadLock in MemControl module
---
 .../db/engine/storagegroup/StorageGroupProcessor.java   | 13 +++++++++++++
 .../iotdb/db/engine/storagegroup/TsFileProcessor.java   | 17 ++++++++---------
 .../java/org/apache/iotdb/db/rescon/SystemInfo.java     |  3 +++
 3 files changed, 24 insertions(+), 9 deletions(-)

diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
index 5ddf1e3..8672d86 100755
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
@@ -106,6 +106,7 @@ import java.util.concurrent.RejectedExecutionException;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.locks.Condition;
 import java.util.concurrent.locks.ReadWriteLock;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 
@@ -160,6 +161,9 @@ public class StorageGroupProcessor {
    * partitionLatestFlushedTimeForEachDevice)
    */
   private final ReadWriteLock insertLock = new ReentrantReadWriteLock();
+
+  private final Condition rejectCondition = 
insertLock.writeLock().newCondition();
+
   /** closeStorageGroupCondition is used to wait for all currently closing 
TsFiles to be done. */
   private final Object closeStorageGroupCondition = new Object();
   /**
@@ -979,6 +983,7 @@ public class StorageGroupProcessor {
 
     try {
       tsFileProcessor.insertTablet(insertTabletPlan, start, end, results);
+      writeLock();
     } catch (WriteProcessRejectException e) {
       logger.warn("insert to TsFileProcessor rejected, {}", e.getMessage());
       return false;
@@ -1585,6 +1590,14 @@ public class StorageGroupProcessor {
     insertLock.writeLock().unlock();
   }
 
+  public void rejectConditionAwait() throws InterruptedException {
+    rejectCondition.await();
+  }
+
+  public void rejectConditionSignal() {
+    rejectCondition.signal();
+  }
+
   /**
    * @param tsFileResources includes sealed and unsealed tsfile resources
    * @return fill unsealed tsfile resources with memory data and 
ChunkMetadataList of data in disk
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
index c20ec7e..73c9e2e 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/TsFileProcessor.java
@@ -22,7 +22,6 @@ import org.apache.iotdb.db.conf.IoTDBConfig;
 import org.apache.iotdb.db.conf.IoTDBConstant;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.conf.adapter.CompressionRatio;
-import org.apache.iotdb.db.engine.StorageEngine;
 import org.apache.iotdb.db.engine.flush.CloseFileListener;
 import org.apache.iotdb.db.engine.flush.FlushListener;
 import org.apache.iotdb.db.engine.flush.FlushManager;
@@ -37,7 +36,6 @@ import 
org.apache.iotdb.db.engine.querycontext.ReadOnlyMemChunk;
 import 
org.apache.iotdb.db.engine.storagegroup.StorageGroupProcessor.UpdateEndTimeCallBack;
 import org.apache.iotdb.db.exception.TsFileProcessorException;
 import org.apache.iotdb.db.exception.WriteProcessException;
-import org.apache.iotdb.db.exception.WriteProcessRejectException;
 import org.apache.iotdb.db.exception.metadata.MetadataException;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
 import org.apache.iotdb.db.metadata.PartialPath;
@@ -387,14 +385,15 @@ public class TsFileProcessor {
     tsFileProcessorInfo.addTSPMemCost(unsealedResourceIncrement + 
chunkMetadataIncrement);
     if (storageGroupInfo.needToReportToSystem()) {
       SystemInfo.getInstance().reportStorageGroupStatus(storageGroupInfo);
-      try {
-        StorageEngine.blockInsertionIfReject();
-      } catch (WriteProcessRejectException e) {
-        storageGroupInfo.releaseStorageGroupMemCost(memTableIncrement);
-        tsFileProcessorInfo.releaseTSPMemCost(unsealedResourceIncrement + 
chunkMetadataIncrement);
-        SystemInfo.getInstance().resetStorageGroupStatus(storageGroupInfo, 
false);
-        throw e;
+      long startTime = System.currentTimeMillis();
+      while (SystemInfo.getInstance().isRejected()) {
+        try {
+          storageGroupInfo.getStorageGroupProcessor().rejectConditionAwait();
+        } catch (InterruptedException e) {
+          Thread.currentThread().interrupt();
+        }
       }
+      logger.debug("Time for Waiting memory release is {}", 
System.currentTimeMillis() - startTime);
     }
     workMemTable.addTVListRamCost(memTableIncrement);
     workMemTable.addTextDataSize(textDataIncrement);
diff --git a/server/src/main/java/org/apache/iotdb/db/rescon/SystemInfo.java 
b/server/src/main/java/org/apache/iotdb/db/rescon/SystemInfo.java
index 9e20d51..eede244 100644
--- a/server/src/main/java/org/apache/iotdb/db/rescon/SystemInfo.java
+++ b/server/src/main/java/org/apache/iotdb/db/rescon/SystemInfo.java
@@ -140,6 +140,9 @@ public class SystemInfo {
         rejected = false;
       }
     }
+    if (!rejected) {
+      storageGroupInfo.getStorageGroupProcessor().rejectConditionSignal();
+    }
     if (shouldInvokeFlush && needForceAsyncFlush) {
       forceAsyncFlush();
     }

Reply via email to