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;

Reply via email to