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

qiaojialin pushed a commit to branch rel/0.11
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/rel/0.11 by this push:
     new 8377aa0  [IOTDB-1084] [To rel/0.11] Fix temporary memory of flushing 
may cause OOM (#2357)
8377aa0 is described below

commit 8377aa04d15e803431514c1cd47c375e57e6f950
Author: Haonan <[email protected]>
AuthorDate: Wed Dec 30 23:40:13 2020 +0800

    [IOTDB-1084] [To rel/0.11] Fix temporary memory of flushing may cause OOM 
(#2357)
---
 .../resources/conf/iotdb-engine.properties         |   6 +
 .../java/org/apache/iotdb/db/conf/IoTDBConfig.java |  26 ++++
 .../org/apache/iotdb/db/conf/IoTDBDescriptor.java  |   8 ++
 .../iotdb/db/engine/flush/MemTableFlushTask.java   | 140 +++++++++++----------
 .../iotdb/tsfile/write/chunk/ChunkWriterImpl.java  |   6 +-
 .../iotdb/tsfile/write/chunk/IChunkWriter.java     |   5 +
 6 files changed, 125 insertions(+), 66 deletions(-)

diff --git a/server/src/assembly/resources/conf/iotdb-engine.properties 
b/server/src/assembly/resources/conf/iotdb-engine.properties
index 3cf1302..931ee25 100644
--- a/server/src/assembly/resources/conf/iotdb-engine.properties
+++ b/server/src/assembly/resources/conf/iotdb-engine.properties
@@ -267,6 +267,12 @@ max_waiting_time_when_insert_blocked=10000
 # estimated metadata size (in byte) of one timeseries in Mtree
 estimated_series_size=300
 
+# size of encodingTaskQueue. The default value is 2147483647
+encoding_task_queue_size_for_flushing=2147483647
+
+# size of ioTaskQueue. The default value is 2147483647
+io_task_queue_size_for_flushing=2147483647
+
 ####################
 ### Upgrade Configurations
 ####################
diff --git a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java 
b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
index 8c5b7e7..4a6fec5 100644
--- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
+++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
@@ -764,6 +764,16 @@ public class IoTDBConfig {
    */
   private boolean debugState = false;
 
+  /**
+   * the size of ioTaskQueue
+   */
+  private int ioTaskQueueSizeForFlushing = Integer.MAX_VALUE;
+
+  /**
+   * the size of encodingTaskQueue
+   */
+  private int encodingTaskQueueSizeForFlushing = Integer.MAX_VALUE;
+
   public IoTDBConfig() {
     // empty constructor
   }
@@ -2058,4 +2068,20 @@ public class IoTDBConfig {
   public void setDebugState(boolean debugState) {
     this.debugState = debugState;
   }
+
+  public int getEncodingTaskQueueSizeForFlushing() {
+    return encodingTaskQueueSizeForFlushing;
+  }
+
+  public void setEncodingTaskQueueSizeForFlushing(int 
encodingTaskQueueSizeForFlushing) {
+    this.encodingTaskQueueSizeForFlushing = encodingTaskQueueSizeForFlushing;
+  }
+
+  public int getIoTaskQueueSizeForFlushing() {
+    return ioTaskQueueSizeForFlushing;
+  }
+
+  public void setIoTaskQueueSizeForFlushing(int ioTaskQueueSizeForFlushing) {
+    this.ioTaskQueueSizeForFlushing = ioTaskQueueSizeForFlushing;
+  }
 }
diff --git a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java 
b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
index dc95ef1..41d3af1 100644
--- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
+++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
@@ -295,6 +295,14 @@ public class IoTDBDescriptor {
           .getProperty("estimated_series_size",
               Integer.toString(conf.getEstimatedSeriesSize()))));
 
+      conf.setIoTaskQueueSizeForFlushing(Integer.parseInt(properties
+          .getProperty("io_task_queue_size_for_flushing",
+              Integer.toString(conf.getIoTaskQueueSizeForFlushing()))));
+
+      conf.setEncodingTaskQueueSizeForFlushing(Integer.parseInt(properties
+          .getProperty("encoding_task_queue_size_for_flushing",
+              Integer.toString(conf.getEncodingTaskQueueSizeForFlushing()))));
+
       conf.setMergeChunkPointNumberThreshold(Integer.parseInt(properties
           .getProperty("merge_chunk_point_number",
               Integer.toString(conf.getMergeChunkPointNumberThreshold()))));
diff --git 
a/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java 
b/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java
index e9f4c92..ce08298 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/engine/flush/MemTableFlushTask.java
@@ -19,10 +19,12 @@
 package org.apache.iotdb.db.engine.flush;
 
 import java.io.IOException;
-import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.concurrent.LinkedBlockingQueue;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.Future;
-import java.util.concurrent.TimeUnit;
+
+import org.apache.iotdb.db.conf.IoTDBConfig;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.engine.flush.pool.FlushSubTaskPoolManager;
 import org.apache.iotdb.db.engine.memtable.IMemTable;
 import org.apache.iotdb.db.engine.memtable.IWritableMemChunk;
@@ -42,18 +44,19 @@ public class MemTableFlushTask {
   private static final Logger logger = 
LoggerFactory.getLogger(MemTableFlushTask.class);
   private static final FlushSubTaskPoolManager subTaskPoolManager = 
FlushSubTaskPoolManager
       .getInstance();
+  private static IoTDBConfig config = 
IoTDBDescriptor.getInstance().getConfig();
   private final Future<?> encodingTaskFuture;
   private final Future<?> ioTaskFuture;
   private RestorableTsFileIOWriter writer;
 
-  private final ConcurrentLinkedQueue<Object> ioTaskQueue = new 
ConcurrentLinkedQueue<>();
-  private final ConcurrentLinkedQueue<Object> encodingTaskQueue = new 
ConcurrentLinkedQueue<>();
+  private final LinkedBlockingQueue<Object> ioTaskQueue =
+      new LinkedBlockingQueue<>(config.getIoTaskQueueSizeForFlushing());
+  private final LinkedBlockingQueue<Object> encodingTaskQueue = 
+      new LinkedBlockingQueue<>(config.getEncodingTaskQueueSizeForFlushing());
   private String storageGroup;
 
   private IMemTable memTable;
 
-  private volatile boolean noMoreEncodingTask = false;
-  private volatile boolean noMoreIOTask = false;
 
   /**
    * @param memTable the memTable to flush
@@ -84,18 +87,18 @@ public class MemTableFlushTask {
     long sortTime = 0;
 
     for (String deviceId : memTable.getMemTableMap().keySet()) {
-      encodingTaskQueue.add(new StartFlushGroupIOTask(deviceId));
+      encodingTaskQueue.put(new StartFlushGroupIOTask(deviceId));
       for (String measurementId : 
memTable.getMemTableMap().get(deviceId).keySet()) {
         long startTime = System.currentTimeMillis();
         IWritableMemChunk series = 
memTable.getMemTableMap().get(deviceId).get(measurementId);
         MeasurementSchema desc = series.getSchema();
         TVList tvList = series.getSortedTVList();
         sortTime += System.currentTimeMillis() - startTime;
-        encodingTaskQueue.add(new Pair<>(tvList, desc));
+        encodingTaskQueue.put(new Pair<>(tvList, desc));
       }
-      encodingTaskQueue.add(new EndChunkGroupIoTask());
+      encodingTaskQueue.put(new EndChunkGroupIoTask());
     }
-    noMoreEncodingTask = true;
+    encodingTaskQueue.put(new TaskEnd());
     logger.debug(
         "Storage group {} memtable {}, flushing into disk: data sort time cost 
{} ms.",
         storageGroup, memTable.getVersion(), sortTime);
@@ -103,8 +106,6 @@ public class MemTableFlushTask {
     try {
       encodingTaskFuture.get();
     } catch (InterruptedException | ExecutionException e) {
-      // avoid ioTask waiting forever
-      noMoreIOTask = true;
       ioTaskFuture.cancel(true);
       throw e;
     }
@@ -164,40 +165,50 @@ public class MemTableFlushTask {
     @Override
     public void run() {
       long memSerializeTime = 0;
-      boolean noMoreMessages = false;
       logger.debug("Storage group {} memtable {}, starts to encoding data.", 
storageGroup,
           memTable.getVersion());
       while (true) {
-        if (noMoreEncodingTask) {
-          noMoreMessages = true;
+
+        Object task = null;
+        try {
+          task = encodingTaskQueue.take();
+        } catch (InterruptedException e1) {
+          logger.error("Take task into ioTaskQueue Interrupted");
+          Thread.currentThread().interrupt();
+          break;
         }
-        Object task = encodingTaskQueue.poll();
-        if (task == null) {
-          if (noMoreMessages) {
-            break;
-          }
+        if (task instanceof StartFlushGroupIOTask || task instanceof 
EndChunkGroupIoTask) {
           try {
-            TimeUnit.MILLISECONDS.sleep(10);
-          } catch (@SuppressWarnings("squid:S2142") InterruptedException e) {
-            logger.error("Storage group {} memtable {}, encoding task is 
interrupted.",
-                storageGroup, memTable.getVersion(), e);
-            // generally it is because the thread pool is shutdown so the task 
should be aborted
+            ioTaskQueue.put(task);
+          } catch (InterruptedException e) {
+            logger.error("Put task into ioTaskQueue Interrupted");
+            Thread.currentThread().interrupt();
             break;
           }
+        } else if (task instanceof TaskEnd) {
+          break;
         } else {
-          if (task instanceof StartFlushGroupIOTask || task instanceof 
EndChunkGroupIoTask) {
-            ioTaskQueue.add(task);
-          } else {
-            long starTime = System.currentTimeMillis();
-            Pair<TVList, MeasurementSchema> encodingMessage = (Pair<TVList, 
MeasurementSchema>) task;
-            IChunkWriter seriesWriter = new 
ChunkWriterImpl(encodingMessage.right);
-            writeOneSeries(encodingMessage.left, seriesWriter, 
encodingMessage.right.getType());
-            ioTaskQueue.add(seriesWriter);
-            memSerializeTime += System.currentTimeMillis() - starTime;
+          long starTime = System.currentTimeMillis();
+          Pair<TVList, MeasurementSchema> encodingMessage = (Pair<TVList, 
MeasurementSchema>) task;
+          IChunkWriter seriesWriter = new 
ChunkWriterImpl(encodingMessage.right);
+          writeOneSeries(encodingMessage.left, seriesWriter, 
encodingMessage.right.getType());
+          seriesWriter.sealCurrentPage();
+          seriesWriter.clearPageWriter();
+          try {
+            ioTaskQueue.put(seriesWriter);
+          } catch (InterruptedException e) {
+            logger.error("Put task into ioTaskQueue Interrupted");
+            Thread.currentThread().interrupt();
           }
+          memSerializeTime += System.currentTimeMillis() - starTime;
         }
       }
-      noMoreIOTask = true;
+      try {
+        ioTaskQueue.put(new TaskEnd());
+      } catch (InterruptedException e) {
+        logger.error("Put task into ioTaskQueue Interrupted");
+        Thread.currentThread().interrupt();
+      }
       logger.debug("Storage group {}, flushing memtable {} into disk: Encoding 
data cost "
               + "{} ms.",
           storageGroup, memTable.getVersion(), memSerializeTime);
@@ -207,47 +218,46 @@ public class MemTableFlushTask {
   @SuppressWarnings("squid:S135")
   private Runnable ioTask = () -> {
     long ioTime = 0;
-    boolean returnWhenNoTask = false;
     logger.debug("Storage group {} memtable {}, start io.", storageGroup, 
memTable.getVersion());
     while (true) {
-      if (noMoreIOTask) {
-        returnWhenNoTask = true;
+      Object ioMessage = null;
+      try {
+        ioMessage = ioTaskQueue.take();
+      } catch (InterruptedException e1) {
+        logger.error("take task from ioTaskQueue Interrupted");
+        Thread.currentThread().interrupt();
+        break;
       }
-      Object ioMessage = ioTaskQueue.poll();
-      if (ioMessage == null) {
-        if (returnWhenNoTask) {
+      long starTime = System.currentTimeMillis();
+      try {
+        if (ioMessage instanceof StartFlushGroupIOTask) {
+          this.writer.startChunkGroup(((StartFlushGroupIOTask) 
ioMessage).deviceId);
+        } else if (ioMessage instanceof TaskEnd) {
           break;
         }
-        try {
-          TimeUnit.MILLISECONDS.sleep(10);
-        } catch (@SuppressWarnings("squid:S2142") InterruptedException e) {
-          logger.error("Storage group {} memtable {}, io task is 
interrupted.", storageGroup
-              , memTable.getVersion(), e);
-          // generally it is because the thread pool is shutdown so the task 
should be aborted
-          break;
-        }
-      } else {
-        long starTime = System.currentTimeMillis();
-        try {
-          if (ioMessage instanceof StartFlushGroupIOTask) {
-            this.writer.startChunkGroup(((StartFlushGroupIOTask) 
ioMessage).deviceId);
-          } else if (ioMessage instanceof IChunkWriter) {
-            ChunkWriterImpl chunkWriter = (ChunkWriterImpl) ioMessage;
-            chunkWriter.writeToFileWriter(this.writer);
-          } else {
-            this.writer.endChunkGroup();
-          }
-        } catch (IOException e) {
-          logger.error("Storage group {} memtable {}, io task meets error.", 
storageGroup,
-              memTable.getVersion(), e);
-          throw new FlushRunTimeException(e);
+        else if (ioMessage instanceof IChunkWriter) {
+          ChunkWriterImpl chunkWriter = (ChunkWriterImpl) ioMessage;
+          chunkWriter.writeToFileWriter(this.writer);
+        } else {
+          this.writer.endChunkGroup();
         }
-        ioTime += System.currentTimeMillis() - starTime;
+      } catch (IOException e) {
+        logger.error("Storage group {} memtable {}, io task meets error.", 
storageGroup,
+            memTable.getVersion(), e);
+        throw new FlushRunTimeException(e);
       }
+      ioTime += System.currentTimeMillis() - starTime;
     }
     logger.debug("flushing a memtable {} in storage group {}, io cost {}ms", 
memTable.getVersion(),
         storageGroup, ioTime);
   };
+  
+  static class TaskEnd {
+    
+    TaskEnd() {
+      
+    }
+  }
 
   static class EndChunkGroupIoTask {
 
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ChunkWriterImpl.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ChunkWriterImpl.java
index 01a5f80..48ab591 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ChunkWriterImpl.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/ChunkWriterImpl.java
@@ -240,10 +240,14 @@ public class ChunkWriterImpl implements IChunkWriter {
 
   @Override
   public void sealCurrentPage() {
-    if (pageWriter.getPointNumber() > 0) {
+    if (pageWriter != null && pageWriter.getPointNumber() > 0) {
       writePageToPageBuffer();
     }
   }
+  
+  public void clearPageWriter() {
+    pageWriter = null;
+  }
 
   @Override
   public int getNumOfPages() {
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/IChunkWriter.java 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/IChunkWriter.java
index 6f1ab46..e1dcc3d 100644
--- a/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/IChunkWriter.java
+++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/write/chunk/IChunkWriter.java
@@ -115,6 +115,11 @@ public interface IChunkWriter {
    * seal the current page which may has not enough data points in force.
    */
   void sealCurrentPage();
+  
+  /**
+   * set the current pageWriter to null, friendly for gc
+   */
+  void clearPageWriter();
 
   int getNumOfPages();
 

Reply via email to