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(); }
