This is an automated email from the ASF dual-hosted git repository. xiangweiwei pushed a commit to branch maxCachedBuffer in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit f39fd50ece913f3af38939ec54f3ae5feee34a03 Author: Alima777 <[email protected]> AuthorDate: Sat Nov 27 11:38:25 2021 +0800 add max cached buffer size --- .../resources/conf/iotdb-engine.properties | 3 +++ server/src/assembly/resources/conf/iotdb-env.bat | 3 +++ server/src/assembly/resources/conf/iotdb-env.sh | 5 ++++- .../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 11 +++++----- .../org/apache/iotdb/db/conf/IoTDBDescriptor.java | 13 +++++++---- .../iotdb/db/query/pool/QueryTaskPoolManager.java | 4 ++-- .../iotdb/tsfile/read/TsFileSequenceReader.java | 25 ++++++++++++++++++++-- 7 files changed, 50 insertions(+), 14 deletions(-) diff --git a/server/src/assembly/resources/conf/iotdb-engine.properties b/server/src/assembly/resources/conf/iotdb-engine.properties index 3e2da93..c8cc653 100644 --- a/server/src/assembly/resources/conf/iotdb-engine.properties +++ b/server/src/assembly/resources/conf/iotdb-engine.properties @@ -530,6 +530,9 @@ timestamp_precision=ms # Datatype: int # query_timeout_threshold=60000 +# The max number of sub query thread +# max_concurrent_sub_query_thread=8 + #################### ### Metadata Cache Configuration #################### diff --git a/server/src/assembly/resources/conf/iotdb-env.bat b/server/src/assembly/resources/conf/iotdb-env.bat index 6a165ba..3653b0e 100644 --- a/server/src/assembly/resources/conf/iotdb-env.bat +++ b/server/src/assembly/resources/conf/iotdb-env.bat @@ -103,9 +103,12 @@ for /f "tokens=1-3" %%j in ('java -version 2^>^&1') do ( @REM maximum direct memory size set MAX_DIRECT_MEMORY_SIZE=%MAX_HEAP_SIZE% +@REM Max cached buffer size, Note: unit can only be B +set MAX_DIRECT_MEMORY_SIZE=%MAX_DIRECT_MEMORY_SIZE% / 16 / 1024 set IOTDB_HEAP_OPTS=-Xmx%MAX_HEAP_SIZE% -Xms%HEAP_NEWSIZE% -Xlog:gc:"..\gc.log" set IOTDB_HEAP_OPTS=%IOTDB_HEAP_OPTS% -XX:MaxDirectMemorySize=%MAX_DIRECT_MEMORY_SIZE% +set IOTDB_HEAP_OPTS=%IOTDB_HEAP_OPTS% -Djdk.nio.maxCachedBufferSize=%MAX_CACHED_BUFFER_SIZE% @REM You can put your env variable here @REM set JAVA_HOME=%JAVA_HOME% diff --git a/server/src/assembly/resources/conf/iotdb-env.sh b/server/src/assembly/resources/conf/iotdb-env.sh index a62f99e..5caf88e 100755 --- a/server/src/assembly/resources/conf/iotdb-env.sh +++ b/server/src/assembly/resources/conf/iotdb-env.sh @@ -205,8 +205,10 @@ calculate_heap_sizes #MAX_HEAP_SIZE="2G" # Minimum heap size #HEAP_NEWSIZE="2G" -# maximum direct memory size +# Maximum direct memory size MAX_DIRECT_MEMORY_SIZE=${MAX_HEAP_SIZE} +# Max cached buffer size, Note: unit can only be B +MAX_CACHED_BUFFER_SIZE=`expr $MAX_DIRECT_MEMORY_SIZE / 16 / 1024` #true or false #DO NOT FORGET TO MODIFY THE PASSWORD FOR SECURITY (${IOTDB_CONF}/jmx.password and ${IOTDB_CONF}/jmx.access) @@ -241,6 +243,7 @@ fi IOTDB_JMX_OPTS="$IOTDB_JMX_OPTS -Xms${HEAP_NEWSIZE}" IOTDB_JMX_OPTS="$IOTDB_JMX_OPTS -Xmx${MAX_HEAP_SIZE}" IOTDB_JMX_OPTS="$IOTDB_JMX_OPTS -XX:MaxDirectMemorySize=${MAX_DIRECT_MEMORY_SIZE}" +IOTDB_JMX_OPTS="$IOTDB_JMX_OPTS -Djdk.nio.maxCachedBufferSize=${MAX_CACHED_BUFFER_SIZE}" echo "Maximum memory allocation pool = ${MAX_HEAP_SIZE}B, initial memory allocation pool = ${HEAP_NEWSIZE}B" echo "If you want to change this configuration, please check conf/iotdb-env.sh(Unix or OS X, if you use Windows, check conf/iotdb-env.bat)." 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 71f9da7..7f4a01b 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 @@ -253,7 +253,7 @@ public class IoTDBConfig { private int concurrentFlushThread = Runtime.getRuntime().availableProcessors(); /** How many threads can concurrently query. When <= 0, use CPU core number. */ - private int concurrentQueryThread = Runtime.getRuntime().availableProcessors(); + private int maxConcurrentSubQueryThread = Runtime.getRuntime().availableProcessors(); /** How many threads can concurrently evaluate windows. When <= 0, use CPU core number. */ private int concurrentWindowEvaluationThread = Runtime.getRuntime().availableProcessors(); @@ -1159,12 +1159,13 @@ public class IoTDBConfig { this.concurrentFlushThread = concurrentFlushThread; } - public int getConcurrentQueryThread() { - return concurrentQueryThread; + public int getMaxConcurrentSubQueryThread() { + return maxConcurrentSubQueryThread; } - void setConcurrentQueryThread(int concurrentQueryThread) { - this.concurrentQueryThread = concurrentQueryThread; + void setMaxConcurrentSubQueryThread(int maxConcurrentSubQueryThread) { + this.maxConcurrentSubQueryThread = + Math.min(Runtime.getRuntime().availableProcessors(), maxConcurrentSubQueryThread); } public int getConcurrentWindowEvaluationThread() { 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 17deb09..5201399 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 @@ -385,6 +385,11 @@ public class IoTDBDescriptor { properties.getProperty( "query_timeout_threshold", Integer.toString(conf.getQueryTimeoutThreshold())))); + conf.setMaxConcurrentSubQueryThread( + Integer.parseInt( + properties.getProperty( + "max_concurrent_sub_query", Integer.toString(conf.getMaxConcurrentSubQueryThread())))); + conf.setSessionTimeoutThreshold( Integer.parseInt( properties.getProperty( @@ -440,13 +445,13 @@ public class IoTDBDescriptor { "index_buffer_size", Long.toString(conf.getIndexBufferSize())))); // end: index parameter setting - conf.setConcurrentQueryThread( + conf.setMaxConcurrentSubQueryThread( Integer.parseInt( properties.getProperty( - "concurrent_query_thread", Integer.toString(conf.getConcurrentQueryThread())))); + "concurrent_query_thread", Integer.toString(conf.getMaxConcurrentSubQueryThread())))); - if (conf.getConcurrentQueryThread() <= 0) { - conf.setConcurrentQueryThread(Runtime.getRuntime().availableProcessors()); + if (conf.getMaxConcurrentSubQueryThread() <= 0) { + conf.setMaxConcurrentSubQueryThread(Runtime.getRuntime().availableProcessors()); } conf.setmManagerCacheSize( diff --git a/server/src/main/java/org/apache/iotdb/db/query/pool/QueryTaskPoolManager.java b/server/src/main/java/org/apache/iotdb/db/query/pool/QueryTaskPoolManager.java index fa284b7..93233dc 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/pool/QueryTaskPoolManager.java +++ b/server/src/main/java/org/apache/iotdb/db/query/pool/QueryTaskPoolManager.java @@ -32,7 +32,7 @@ public class QueryTaskPoolManager extends AbstractPoolManager { private static final Logger LOGGER = LoggerFactory.getLogger(QueryTaskPoolManager.class); private QueryTaskPoolManager() { - int threadCnt = IoTDBDescriptor.getInstance().getConfig().getConcurrentQueryThread(); + int threadCnt = IoTDBDescriptor.getInstance().getConfig().getMaxConcurrentSubQueryThread(); pool = IoTDBThreadPoolFactory.newFixedThreadPool(threadCnt, ThreadName.QUERY_SERVICE.getName()); } @@ -53,7 +53,7 @@ public class QueryTaskPoolManager extends AbstractPoolManager { @Override public void start() { if (pool == null) { - int threadCnt = IoTDBDescriptor.getInstance().getConfig().getConcurrentQueryThread(); + int threadCnt = IoTDBDescriptor.getInstance().getConfig().getMaxConcurrentSubQueryThread(); pool = IoTDBThreadPoolFactory.newFixedThreadPool(threadCnt, ThreadName.QUERY_SERVICE.getName()); } diff --git a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java index 6f5e088..1e7752a 100644 --- a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java +++ b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/TsFileSequenceReader.java @@ -27,7 +27,15 @@ import org.apache.iotdb.tsfile.file.MetaMarker; import org.apache.iotdb.tsfile.file.header.ChunkGroupHeader; import org.apache.iotdb.tsfile.file.header.ChunkHeader; import org.apache.iotdb.tsfile.file.header.PageHeader; -import org.apache.iotdb.tsfile.file.metadata.*; +import org.apache.iotdb.tsfile.file.metadata.AlignedTimeSeriesMetadata; +import org.apache.iotdb.tsfile.file.metadata.ChunkGroupMetadata; +import org.apache.iotdb.tsfile.file.metadata.ChunkMetadata; +import org.apache.iotdb.tsfile.file.metadata.IChunkMetadata; +import org.apache.iotdb.tsfile.file.metadata.ITimeSeriesMetadata; +import org.apache.iotdb.tsfile.file.metadata.MetadataIndexEntry; +import org.apache.iotdb.tsfile.file.metadata.MetadataIndexNode; +import org.apache.iotdb.tsfile.file.metadata.TimeseriesMetadata; +import org.apache.iotdb.tsfile.file.metadata.TsFileMetadata; import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType; import org.apache.iotdb.tsfile.file.metadata.enums.MetadataIndexNodeType; import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType; @@ -54,7 +62,20 @@ import java.io.IOException; import java.io.Serializable; import java.nio.BufferOverflowException; import java.nio.ByteBuffer; -import java.util.*; +import java.util.ArrayList; +import java.util.Collections; +import java.util.Comparator; +import java.util.HashMap; +import java.util.HashSet; +import java.util.Iterator; +import java.util.LinkedHashMap; +import java.util.LinkedList; +import java.util.List; +import java.util.Map; +import java.util.NoSuchElementException; +import java.util.Queue; +import java.util.Set; +import java.util.TreeMap; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantReadWriteLock;
