This is an automated email from the ASF dual-hosted git repository.
jackietien 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 a07b5ec [IOTDB-1084] Fix temporary memory of flushing may cause OOM
(#2358)
a07b5ec is described below
commit a07b5ec7a48131d2185577fe578bed4aeea71507
Author: Haonan <[email protected]>
AuthorDate: Tue Jan 19 21:04:02 2021 +0800
[IOTDB-1084] Fix temporary memory of flushing may cause OOM (#2358)
Fix temporary memory of flushing may cause OOM
---
.../resources/conf/iotdb-engine.properties | 3 +
.../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 13 ++
.../org/apache/iotdb/db/conf/IoTDBDescriptor.java | 4 +
.../iotdb/db/engine/flush/MemTableFlushTask.java | 169 ++++++++++++---------
.../org/apache/iotdb/db/rescon/SystemInfo.java | 28 +++-
.../iotdb/tsfile/write/chunk/ChunkWriterImpl.java | 6 +-
.../iotdb/tsfile/write/chunk/IChunkWriter.java | 5 +
7 files changed, 152 insertions(+), 76 deletions(-)
diff --git a/server/src/assembly/resources/conf/iotdb-engine.properties
b/server/src/assembly/resources/conf/iotdb-engine.properties
index e9296a6..1857dac 100644
--- a/server/src/assembly/resources/conf/iotdb-engine.properties
+++ b/server/src/assembly/resources/conf/iotdb-engine.properties
@@ -285,6 +285,9 @@ max_waiting_time_when_insert_blocked=10000
# estimated metadata size (in byte) of one timeseries in Mtree
estimated_series_size=300
+# size of ioTaskQueue. The default value is 10
+io_task_queue_size_for_flushing=10
+
####################
### 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 d03114d..9d0d828 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
@@ -846,6 +846,11 @@ public class IoTDBConfig {
private boolean debugState = false;
/**
+ * the size of ioTaskQueue
+ */
+ private int ioTaskQueueSizeForFlushing = 10;
+
+ /**
* the number of virtual storage groups per user-defined storage group
*/
private int virtualStorageGroupNum = 1;
@@ -2269,4 +2274,12 @@ public class IoTDBConfig {
public void setMlogBufferSize(int mlogBufferSize) {
this.mlogBufferSize = mlogBufferSize;
}
+
+ 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 a2f7d41..ef56f50 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
@@ -288,6 +288,10 @@ 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.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 916d2d9..c643f59 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,16 +19,19 @@
package org.apache.iotdb.db.engine.flush;
import java.io.IOException;
+import java.util.concurrent.LinkedBlockingQueue;
import java.util.Map;
-import java.util.concurrent.ConcurrentLinkedQueue;
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;
import org.apache.iotdb.db.exception.runtime.FlushRunTimeException;
import org.apache.iotdb.db.utils.datastructure.TVList;
+import org.apache.iotdb.db.rescon.SystemInfo;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import org.apache.iotdb.tsfile.utils.Pair;
import org.apache.iotdb.tsfile.write.chunk.ChunkWriterImpl;
@@ -43,18 +46,23 @@ public class MemTableFlushTask {
private static final Logger LOGGER =
LoggerFactory.getLogger(MemTableFlushTask.class);
private static final FlushSubTaskPoolManager SUB_TASK_POOL_MANAGER =
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> encodingTaskQueue = new
LinkedBlockingQueue<>();
+ private final LinkedBlockingQueue<Object> ioTaskQueue =
(config.isEnableMemControl()
+ && SystemInfo.getInstance().isEncodingFasterThanIo())
+ ? new LinkedBlockingQueue<>(config.getIoTaskQueueSizeForFlushing())
+ : new LinkedBlockingQueue<>();
+
private String storageGroup;
private IMemTable memTable;
- private volatile boolean noMoreEncodingTask = false;
- private volatile boolean noMoreIOTask = false;
+ private volatile long memSerializeTime = 0L;
+ private volatile long ioTime = 0L;
/**
* @param memTable the memTable to flush
@@ -81,12 +89,19 @@ public class MemTableFlushTask {
storageGroup,
memTable.memSize(),
memTable.getTotalPointsNum() / memTable.getSeriesNumber());
+
+ long estimatedTemporaryMemSize = 0L;
+ if (config.isEnableMemControl() &&
SystemInfo.getInstance().isEncodingFasterThanIo()) {
+ estimatedTemporaryMemSize = memTable.memSize() /
memTable.getSeriesNumber()
+ * config.getIoTaskQueueSizeForFlushing();
+
SystemInfo.getInstance().applyTemporaryMemoryForFlushing(estimatedTemporaryMemSize);
+ }
long start = System.currentTimeMillis();
long sortTime = 0;
//for map do not use get(key) to iteratate
for (Map.Entry<String, Map<String, IWritableMemChunk>> memTableEntry :
memTable.getMemTableMap().entrySet()) {
- encodingTaskQueue.add(new StartFlushGroupIOTask(memTableEntry.getKey()));
+ encodingTaskQueue.put(new StartFlushGroupIOTask(memTableEntry.getKey()));
final Map<String, IWritableMemChunk> value = memTableEntry.getValue();
for (Map.Entry<String, IWritableMemChunk> iWritableMemChunkEntry :
value.entrySet()) {
@@ -95,13 +110,12 @@ public class MemTableFlushTask {
MeasurementSchema desc = series.getSchema();
TVList tvList = series.getSortedTVListForFlush();
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 file {}: data sort time cost
{} ms.",
storageGroup, writer.getFile().getName(), sortTime);
@@ -109,8 +123,6 @@ public class MemTableFlushTask {
try {
encodingTaskFuture.get();
} catch (InterruptedException | ExecutionException e) {
- // avoid ioTask waiting forever
- noMoreIOTask = true;
ioTaskFuture.cancel(true);
throw e;
}
@@ -123,6 +135,13 @@ public class MemTableFlushTask {
throw new ExecutionException(e);
}
+ if (config.isEnableMemControl()) {
+ if (estimatedTemporaryMemSize != 0) {
+
SystemInfo.getInstance().releaseTemporaryMemoryForFlushing(estimatedTemporaryMemSize);
+ }
+ SystemInfo.getInstance().setEncodingFasterThanIo(ioTime >=
memSerializeTime);
+ }
+
LOGGER.info(
"Storage group {} memtable {} flushing a memtable has finished! Time
consumption: {}ms",
storageGroup, memTable, System.currentTimeMillis() - start);
@@ -169,95 +188,103 @@ public class MemTableFlushTask {
@SuppressWarnings("squid:S135")
@Override
public void run() {
- long memSerializeTime = 0;
- boolean noMoreMessages = false;
LOGGER.debug("Storage group {} memtable flushing to file {} starts to
encoding data.",
- storageGroup, writer.getFile().getName());
+ storageGroup, writer.getFile().getName());
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);
+ ioTaskQueue.put(task);
} catch (@SuppressWarnings("squid:S2142") InterruptedException e) {
LOGGER.error("Storage group {} memtable flushing to file {},
encoding task is interrupted.",
storageGroup, writer.getFile().getName(), e);
// generally it is because the thread pool is shutdown so the task
should be aborted
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;
- LOGGER.debug("Storage group {}, flushing memtable into file {}: Encoding
data cost "
- + "{} ms.",
+ 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, writer.getFile().getName(), memSerializeTime);
}
};
@SuppressWarnings("squid:S135")
private Runnable ioTask = () -> {
- long ioTime = 0;
- boolean returnWhenNoTask = false;
LOGGER.debug("Storage group {} memtable flushing to file {} start io.",
- storageGroup, writer.getFile().getName());
+ storageGroup, writer.getFile().getName());
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;
+ } else if (ioMessage instanceof IChunkWriter) {
+ ChunkWriterImpl chunkWriter = (ChunkWriterImpl) ioMessage;
+ chunkWriter.writeToFileWriter(this.writer);
+ } else {
+ this.writer.setMinPlanIndex(memTable.getMinPlanIndex());
+ this.writer.setMaxPlanIndex(memTable.getMaxPlanIndex());
+ this.writer.endChunkGroup();
}
- try {
- TimeUnit.MILLISECONDS.sleep(10);
- } catch (@SuppressWarnings("squid:S2142") InterruptedException e) {
- LOGGER.error("Storage group {} memtable flushing to file {}, io task
is interrupted.",
- storageGroup, writer.getFile().getName());
- // 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.setMinPlanIndex(memTable.getMinPlanIndex());
- this.writer.setMaxPlanIndex(memTable.getMaxPlanIndex());
- this.writer.endChunkGroup();
- }
- } catch (IOException e) {
- LOGGER.error("Storage group {} memtable flushing to file {}, io task
meets error.",
- storageGroup, writer.getFile().getName(), e);
- throw new FlushRunTimeException(e);
- }
- ioTime += System.currentTimeMillis() - starTime;
+ } catch (IOException e) {
+ LOGGER.error("Storage group {} memtable {}, io task meets error.",
storageGroup,
+ memTable, e);
+ throw new FlushRunTimeException(e);
}
+ ioTime += System.currentTimeMillis() - starTime;
}
LOGGER.debug("flushing a memtable to file {} in storage group {}, io cost
{}ms",
writer.getFile().getName(), storageGroup, ioTime);
};
+ static class TaskEnd {
+
+ TaskEnd() {
+
+ }
+ }
+
static class EndChunkGroupIoTask {
EndChunkGroupIoTask() {
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 cbf76bd..05238a5 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
@@ -41,13 +41,13 @@ public class SystemInfo {
private long totalSgMemCost = 0L;
private volatile boolean rejected = false;
+ private static long memorySizeForWrite = config.getAllocateMemoryForWrite();
private Map<StorageGroupInfo, Long> reportedSgMemCostMap = new HashMap<>();
- private static final double FLUSH_THERSHOLD =
- config.getAllocateMemoryForWrite() * config.getFlushProportion();
- private static final double REJECT_THERSHOLD =
- config.getAllocateMemoryForWrite() * config.getRejectProportion();
+ private static double FLUSH_THERSHOLD = memorySizeForWrite *
config.getFlushProportion();
+ private static double REJECT_THERSHOLD = memorySizeForWrite *
config.getRejectProportion();
+ private boolean isEncodingFasterThanIo = true;
/**
* Report current mem cost of storage group to system. Called when the
memory of
@@ -185,6 +185,14 @@ public class SystemInfo {
return rejected;
}
+ public void setEncodingFasterThanIo(boolean isEncodingFasterThanIo) {
+ this.isEncodingFasterThanIo = isEncodingFasterThanIo;
+ }
+
+ public boolean isEncodingFasterThanIo() {
+ return isEncodingFasterThanIo;
+ }
+
public void close() {
reportedSgMemCostMap.clear();
totalSgMemCost = 0;
@@ -202,4 +210,16 @@ public class SystemInfo {
private static SystemInfo instance = new SystemInfo();
}
+
+ public synchronized void applyTemporaryMemoryForFlushing(long
estimatedTemporaryMemSize) {
+ memorySizeForWrite -= estimatedTemporaryMemSize;
+ FLUSH_THERSHOLD = memorySizeForWrite * config.getFlushProportion();
+ REJECT_THERSHOLD = memorySizeForWrite * config.getRejectProportion();
+ }
+
+ public synchronized void releaseTemporaryMemoryForFlushing(long
estimatedTemporaryMemSize) {
+ memorySizeForWrite += estimatedTemporaryMemSize;
+ FLUSH_THERSHOLD = memorySizeForWrite * config.getFlushProportion();
+ REJECT_THERSHOLD = memorySizeForWrite * config.getRejectProportion();
+ }
}
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 5977155..98f9403 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
@@ -339,10 +339,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();