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

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


The following commit(s) were added to refs/heads/rel/0.12 by this push:
     new 18f0144  [To rel/0.12] Cherry pick from autoai (#3637)
18f0144 is described below

commit 18f0144d099712b4323f9c836e878a87820a80e7
Author: Jialin Qiao <[email protected]>
AuthorDate: Tue Jul 27 06:54:44 2021 -0500

    [To rel/0.12] Cherry pick from autoai (#3637)
---
 .../iotdb/cluster/client/DataClientProvider.java   |   9 +-
 .../iotdb/cluster/client/sync/SyncClientPool.java  |   2 +-
 .../iotdb/cluster/coordinator/Coordinator.java     |  17 +-
 .../apache/iotdb/cluster/metadata/CMManager.java   |  71 ++++++--
 .../apache/iotdb/cluster/metadata/MetaPuller.java  |  18 +-
 .../iotdb/cluster/query/ClusterPlanExecutor.java   |  39 +++-
 .../cluster/query/aggregate/ClusterAggregator.java |   9 +-
 .../cluster/query/fill/ClusterPreviousFill.java    |  25 ++-
 .../query/groupby/RemoteGroupByExecutor.java       |  21 ++-
 .../query/last/ClusterLastQueryExecutor.java       |  26 ++-
 .../cluster/query/reader/ClusterReaderFactory.java |   8 +-
 .../iotdb/cluster/query/reader/DataSourceInfo.java |  36 ++--
 .../reader/RemoteSeriesReaderByTimestamp.java      |   2 +
 .../query/reader/RemoteSimpleSeriesReader.java     |   2 +
 .../query/reader/mult/MultDataSourceInfo.java      |  15 +-
 .../query/reader/mult/RemoteMultSeriesReader.java  |  17 +-
 .../apache/iotdb/cluster/server/ClientServer.java  |   8 +-
 .../cluster/server/heartbeat/HeartbeatThread.java  |   7 +
 .../cluster/server/member/DataGroupMember.java     |   8 +-
 .../cluster/server/member/MetaGroupMember.java     |  19 +-
 .../iotdb/cluster/server/member/RaftMember.java    |   9 +-
 .../cluster/client/DataClientProviderTest.java     |   5 +-
 .../org/apache/iotdb/SessionConcurrentExample.java | 198 +++++++++++++++++++++
 .../db/engine/storagegroup/TsFileProcessor.java    |  10 ++
 .../org/apache/iotdb/db/metadata/MManager.java     |  10 +-
 .../iotdb/db/rescon/PrimitiveArrayManager.java     |  10 +-
 .../org/apache/iotdb/db/service/TSServiceImpl.java |   8 +-
 .../db/writelog/node/ExclusiveWriteLogNode.java    |  81 +++------
 .../writelog/recover/TsFileRecoverPerformer.java   |  25 ++-
 .../iotdb/db/integration/IoTDBRestartIT.java       |  48 +++++
 .../read/expression/util/ExpressionOptimizer.java  |  34 ++--
 31 files changed, 601 insertions(+), 196 deletions(-)

diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/client/DataClientProvider.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/client/DataClientProvider.java
index 106705f..0950958 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/client/DataClientProvider.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/client/DataClientProvider.java
@@ -28,7 +28,6 @@ import org.apache.iotdb.cluster.config.ClusterDescriptor;
 import org.apache.iotdb.cluster.rpc.thrift.Node;
 import org.apache.iotdb.cluster.rpc.thrift.RaftService.Client;
 
-import org.apache.thrift.TException;
 import org.apache.thrift.protocol.TProtocolFactory;
 
 import java.io.IOException;
@@ -102,10 +101,10 @@ public class DataClientProvider {
    * @param node the node to be connected
    * @param timeout timeout threshold of connection
    */
-  public SyncDataClient getSyncDataClient(Node node, int timeout) throws 
TException {
+  public SyncDataClient getSyncDataClient(Node node, int timeout) throws 
IOException {
     SyncDataClient client = (SyncDataClient) 
getDataSyncClientPool().getClient(node);
     if (client == null) {
-      throw new TException(GET_CLIENT_FAILED_MSG + node);
+      throw new IOException(GET_CLIENT_FAILED_MSG + node);
     }
     client.setTimeout(timeout);
     return client;
@@ -121,10 +120,10 @@ public class DataClientProvider {
    * @param node the node to be connected
    * @param timeout timeout threshold of connection
    */
-  public SyncDataClient getSyncDataClientForRefresh(Node node, int timeout) 
throws TException {
+  public SyncDataClient getSyncDataClientForRefresh(Node node, int timeout) 
throws IOException {
     SyncDataClient client = (SyncDataClient) 
getDataSyncClientPool().getClientForRefresh(node);
     if (client == null) {
-      throw new TException(GET_CLIENT_FAILED_MSG + node);
+      throw new IOException(GET_CLIENT_FAILED_MSG + node);
     }
     client.setTimeout(timeout);
     return client;
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/client/sync/SyncClientPool.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/client/sync/SyncClientPool.java
index 67dbae5..f607fa3 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/client/sync/SyncClientPool.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/client/sync/SyncClientPool.java
@@ -67,7 +67,7 @@ public class SyncClientPool {
       if (clientStack.isEmpty()) {
         return null;
       } else {
-        return clientStack.pollLast();
+        return clientStack.poll();
       }
     }
   }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/coordinator/Coordinator.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/coordinator/Coordinator.java
index 11f99e8..db0fce3 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/coordinator/Coordinator.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/coordinator/Coordinator.java
@@ -58,7 +58,6 @@ import org.apache.iotdb.rpc.TSStatusCode;
 import org.apache.iotdb.service.rpc.thrift.EndPoint;
 import org.apache.iotdb.service.rpc.thrift.TSStatus;
 
-import org.apache.thrift.TException;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -703,15 +702,11 @@ public class Coordinator {
 
   private TSStatus forwardDataPlanSync(PhysicalPlan plan, Node receiver, Node 
header)
       throws IOException {
-    RaftService.Client client = null;
-    try {
-      client =
-          metaGroupMember
-              .getClientProvider()
-              .getSyncDataClient(receiver, 
RaftServer.getWriteOperationTimeoutMS());
-    } catch (TException e) {
-      throw new IOException(e);
-    }
+    RaftService.Client client =
+        metaGroupMember
+            .getClientProvider()
+            .getSyncDataClient(receiver, 
RaftServer.getWriteOperationTimeoutMS());
+
     return this.metaGroupMember.forwardPlanSync(plan, receiver, header, 
client);
   }
 
@@ -735,7 +730,7 @@ public class Coordinator {
    * @param node the node to be connected
    * @param timeout timeout threshold of connection
    */
-  public SyncDataClient getSyncDataClient(Node node, int timeout) throws 
TException {
+  public SyncDataClient getSyncDataClient(Node node, int timeout) throws 
IOException {
     return metaGroupMember.getClientProvider().getSyncDataClient(node, 
timeout);
   }
 }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java
index 319159a..05d61a4 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java
@@ -794,8 +794,14 @@ public class CMManager extends MManager {
               metaGroupMember
                   .getClientProvider()
                   .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
-            result =
-                
syncDataClient.getUnregisteredTimeseries(partitionGroup.getHeader(), 
seriesList);
+            try {
+              result =
+                  
syncDataClient.getUnregisteredTimeseries(partitionGroup.getHeader(), 
seriesList);
+            } catch (TException e) {
+              // the connection may be broken, close it to avoid it being 
reused
+              syncDataClient.getInputProtocol().getTransport().close();
+              throw e;
+            }
           }
         }
         if (result != null) {
@@ -978,13 +984,18 @@ public class CMManager extends MManager {
           metaGroupMember
               .getClientProvider()
               .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
-
-        PullSchemaResp pullSchemaResp = 
syncDataClient.pullTimeSeriesSchema(request);
-        ByteBuffer buffer = pullSchemaResp.schemaBytes;
-        int size = buffer.getInt();
-        schemas = new ArrayList<>(size);
-        for (int i = 0; i < size; i++) {
-          schemas.add(TimeseriesSchema.deserializeFrom(buffer));
+        try {
+          PullSchemaResp pullSchemaResp = 
syncDataClient.pullTimeSeriesSchema(request);
+          ByteBuffer buffer = pullSchemaResp.schemaBytes;
+          int size = buffer.getInt();
+          schemas = new ArrayList<>(size);
+          for (int i = 0; i < size; i++) {
+            schemas.add(TimeseriesSchema.deserializeFrom(buffer));
+          }
+        } catch (TException e) {
+          // the connection may be broken, close it to avoid it being reused
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
         }
       }
     }
@@ -1212,8 +1223,13 @@ public class CMManager extends MManager {
           metaGroupMember
               .getClientProvider()
               .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
-
-        result = syncDataClient.getAllPaths(header, pathsToQuery, withAlias);
+        try {
+          result = syncDataClient.getAllPaths(header, pathsToQuery, withAlias);
+        } catch (TException e) {
+          // the connection may be broken, close it to avoid it being reused
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
+        }
       }
     }
 
@@ -1338,8 +1354,13 @@ public class CMManager extends MManager {
           metaGroupMember
               .getClientProvider()
               .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
-
-        paths = syncDataClient.getAllDevices(header, pathsToQuery);
+        try {
+          paths = syncDataClient.getAllDevices(header, pathsToQuery);
+        } catch (TException e) {
+          // the connection may be broken, close it to avoid it being reused
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
+        }
       }
     }
     return paths;
@@ -1792,9 +1813,15 @@ public class CMManager extends MManager {
                   .getClientProvider()
                   .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
         plan.serialize(dataOutputStream);
-        resultBinary =
-            syncDataClient.getAllMeasurementSchema(
-                group.getHeader(), 
ByteBuffer.wrap(byteArrayOutputStream.toByteArray()));
+        try {
+          resultBinary =
+              syncDataClient.getAllMeasurementSchema(
+                  group.getHeader(), 
ByteBuffer.wrap(byteArrayOutputStream.toByteArray()));
+        } catch (TException e) {
+          // the connection may be broken, close it to avoid it being reused
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
+        }
       }
     }
     return resultBinary;
@@ -1818,9 +1845,15 @@ public class CMManager extends MManager {
                   .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
 
         plan.serialize(dataOutputStream);
-        resultBinary =
-            syncDataClient.getDevices(
-                group.getHeader(), 
ByteBuffer.wrap(byteArrayOutputStream.toByteArray()));
+        try {
+          resultBinary =
+              syncDataClient.getDevices(
+                  group.getHeader(), 
ByteBuffer.wrap(byteArrayOutputStream.toByteArray()));
+        } catch (TException e) {
+          // the connection may be broken, close it to avoid it being reused
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
+        }
       }
     }
     return resultBinary;
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/metadata/MetaPuller.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/metadata/MetaPuller.java
index e524772..9991c5a 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/metadata/MetaPuller.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/metadata/MetaPuller.java
@@ -231,12 +231,18 @@ public class MetaPuller {
               .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
 
         // only need measurement name
-        PullSchemaResp pullSchemaResp = 
syncDataClient.pullMeasurementSchema(request);
-        ByteBuffer buffer = pullSchemaResp.schemaBytes;
-        int size = buffer.getInt();
-        schemas = new ArrayList<>(size);
-        for (int i = 0; i < size; i++) {
-          schemas.add(MeasurementSchema.deserializeFrom(buffer));
+        try {
+          PullSchemaResp pullSchemaResp = 
syncDataClient.pullMeasurementSchema(request);
+          ByteBuffer buffer = pullSchemaResp.schemaBytes;
+          int size = buffer.getInt();
+          schemas = new ArrayList<>(size);
+          for (int i = 0; i < size; i++) {
+            schemas.add(MeasurementSchema.deserializeFrom(buffer));
+          }
+        } catch (TException e) {
+          // the connection may be broken, close it to avoid it being reused
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
         }
       }
     }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterPlanExecutor.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterPlanExecutor.java
index ca16cfb..34e9904 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterPlanExecutor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/ClusterPlanExecutor.java
@@ -259,8 +259,13 @@ public class ClusterPlanExecutor extends PlanExecutor {
               metaGroupMember
                   .getClientProvider()
                   .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
-            syncDataClient.setTimeout(RaftServer.getReadOperationTimeoutMS());
-            count = syncDataClient.getPathCount(partitionGroup.getHeader(), 
pathsToQuery, level);
+            try {
+              count = syncDataClient.getPathCount(partitionGroup.getHeader(), 
pathsToQuery, level);
+            } catch (TException e) {
+              // the connection may be broken, close it to avoid it being 
reused
+              syncDataClient.getInputProtocol().getTransport().close();
+              throw e;
+            }
           }
         }
         logger.debug(
@@ -363,8 +368,14 @@ public class ClusterPlanExecutor extends PlanExecutor {
               metaGroupMember
                   .getClientProvider()
                   .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
-            paths =
-                syncDataClient.getNodeList(group.getHeader(), 
schemaPattern.getFullPath(), level);
+            try {
+              paths =
+                  syncDataClient.getNodeList(group.getHeader(), 
schemaPattern.getFullPath(), level);
+            } catch (TException e) {
+              // the connection may be broken, close it to avoid it being 
reused
+              syncDataClient.getInputProtocol().getTransport().close();
+              throw e;
+            }
           }
         }
         if (paths != null) {
@@ -449,8 +460,14 @@ public class ClusterPlanExecutor extends PlanExecutor {
               metaGroupMember
                   .getClientProvider()
                   .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
-            nextChildrenNodes =
-                syncDataClient.getChildNodeInNextLevel(group.getHeader(), 
path.getFullPath());
+            try {
+              nextChildrenNodes =
+                  syncDataClient.getChildNodeInNextLevel(group.getHeader(), 
path.getFullPath());
+            } catch (TException e) {
+              // the connection may be broken, close it to avoid it being 
reused
+              syncDataClient.getInputProtocol().getTransport().close();
+              throw e;
+            }
           }
         }
         if (nextChildrenNodes != null) {
@@ -558,8 +575,14 @@ public class ClusterPlanExecutor extends PlanExecutor {
               metaGroupMember
                   .getClientProvider()
                   .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
-            nextChildren =
-                syncDataClient.getChildNodePathInNextLevel(group.getHeader(), 
path.getFullPath());
+            try {
+              nextChildren =
+                  
syncDataClient.getChildNodePathInNextLevel(group.getHeader(), 
path.getFullPath());
+            } catch (TException e) {
+              // the connection may be broken, close it to avoid it being 
reused
+              syncDataClient.getInputProtocol().getTransport().close();
+              throw e;
+            }
           }
         }
         if (nextChildren != null) {
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java
index 6c33f50..dc52923 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java
@@ -274,8 +274,13 @@ public class ClusterAggregator {
           metaGroupMember
               .getClientProvider()
               .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
-
-        resultBuffers = syncDataClient.getAggrResult(request);
+        try {
+          resultBuffers = syncDataClient.getAggrResult(request);
+        } catch (TException e) {
+          // the connection may be broken, close it to avoid it being reused
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
+        }
       }
     }
     return resultBuffers;
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterPreviousFill.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterPreviousFill.java
index 33274e3..9af7082 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterPreviousFill.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/fill/ClusterPreviousFill.java
@@ -32,6 +32,7 @@ import org.apache.iotdb.cluster.server.RaftServer;
 import org.apache.iotdb.cluster.server.handlers.caller.PreviousFillHandler;
 import org.apache.iotdb.cluster.server.member.DataGroupMember;
 import org.apache.iotdb.cluster.server.member.MetaGroupMember;
+import org.apache.iotdb.cluster.utils.ClientUtils;
 import org.apache.iotdb.cluster.utils.PartitionUtils.Intervals;
 import org.apache.iotdb.db.exception.StorageEngineException;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
@@ -41,6 +42,7 @@ import org.apache.iotdb.db.query.executor.fill.PreviousFill;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
 import org.apache.iotdb.tsfile.read.TimeValuePair;
 
+import org.apache.thrift.TException;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -240,19 +242,28 @@ public class ClusterPreviousFill extends PreviousFill {
   private ByteBuffer remoteSyncPreviousFill(
       Node node, PreviousFillRequest request, PreviousFillArguments arguments) 
{
     ByteBuffer byteBuffer = null;
-    try (SyncDataClient syncDataClient =
-        metaGroupMember
-            .getClientProvider()
-            .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) {
-
-      byteBuffer = syncDataClient.previousFill(request);
-    } catch (Exception e) {
+    SyncDataClient client = null;
+    try {
+      client =
+          metaGroupMember
+              .getClientProvider()
+              .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS());
+      byteBuffer = client.previousFill(request);
+    } catch (IOException e) {
+      logger.warn("{}: Cannot connect to {} during previous fill", 
metaGroupMember, node);
+    } catch (TException e) {
       logger.error(
           "{}: Cannot perform previous fill of {} to {}",
           metaGroupMember.getName(),
           arguments.getPath(),
           node,
           e);
+      // the connection may be broken, close it to avoid it being reused
+      client.getInputProtocol().getTransport().close();
+    } finally {
+      if (client != null) {
+        ClientUtils.putBackSyncClient(client);
+      }
     }
     return byteBuffer;
   }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/groupby/RemoteGroupByExecutor.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/groupby/RemoteGroupByExecutor.java
index 02df747..468ec57 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/groupby/RemoteGroupByExecutor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/groupby/RemoteGroupByExecutor.java
@@ -89,8 +89,14 @@ public class RemoteGroupByExecutor implements 
GroupByExecutor {
                 .getClientProvider()
                 .getSyncDataClient(source, 
RaftServer.getReadOperationTimeoutMS())) {
 
-          aggrBuffers =
-              syncDataClient.getGroupByResult(header, executorId, 
curStartTime, curEndTime);
+          try {
+            aggrBuffers =
+                syncDataClient.getGroupByResult(header, executorId, 
curStartTime, curEndTime);
+          } catch (TException e) {
+            // the connection may be broken, close it to avoid it being reused
+            syncDataClient.getInputProtocol().getTransport().close();
+            throw e;
+          }
         }
       }
     } catch (TException e) {
@@ -133,9 +139,14 @@ public class RemoteGroupByExecutor implements 
GroupByExecutor {
             metaGroupMember
                 .getClientProvider()
                 .getSyncDataClient(source, 
RaftServer.getReadOperationTimeoutMS())) {
-
-          aggrBuffer =
-              syncDataClient.peekNextNotNullValue(header, executorId, 
nextStartTime, nextEndTime);
+          try {
+            aggrBuffer =
+                syncDataClient.peekNextNotNullValue(header, executorId, 
nextStartTime, nextEndTime);
+          } catch (TException e) {
+            // the connection may be broken, close it to avoid it being reused
+            syncDataClient.getInputProtocol().getTransport().close();
+            throw e;
+          }
         }
       }
     } catch (TException e) {
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java
index d5ec324..a30db6c 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/last/ClusterLastQueryExecutor.java
@@ -30,6 +30,7 @@ import org.apache.iotdb.cluster.rpc.thrift.Node;
 import org.apache.iotdb.cluster.server.RaftServer;
 import org.apache.iotdb.cluster.server.member.DataGroupMember;
 import org.apache.iotdb.cluster.server.member.MetaGroupMember;
+import org.apache.iotdb.cluster.utils.ClientUtils;
 import org.apache.iotdb.cluster.utils.ClusterQueryUtils;
 import org.apache.iotdb.db.exception.StorageEngineException;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
@@ -258,19 +259,30 @@ public class ClusterLastQueryExecutor extends 
LastQueryExecutor {
     }
 
     private ByteBuffer lastSync(Node node, QueryContext context) throws 
TException {
-      try (SyncDataClient syncDataClient =
-          metaGroupMember
-              .getClientProvider()
-              .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
-
-        return syncDataClient.last(
+      SyncDataClient client = null;
+      try {
+        client =
+            metaGroupMember
+                .getClientProvider()
+                .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS());
+        return client.last(
             new LastQueryRequest(
                 PartialPath.toStringList(seriesPaths),
                 dataTypeOrdinals,
                 context.getQueryId(),
                 queryPlan.getDeviceToMeasurements(),
                 group.getHeader(),
-                syncDataClient.getNode()));
+                client.getNode()));
+      } catch (IOException e) {
+        return null;
+      } catch (TException e) {
+        // the connection may be broken, close it to avoid it being reused
+        client.getInputProtocol().getTransport().close();
+        throw e;
+      } finally {
+        if (client != null) {
+          ClientUtils.putBackSyncClient(client);
+        }
       }
     }
   }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java
index 64a2e3b..e3e0d53 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java
@@ -896,7 +896,13 @@ public class ClusterReaderFactory {
               .getClientProvider()
               .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
 
-        executorId = syncDataClient.getGroupByExecutor(request);
+        try {
+          executorId = syncDataClient.getGroupByExecutor(request);
+        } catch (TException e) {
+          // the connection may be broken, close it to avoid it being reused
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
+        }
       }
     }
     return executorId;
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/DataSourceInfo.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/DataSourceInfo.java
index 8889535..4ba11e4 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/DataSourceInfo.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/DataSourceInfo.java
@@ -153,27 +153,31 @@ public class DataSourceInfo {
   }
 
   private Long applyForReaderIdSync(Node node, boolean byTimestamp, long 
timestamp)
-      throws TException {
-
-    Long newReaderId;
+      throws TException, IOException {
+    long newReaderId;
     try (SyncDataClient client =
         this.metaGroupMember
             .getClientProvider()
             .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) {
-
-      if (byTimestamp) {
-        newReaderId = client.querySingleSeriesByTimestamp(request);
-      } else {
-        Filter newFilter;
-        // add timestamp to as a timeFilter to skip the data which has been 
read
-        if (request.isSetTimeFilterBytes()) {
-          Filter timeFilter = 
FilterFactory.deserialize(request.timeFilterBytes);
-          newFilter = new AndFilter(timeFilter, TimeFilter.gt(timestamp));
+      try {
+        if (byTimestamp) {
+          newReaderId = client.querySingleSeriesByTimestamp(request);
         } else {
-          newFilter = TimeFilter.gt(timestamp);
+          Filter newFilter;
+          // add timestamp to as a timeFilter to skip the data which has been 
read
+          if (request.isSetTimeFilterBytes()) {
+            Filter timeFilter = 
FilterFactory.deserialize(request.timeFilterBytes);
+            newFilter = new AndFilter(timeFilter, TimeFilter.gt(timestamp));
+          } else {
+            newFilter = TimeFilter.gt(timestamp);
+          }
+          
request.setTimeFilterBytes(SerializeUtils.serializeFilter(newFilter));
+          newReaderId = client.querySingleSeries(request);
         }
-        request.setTimeFilterBytes(SerializeUtils.serializeFilter(newFilter));
-        newReaderId = client.querySingleSeries(request);
+      } catch (TException e) {
+        // the connection may be broken, close it to avoid it being reused
+        client.getInputProtocol().getTransport().close();
+        throw e;
       }
       return newReaderId;
     }
@@ -201,7 +205,7 @@ public class DataSourceInfo {
         : 
metaGroupMember.getClientProvider().getAsyncDataClient(this.curSource, timeout);
   }
 
-  SyncDataClient getCurSyncClient(int timeout) throws TException {
+  SyncDataClient getCurSyncClient(int timeout) throws IOException {
     return isNoClient
         ? null
         : 
metaGroupMember.getClientProvider().getSyncDataClient(this.curSource, timeout);
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSeriesReaderByTimestamp.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSeriesReaderByTimestamp.java
index d077f02..e266af8 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSeriesReaderByTimestamp.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSeriesReaderByTimestamp.java
@@ -108,6 +108,8 @@ public class RemoteSeriesReaderByTimestamp implements 
IReaderByTimestamp {
       return curSyncClient.fetchSingleSeriesByTimestamps(
           sourceInfo.getHeader(), sourceInfo.getReaderId(), timestampList);
     } catch (TException e) {
+      // the connection may be broken, close it to avoid it being reused
+      curSyncClient.getInputProtocol().getTransport().close();
       // try other node
       if (!sourceInfo.switchNode(true, timestamps[0])) {
         return null;
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSimpleSeriesReader.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSimpleSeriesReader.java
index 2dcc1b7..f53f2bc 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSimpleSeriesReader.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/RemoteSimpleSeriesReader.java
@@ -149,6 +149,8 @@ public class RemoteSimpleSeriesReader implements 
IPointReader {
       curSyncClient = 
sourceInfo.getCurSyncClient(RaftServer.getReadOperationTimeoutMS());
       return curSyncClient.fetchSingleSeries(sourceInfo.getHeader(), 
sourceInfo.getReaderId());
     } catch (TException e) {
+      // the connection may be broken, close it to avoid it being reused
+      curSyncClient.getInputProtocol().getTransport().close();
       // try other node
       if (!sourceInfo.switchNode(false, lastTimestamp)) {
         return null;
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/MultDataSourceInfo.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/MultDataSourceInfo.java
index 27c9f59..a4488aa 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/MultDataSourceInfo.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/MultDataSourceInfo.java
@@ -172,9 +172,8 @@ public class MultDataSourceInfo {
     return result.get();
   }
 
-  private Long applyForReaderIdSync(Node node, long timestamp) throws 
TException {
-
-    Long newReaderId;
+  private Long applyForReaderIdSync(Node node, long timestamp) throws 
TException, IOException {
+    long newReaderId;
     try (SyncDataClient client =
         this.metaGroupMember
             .getClientProvider()
@@ -189,7 +188,13 @@ public class MultDataSourceInfo {
         newFilter = TimeFilter.gt(timestamp);
       }
       request.setTimeFilterBytes(SerializeUtils.serializeFilter(newFilter));
-      newReaderId = client.queryMultSeries(request);
+      try {
+        newReaderId = client.queryMultSeries(request);
+      } catch (TException e) {
+        // the connection may be broken, close it to avoid it being reused
+        client.getInputProtocol().getTransport().close();
+        throw e;
+      }
       return newReaderId;
     }
   }
@@ -212,7 +217,7 @@ public class MultDataSourceInfo {
         : 
metaGroupMember.getClientProvider().getAsyncDataClient(this.curSource, timeout);
   }
 
-  SyncDataClient getCurSyncClient(int timeout) throws TException {
+  SyncDataClient getCurSyncClient(int timeout) throws IOException {
     return isNoClient
         ? null
         : 
metaGroupMember.getClientProvider().getSyncDataClient(this.curSource, timeout);
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/RemoteMultSeriesReader.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/RemoteMultSeriesReader.java
index bf20b35..36513d4 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/RemoteMultSeriesReader.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/mult/RemoteMultSeriesReader.java
@@ -181,15 +181,18 @@ public class RemoteMultSeriesReader extends 
AbstractMultPointReader {
   }
 
   private Map<String, ByteBuffer> fetchResultSync(List<String> paths) throws 
IOException {
-
     try (SyncDataClient curSyncClient =
-        sourceInfo.getCurSyncClient(RaftServer.getReadOperationTimeoutMS()); ) 
{
-
-      return curSyncClient.fetchMultSeries(sourceInfo.getHeader(), 
sourceInfo.getReaderId(), paths);
-    } catch (TException e) {
-      logger.error("Failed to fetch result sync, connect to {}", sourceInfo, 
e);
-      return null;
+        sourceInfo.getCurSyncClient(RaftServer.getReadOperationTimeoutMS())) {
+      try {
+        return curSyncClient.fetchMultSeries(
+            sourceInfo.getHeader(), sourceInfo.getReaderId(), paths);
+      } catch (TException e) {
+        // the connection may be broken, close it to avoid it being reused
+        curSyncClient.getInputProtocol().getTransport().close();
+        logger.error("Failed to fetch result sync, connect to {}", sourceInfo, 
e);
+      }
     }
+    return null;
   }
 
   /** select path, which could batch-fetch result */
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java
index c627373..722aeaf 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java
@@ -317,7 +317,13 @@ public class ClientServer extends TSServiceImpl {
               try (SyncDataClient syncDataClient =
                   coordinator.getSyncDataClient(
                       queriedNode, RaftServer.getReadOperationTimeoutMS())) {
-                syncDataClient.endQuery(header, coordinator.getThisNode(), 
queryId);
+                try {
+                  syncDataClient.endQuery(header, coordinator.getThisNode(), 
queryId);
+                } catch (TException e) {
+                  // the connection may be broken, close it to avoid it being 
reused
+                  syncDataClient.getInputProtocol().getTransport().close();
+                  throw e;
+                }
               }
             }
           } catch (IOException | TException e) {
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
index 0459fef..67acc5f 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
@@ -34,6 +34,7 @@ import 
org.apache.iotdb.cluster.server.handlers.caller.HeartbeatHandler;
 import org.apache.iotdb.cluster.server.member.RaftMember;
 import org.apache.iotdb.cluster.utils.ClientUtils;
 
+import org.apache.thrift.TException;
 import org.apache.thrift.transport.TTransportException;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -222,6 +223,7 @@ public class HeartbeatThread implements Runnable {
                 } catch (TTransportException e) {
                   logger.warn(
                       "{}: Cannot send heart beat to node {} due to network", 
memberName, node, e);
+                  // the connection may be broken, close it to avoid it being 
reused
                   client.getInputProtocol().getTransport().close();
                 } catch (Exception e) {
                   logger.warn("{}: Cannot send heart beat to node {}", 
memberName, node, e);
@@ -401,6 +403,11 @@ public class HeartbeatThread implements Runnable {
                 try {
                   long result = client.startElection(request);
                   handler.onComplete(result);
+                } catch (TException e) {
+                  // the connection may be broken, close it to avoid it being 
reused
+                  client.getInputProtocol().getTransport().close();
+                  logger.warn("{}: Cannot request a vote from {}", memberName, 
node, e);
+                  handler.onError(e);
                 } catch (Exception e) {
                   handler.onError(e);
                 } finally {
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
index 6f6936a..89d933a 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
@@ -696,6 +696,7 @@ public class DataGroupMember extends RaftMember {
    */
   @Override
   public TSStatus executeNonQueryPlan(PhysicalPlan plan) {
+    long startTime = System.currentTimeMillis();
     if (ClusterDescriptor.getInstance().getConfig().getReplicationNum() == 1) {
       try {
         getLocalExecutor().processNonQuery(plan);
@@ -718,6 +719,11 @@ public class DataGroupMember extends RaftMember {
           }
         }
         return handleLogExecutionException(plan, cause);
+      } finally {
+        long elapsed = System.currentTimeMillis() - startTime;
+        if (elapsed > 5000) {
+          logger.error("PlanExecutor execute slowly : time cost : {}ms", 
elapsed);
+        }
       }
     } else {
       TSStatus status = executeNonQueryPlanWithKnownLeader(plan);
@@ -725,7 +731,7 @@ public class DataGroupMember extends RaftMember {
         return status;
       }
 
-      long startTime = 
Timer.Statistic.DATA_GROUP_MEMBER_WAIT_LEADER.getOperationStartTime();
+      startTime = 
Timer.Statistic.DATA_GROUP_MEMBER_WAIT_LEADER.getOperationStartTime();
       waitLeader();
       
Timer.Statistic.DATA_GROUP_MEMBER_WAIT_LEADER.calOperationCostTimeFromStart(startTime);
 
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java
index cede9a3..1cc4ad9 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java
@@ -169,7 +169,7 @@ public class MetaGroupMember extends RaftMember {
    * every "REFRESH_CLIENT_SEC" seconds, a dataClientRefresher thread will try 
to refresh one thrift
    * connection for each nodes other than itself.
    */
-  private static final int REFRESH_CLIENT_SEC = 5;
+  private static final int REFRESH_CLIENT_SEC = 1;
 
   /** how many times is a data record replicated, also the number of nodes in 
a data group */
   private static final int REPLICATION_NUM =
@@ -498,23 +498,22 @@ public class MetaGroupMember extends RaftMember {
   }
 
   private void refreshClientOnceSync(Node receiver) {
-    RaftService.Client client;
+    RaftService.Client client = null;
     try {
       client =
           getClientProvider()
               .getSyncDataClientForRefresh(receiver, 
RaftServer.getWriteOperationTimeoutMS());
-    } catch (TException e) {
-      return;
-    }
-    try {
       RefreshReuqest req = new RefreshReuqest();
       client.refreshConnection(req);
+    } catch (IOException ignored) {
     } catch (TException e) {
-      logger.warn("encounter refreshing client timeout, throw broken 
connection", e);
+      logger.info("encounter refreshing client timeout, throw broken 
connection", e);
       // the connection may be broken, close it to avoid it being reused
       client.getInputProtocol().getTransport().close();
     } finally {
-      ClientUtils.putBackSyncClient(client);
+      if (client != null) {
+        ClientUtils.putBackSyncClient(client);
+      }
     }
   }
 
@@ -530,13 +529,13 @@ public class MetaGroupMember extends RaftMember {
     try {
       client.refreshConnection(new RefreshReuqest(), new 
GenericHandler<>(receiver, null));
     } catch (TException e) {
-      logger.warn("encounter refreshing client timeout, throw broken 
connection", e);
+      logger.info("encounter refreshing client timeout, throw broken 
connection", e);
     }
   }
 
   private void generateNodeReport() {
     try {
-      if (logger.isInfoEnabled()) {
+      if (logger.isDebugEnabled()) {
         NodeReport report = genNodeReport();
         logger.debug(report.toString());
       }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
index f2187f9..3078013 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
@@ -1317,13 +1317,13 @@ public abstract class RaftMember {
   }
 
   public TSStatus forwardPlanSync(PhysicalPlan plan, Node receiver, Node 
header, Client client) {
+    long startTime = System.currentTimeMillis();
     try {
       ExecutNonQueryReq req = new ExecutNonQueryReq();
       req.setPlanBytes(PlanSerializer.getInstance().serialize(plan));
       if (header != null) {
         req.setHeader(header);
       }
-
       TSStatus tsStatus = client.executeNonQueryPlan(req);
       if (tsStatus == null) {
         tsStatus = StatusUtils.TIME_OUT;
@@ -1337,7 +1337,12 @@ public abstract class RaftMember {
       TSStatus status;
       if (e.getCause() instanceof SocketTimeoutException) {
         status = StatusUtils.TIME_OUT;
-        logger.warn(MSG_FORWARD_TIMEOUT, name, plan, receiver);
+        logger.warn(
+            MSG_FORWARD_TIMEOUT + ": {}ms",
+            name,
+            plan,
+            receiver,
+            System.currentTimeMillis() - startTime);
       } else {
         logger.error(MSG_FORWARD_ERROR, name, plan, receiver, e);
         status = StatusUtils.getStatus(StatusUtils.INTERNAL_ERROR, 
e.getMessage());
diff --git 
a/cluster/src/test/java/org/apache/iotdb/cluster/client/DataClientProviderTest.java
 
b/cluster/src/test/java/org/apache/iotdb/cluster/client/DataClientProviderTest.java
index d2ee0b7..3987450 100644
--- 
a/cluster/src/test/java/org/apache/iotdb/cluster/client/DataClientProviderTest.java
+++ 
b/cluster/src/test/java/org/apache/iotdb/cluster/client/DataClientProviderTest.java
@@ -28,7 +28,6 @@ import org.apache.iotdb.cluster.rpc.thrift.Node;
 import org.apache.iotdb.cluster.utils.ClientUtils;
 import org.apache.iotdb.cluster.utils.ClusterNode;
 
-import org.apache.thrift.TException;
 import org.apache.thrift.protocol.TBinaryProtocol.Factory;
 import org.junit.After;
 import org.junit.Assert;
@@ -96,7 +95,7 @@ public class DataClientProviderTest {
       SyncDataClient client = null;
       try {
         client = provider.getSyncDataClient(node, 100);
-      } catch (TException e) {
+      } catch (IOException e) {
         Assert.fail(e.getMessage());
       } finally {
         ClientUtils.putBackSyncClient(client);
@@ -135,7 +134,7 @@ public class DataClientProviderTest {
       SyncDataClient client = null;
       try {
         client = provider.getSyncDataClient(node, 100);
-      } catch (TException e) {
+      } catch (IOException e) {
         Assert.fail(e.getMessage());
       }
       assertNotNull(client);
diff --git 
a/example/session/src/main/java/org/apache/iotdb/SessionConcurrentExample.java 
b/example/session/src/main/java/org/apache/iotdb/SessionConcurrentExample.java
new file mode 100644
index 0000000..fb0f793
--- /dev/null
+++ 
b/example/session/src/main/java/org/apache/iotdb/SessionConcurrentExample.java
@@ -0,0 +1,198 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb;
+
+import org.apache.iotdb.rpc.IoTDBConnectionException;
+import org.apache.iotdb.rpc.StatementExecutionException;
+import org.apache.iotdb.session.Session;
+import org.apache.iotdb.tsfile.file.metadata.enums.CompressionType;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
+import org.apache.iotdb.tsfile.file.metadata.enums.TSEncoding;
+import org.apache.iotdb.tsfile.write.record.Tablet;
+import org.apache.iotdb.tsfile.write.schema.MeasurementSchema;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Random;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+
+public class SessionConcurrentExample {
+
+  private static final int sgNum = 20;
+  private static final int deviceNum = 100;
+  private static final int parallelDegreeForOneSG = 3;
+
+  public static void main(String[] args)
+      throws IoTDBConnectionException, StatementExecutionException {
+
+    Session session = new Session("127.0.0.1", 6667, "root", "root");
+    session.open(false);
+    createTemplate(session);
+    session.close();
+
+    CountDownLatch latch = new CountDownLatch(sgNum * parallelDegreeForOneSG);
+    ExecutorService es = Executors.newFixedThreadPool(sgNum * 
parallelDegreeForOneSG);
+
+    for (int i = 0; i < sgNum * parallelDegreeForOneSG; i++) {
+      int currentIndex = i;
+      es.execute(() -> concurrentOperation(latch, currentIndex));
+    }
+
+    es.shutdown();
+
+    try {
+      latch.await();
+    } catch (InterruptedException e) {
+      e.printStackTrace();
+    }
+  }
+
+  private static void concurrentOperation(CountDownLatch latch, int 
currentIndex) {
+
+    Session tmpSession = new Session("127.0.0.1", 6667, "root", "root");
+    try {
+      tmpSession.open(false);
+    } catch (IoTDBConnectionException e) {
+      e.printStackTrace();
+    }
+
+    for (int j = 0; j < deviceNum; j++) {
+      try {
+        insertTablet(
+            tmpSession, String.format("root.sg_%d.d_%d", currentIndex / 
parallelDegreeForOneSG, j));
+      } catch (IoTDBConnectionException | StatementExecutionException e) {
+        e.printStackTrace();
+      }
+    }
+
+    try {
+      tmpSession.close();
+    } catch (IoTDBConnectionException e) {
+      e.printStackTrace();
+    }
+
+    latch.countDown();
+  }
+
+  private static void createTemplate(Session session)
+      throws IoTDBConnectionException, StatementExecutionException {
+    List<List<String>> measurementList = new ArrayList<>();
+    measurementList.add(Collections.singletonList("s1"));
+    measurementList.add(Collections.singletonList("s2"));
+    measurementList.add(Collections.singletonList("s3"));
+
+    List<List<TSDataType>> dataTypeList = new ArrayList<>();
+    dataTypeList.add(Collections.singletonList(TSDataType.INT64));
+    dataTypeList.add(Collections.singletonList(TSDataType.INT64));
+    dataTypeList.add(Collections.singletonList(TSDataType.INT64));
+
+    List<List<TSEncoding>> encodingList = new ArrayList<>();
+    encodingList.add(Collections.singletonList(TSEncoding.RLE));
+    encodingList.add(Collections.singletonList(TSEncoding.RLE));
+    encodingList.add(Collections.singletonList(TSEncoding.RLE));
+
+    List<CompressionType> compressionTypes = new ArrayList<>();
+    for (int i = 0; i < 3; i++) {
+      compressionTypes.add(CompressionType.SNAPPY);
+    }
+    List<String> schemaNames = new ArrayList<>();
+    schemaNames.add("s1");
+    schemaNames.add("s2");
+    schemaNames.add("s3");
+
+    session.createSchemaTemplate(
+        "template1", schemaNames, measurementList, dataTypeList, encodingList, 
compressionTypes);
+    for (int i = 0; i < sgNum; i++) {
+      session.setSchemaTemplate("template1", "root.sg_" + i);
+    }
+  }
+
+  /**
+   * insert the data of a device. For each timestamp, the number of 
measurements is the same.
+   *
+   * <p>Users need to control the count of Tablet and write a batch when it 
reaches the maxBatchSize
+   */
+  private static void insertTablet(Session session, String deviceId)
+      throws IoTDBConnectionException, StatementExecutionException {
+    /*
+     * A Tablet example:
+     *      device1
+     * time s1, s2, s3
+     * 1,   1,  1,  1
+     * 2,   2,  2,  2
+     * 3,   3,  3,  3
+     */
+    // The schema of measurements of one device
+    // only measurementId and data type in MeasurementSchema take effects in 
Tablet
+    List<MeasurementSchema> schemaList = new ArrayList<>();
+    schemaList.add(new MeasurementSchema("s1", TSDataType.INT64));
+    schemaList.add(new MeasurementSchema("s2", TSDataType.INT64));
+    schemaList.add(new MeasurementSchema("s3", TSDataType.INT64));
+
+    Tablet tablet = new Tablet(deviceId, schemaList, 100);
+
+    // Method 1 to add tablet data
+    long timestamp = System.currentTimeMillis();
+
+    for (long row = 0; row < 100; row++) {
+      int rowIndex = tablet.rowSize++;
+      tablet.addTimestamp(rowIndex, timestamp);
+      for (int s = 0; s < 3; s++) {
+        long value = new Random().nextLong();
+        tablet.addValue(schemaList.get(s).getMeasurementId(), rowIndex, value);
+      }
+      if (tablet.rowSize == tablet.getMaxRowNumber()) {
+        session.insertTablet(tablet, true);
+        tablet.reset();
+      }
+      timestamp++;
+    }
+
+    if (tablet.rowSize != 0) {
+      session.insertTablet(tablet);
+      tablet.reset();
+    }
+
+    // Method 2 to add tablet data
+    long[] timestamps = tablet.timestamps;
+    Object[] values = tablet.values;
+
+    for (long time = 0; time < 100; time++) {
+      int row = tablet.rowSize++;
+      timestamps[row] = time;
+      for (int i = 0; i < 3; i++) {
+        long[] sensor = (long[]) values[i];
+        sensor[row] = i;
+      }
+      if (tablet.rowSize == tablet.getMaxRowNumber()) {
+        session.insertTablet(tablet, true);
+        tablet.reset();
+      }
+    }
+
+    if (tablet.rowSize != 0) {
+      session.insertTablet(tablet);
+      tablet.reset();
+    }
+  }
+}
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 0402f4b..94c63f0 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
@@ -192,7 +192,12 @@ public class TsFileProcessor {
 
     if (IoTDBDescriptor.getInstance().getConfig().isEnableWal()) {
       try {
+        long startTime = System.currentTimeMillis();
         getLogNode().write(insertRowPlan);
+        long elapsed = System.currentTimeMillis() - startTime;
+        if (elapsed > 5000) {
+          logger.error("write wal slowly : cost {}ms", elapsed);
+        }
       } catch (Exception e) {
         throw new WriteProcessException(
             String.format(
@@ -248,11 +253,16 @@ public class TsFileProcessor {
     }
     try {
       workMemTable.insertTablet(insertTabletPlan, start, end);
+      long startTime = System.currentTimeMillis();
       if (IoTDBDescriptor.getInstance().getConfig().isEnableWal()) {
         insertTabletPlan.setStart(start);
         insertTabletPlan.setEnd(end);
         getLogNode().write(insertTabletPlan);
       }
+      long elapsed = System.currentTimeMillis() - startTime;
+      if (elapsed > 5000) {
+        logger.error("write wal slowly : cost {}ms", elapsed);
+      }
     } catch (Exception e) {
       for (int i = start; i < end; i++) {
         results[i] = RpcUtils.getStatus(TSStatusCode.INTERNAL_SERVER_ERROR, 
e.getMessage());
diff --git a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java 
b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java
index 674c90e..7dad6a9 100644
--- a/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java
+++ b/server/src/main/java/org/apache/iotdb/db/metadata/MManager.java
@@ -2113,7 +2113,15 @@ public class MManager {
   }
 
   private void setUsingDeviceTemplate(SetUsingDeviceTemplatePlan plan) throws 
MetadataException {
-    getDeviceNode(plan.getPrefixPath()).setUseTemplate(true);
+    try {
+      getDeviceNode(plan.getPrefixPath()).setUseTemplate(true);
+    } catch (PathNotExistException e) {
+      // the order of SetUsingDeviceTemplatePlan and AutoCreateDeviceMNodePlan 
cannot be guaranteed
+      // during writing currently, so we need a auto-create mechanism here
+      mtree.getDeviceNodeWithAutoCreating(
+          plan.getPrefixPath(), config.getDefaultStorageGroupLevel());
+      getDeviceNode(plan.getPrefixPath()).setUseTemplate(true);
+    }
   }
 
   public long getTotalSeriesNumber() {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/rescon/PrimitiveArrayManager.java 
b/server/src/main/java/org/apache/iotdb/db/rescon/PrimitiveArrayManager.java
index 9881e56..76d2128 100644
--- a/server/src/main/java/org/apache/iotdb/db/rescon/PrimitiveArrayManager.java
+++ b/server/src/main/java/org/apache/iotdb/db/rescon/PrimitiveArrayManager.java
@@ -40,9 +40,17 @@ public class PrimitiveArrayManager {
 
   public static final int ARRAY_SIZE = CONFIG.getPrimitiveArraySize();
 
+  /**
+   * The actual used memory will be 50% larger than the statistic, so we need 
to limit the size of
+   * POOLED_ARRAYS_MEMORY_THRESHOLD, make it smaller than its actual allowed 
value.
+   */
+  private static final double AMPLIFICATION_FACTOR = 1.5;
+
   /** threshold total size of arrays for all data types */
   private static final double POOLED_ARRAYS_MEMORY_THRESHOLD =
-      CONFIG.getAllocateMemoryForWrite() * 
CONFIG.getBufferedArraysMemoryProportion();
+      CONFIG.getAllocateMemoryForWrite()
+          * CONFIG.getBufferedArraysMemoryProportion()
+          / AMPLIFICATION_FACTOR;
 
   /** TSDataType#serialize() -> ArrayDeque<Array> */
   private static final ArrayDeque[] POOLED_ARRAYS = new 
ArrayDeque[TSDataType.values().length];
diff --git 
a/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java 
b/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java
index 5be9dca..d1bf262 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java
@@ -575,12 +575,12 @@ public class TSServiceImpl implements TSIService.Iface {
 
   @Override
   public TSExecuteStatementResp executeStatement(TSExecuteStatementReq req) {
+    String statement = req.getStatement();
     try {
       if (!checkLogin(req.getSessionId())) {
         return 
RpcUtils.getTSExecuteStatementResp(TSStatusCode.NOT_LOGIN_ERROR);
       }
 
-      String statement = req.getStatement();
       PhysicalPlan physicalPlan =
           processor.parseSQLToPhysicalPlan(
               statement, sessionManager.getZoneId(req.getSessionId()), 
req.fetchSize);
@@ -598,9 +598,11 @@ public class TSServiceImpl implements TSIService.Iface {
     } catch (InterruptedException e) {
       LOGGER.error(INFO_INTERRUPT_ERROR, req, e);
       Thread.currentThread().interrupt();
-      return RpcUtils.getTSExecuteStatementResp(onQueryException(e, "executing 
executeStatement"));
+      return RpcUtils.getTSExecuteStatementResp(
+          onQueryException(e, "executing \"" + statement + "\""));
     } catch (Exception e) {
-      return RpcUtils.getTSExecuteStatementResp(onQueryException(e, "executing 
executeStatement"));
+      return RpcUtils.getTSExecuteStatementResp(
+          onQueryException(e, "executing \"" + statement + "\""));
     }
   }
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
 
b/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
index c947963..e982c77 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/writelog/node/ExclusiveWriteLogNode.java
@@ -28,7 +28,6 @@ import org.apache.iotdb.db.writelog.io.ILogWriter;
 import org.apache.iotdb.db.writelog.io.LogWriter;
 import org.apache.iotdb.db.writelog.io.MultiFileLogReader;
 
-import com.google.common.util.concurrent.ThreadFactoryBuilder;
 import org.apache.commons.io.FileUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -43,6 +42,7 @@ import java.util.Arrays;
 import java.util.Comparator;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
+import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.locks.ReentrantLock;
 
 /** This WriteLogNode is used to manage insert ahead logs of a TsFile. */
@@ -51,33 +51,32 @@ public class ExclusiveWriteLogNode implements WriteLogNode, 
Comparable<Exclusive
   public static final String WAL_FILE_NAME = "wal";
   private static final Logger logger = 
LoggerFactory.getLogger(ExclusiveWriteLogNode.class);
 
-  private String identifier;
+  private final String identifier;
 
-  private String logDirectory;
+  private final String logDirectory;
 
   private ILogWriter currentFileWriter;
 
-  private IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
+  private final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
 
-  private ByteBuffer logBufferWorking;
-  private ByteBuffer logBufferIdle;
-  private ByteBuffer logBufferFlushing;
+  private volatile ByteBuffer logBufferWorking;
+  private volatile ByteBuffer logBufferIdle;
+  private volatile ByteBuffer logBufferFlushing;
 
   // used for the convenience of deletion
-  private ByteBuffer[] bufferArray;
+  private volatile ByteBuffer[] bufferArray;
 
   private final Object switchBufferCondition = new Object();
-  private ReentrantLock lock = new ReentrantLock();
-  private static final ExecutorService FLUSH_BUFFER_THREAD_POOL =
-      Executors.newCachedThreadPool(
-          new 
ThreadFactoryBuilder().setNameFormat("Flush-WAL-Thread-%d").setDaemon(true).build());
+  private final ReentrantLock lock = new ReentrantLock();
+  private final ExecutorService FLUSH_BUFFER_THREAD_POOL =
+      Executors.newSingleThreadExecutor(r -> new Thread(r, "Flush-WAL-Thread-" 
+ this.hashCode()));
 
   private long fileId = 0;
   private long lastFlushedId = 0;
 
   private int bufferedLogNum = 0;
 
-  private boolean deleted;
+  private final AtomicBoolean deleted = new AtomicBoolean(false);
 
   /**
    * constructor of ExclusiveWriteLogNode.
@@ -102,7 +101,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode, 
Comparable<Exclusive
 
   @Override
   public void write(PhysicalPlan plan) throws IOException {
-    if (deleted) {
+    if (deleted.get()) {
       throw new IOException("WAL node deleted");
     }
     lock.lock();
@@ -138,7 +137,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode, 
Comparable<Exclusive
     lock.lock();
     try {
       synchronized (switchBufferCondition) {
-        while (logBufferFlushing != null && !deleted) {
+        while (logBufferFlushing != null && !deleted.get()) {
           switchBufferCondition.wait();
         }
         switchBufferCondition.notifyAll();
@@ -151,7 +150,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode, 
Comparable<Exclusive
       }
       logger.debug("Log node {} closed successfully", identifier);
     } catch (IOException e) {
-      logger.error("Cannot close log node {} because:", identifier, e);
+      logger.warn("Cannot close log node {} because:", identifier, e);
     } catch (InterruptedException e) {
       Thread.currentThread().interrupt();
       logger.warn("Waiting for current buffer being flushed interrupted");
@@ -162,7 +161,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode, 
Comparable<Exclusive
 
   @Override
   public void forceSync() {
-    if (deleted) {
+    if (deleted.get()) {
       return;
     }
     sync();
@@ -208,7 +207,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode, 
Comparable<Exclusive
     try {
       close();
       
FileUtils.deleteDirectory(SystemFileFactory.INSTANCE.getFile(logDirectory));
-      deleted = true;
+      deleted.set(true);
       return this.bufferArray;
     } finally {
       lock.unlock();
@@ -232,7 +231,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode, 
Comparable<Exclusive
         FileUtils.forceDelete(logFile);
         logger.info("Log node {} cleaned old file", identifier);
       } catch (IOException e) {
-        logger.error("Old log file {} of {} cannot be deleted", 
logFile.getName(), identifier, e);
+        logger.warn("Old log file {} of {} cannot be deleted", 
logFile.getName(), identifier, e);
       }
     }
   }
@@ -245,7 +244,7 @@ public class ExclusiveWriteLogNode implements WriteLogNode, 
Comparable<Exclusive
           currentFileWriter.force();
         }
       } catch (IOException e) {
-        logger.error("Log node {} force failed.", identifier, e);
+        logger.warn("Log node {} force failed.", identifier, e);
       }
     } finally {
       lock.unlock();
@@ -261,8 +260,6 @@ public class ExclusiveWriteLogNode implements WriteLogNode, 
Comparable<Exclusive
       switchBufferWorkingToFlushing();
       ILogWriter currWriter = getCurrentFileWriter();
       FLUSH_BUFFER_THREAD_POOL.submit(() -> flushBuffer(currWriter));
-      switchBufferIdleToWorking();
-
       bufferedLogNum = 0;
       logger.debug("Log node {} ends sync.", identifier);
     } catch (InterruptedException e) {
@@ -281,50 +278,28 @@ public class ExclusiveWriteLogNode implements 
WriteLogNode, Comparable<Exclusive
     } catch (ClosedChannelException e) {
       // ignore
     } catch (IOException e) {
-      logger.error("Log node {} sync failed, change system mode to read-only", 
identifier, e);
+      logger.warn("Log node {} sync failed, change system mode to read-only", 
identifier, e);
       IoTDBDescriptor.getInstance().getConfig().setReadOnly(true);
       return;
     }
-    logBufferFlushing.clear();
-
-    try {
-      switchBufferFlushingToIdle();
-    } catch (InterruptedException e) {
-      Thread.currentThread().interrupt();
-    }
-  }
 
-  private void switchBufferWorkingToFlushing() throws InterruptedException {
+    // switch buffer flushing to idle and notify the sync thread
     synchronized (switchBufferCondition) {
-      while (logBufferFlushing != null && !deleted) {
-        switchBufferCondition.wait();
-      }
-      logBufferFlushing = logBufferWorking;
-      logBufferWorking = null;
+      logBufferIdle = logBufferFlushing;
+      logBufferFlushing = null;
       switchBufferCondition.notifyAll();
     }
   }
 
-  private void switchBufferIdleToWorking() throws InterruptedException {
+  private void switchBufferWorkingToFlushing() throws InterruptedException {
     synchronized (switchBufferCondition) {
-      while (logBufferIdle == null && !deleted) {
-        switchBufferCondition.wait();
+      while (logBufferFlushing != null && !deleted.get()) {
+        switchBufferCondition.wait(100);
       }
+      logBufferFlushing = logBufferWorking;
       logBufferWorking = logBufferIdle;
+      logBufferWorking.clear();
       logBufferIdle = null;
-      switchBufferCondition.notifyAll();
-    }
-  }
-
-  private void switchBufferFlushingToIdle() throws InterruptedException {
-    synchronized (switchBufferCondition) {
-      while (logBufferIdle != null && !deleted) {
-        switchBufferCondition.wait();
-      }
-      logBufferIdle = logBufferFlushing;
-      logBufferIdle.clear();
-      logBufferFlushing = null;
-      switchBufferCondition.notifyAll();
     }
   }
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/writelog/recover/TsFileRecoverPerformer.java
 
b/server/src/main/java/org/apache/iotdb/db/writelog/recover/TsFileRecoverPerformer.java
index 72ea41c..da53108 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/writelog/recover/TsFileRecoverPerformer.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/writelog/recover/TsFileRecoverPerformer.java
@@ -41,6 +41,8 @@ import org.slf4j.LoggerFactory;
 import java.io.File;
 import java.io.IOException;
 import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Map.Entry;
@@ -188,13 +190,24 @@ public class TsFileRecoverPerformer {
     for (Map.Entry<String, List<ChunkMetadata>> entry : 
deviceChunkMetaDataMap.entrySet()) {
       String deviceId = entry.getKey();
       List<ChunkMetadata> chunkMetadataList = entry.getValue();
-      TSDataType dataType = entry.getValue().get(entry.getValue().size() - 
1).getDataType();
-      for (ChunkMetadata chunkMetaData : chunkMetadataList) {
-        if (!chunkMetaData.getDataType().equals(dataType)) {
-          continue;
+
+      Map<String, List<ChunkMetadata>> measurementToChunkMetadatas = new 
HashMap<>();
+      for (ChunkMetadata chunkMetadata : chunkMetadataList) {
+        List<ChunkMetadata> list =
+            measurementToChunkMetadatas.computeIfAbsent(
+                chunkMetadata.getMeasurementUid(), n -> new ArrayList<>());
+        list.add(chunkMetadata);
+      }
+
+      for (List<ChunkMetadata> metadataList : 
measurementToChunkMetadatas.values()) {
+        TSDataType dataType = metadataList.get(metadataList.size() - 
1).getDataType();
+        for (ChunkMetadata chunkMetaData : chunkMetadataList) {
+          if (!chunkMetaData.getDataType().equals(dataType)) {
+            continue;
+          }
+          tsFileResource.updateStartTime(deviceId, 
chunkMetaData.getStartTime());
+          tsFileResource.updateEndTime(deviceId, chunkMetaData.getEndTime());
         }
-        tsFileResource.updateStartTime(deviceId, chunkMetaData.getStartTime());
-        tsFileResource.updateEndTime(deviceId, chunkMetaData.getEndTime());
       }
     }
     
tsFileResource.updatePlanIndexes(restorableTsFileIOWriter.getMinPlanIndex());
diff --git 
a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java 
b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java
index d01b9d9..fa4f7f5 100644
--- a/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java
+++ b/server/src/test/java/org/apache/iotdb/db/integration/IoTDBRestartIT.java
@@ -18,6 +18,8 @@
  */
 package org.apache.iotdb.db.integration;
 
+import org.apache.iotdb.db.conf.IoTDBConfig;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.engine.StorageEngine;
 import org.apache.iotdb.db.engine.compaction.CompactionMergeTaskPoolManager;
 import org.apache.iotdb.db.exception.StorageEngineException;
@@ -353,6 +355,52 @@ public class IoTDBRestartIT {
   }
 
   @Test
+  public void testRecoverWALDeleteSchemaCheckResourceTime() throws Exception {
+    EnvironmentUtils.envSetUp();
+    Class.forName(Config.JDBC_DRIVER_NAME);
+    IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
+    int avgSeriesPointNumberThreshold = 
config.getAvgSeriesPointNumberThreshold();
+    config.setAvgSeriesPointNumberThreshold(2);
+    long tsfileSize = config.getSeqTsFileSize();
+    config.setSeqTsFileSize(10000000);
+
+    try (Connection connection =
+            DriverManager.getConnection(
+                Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+        Statement statement = connection.createStatement()) {
+      statement.execute("create timeseries root.turbine1.d1.s1 with 
datatype=INT64");
+      statement.execute("insert into root.turbine1.d1(timestamp,s1) 
values(1,1)");
+      statement.execute("insert into root.turbine1.d1(timestamp,s1) 
values(2,1)");
+      statement.execute("create timeseries root.turbine1.d1.s2 with 
datatype=BOOLEAN");
+      statement.execute("insert into root.turbine1.d1(timestamp,s2) 
values(3,true)");
+      statement.execute("insert into root.turbine1.d1(timestamp,s2) 
values(4,true)");
+    }
+
+    Thread.sleep(1000);
+    EnvironmentUtils.restartDaemon();
+
+    try (Connection connection =
+            DriverManager.getConnection(
+                Config.IOTDB_URL_PREFIX + "127.0.0.1:6667/", "root", "root");
+        Statement statement = connection.createStatement()) {
+
+      long[] result = new long[] {1L, 2L};
+      statement.execute("select s1 from root.turbine1.d1 where time < 3");
+      ResultSet resultSet = statement.getResultSet();
+      int cnt = 0;
+      while (resultSet.next()) {
+        assertEquals(resultSet.getLong(1), result[cnt]);
+        cnt++;
+      }
+      assertEquals(2, cnt);
+    }
+
+    config.setAvgSeriesPointNumberThreshold(avgSeriesPointNumberThreshold);
+    config.setSeqTsFileSize(tsfileSize);
+    EnvironmentUtils.cleanEnv();
+  }
+
+  @Test
   public void testRestartCompaction()
       throws SQLException, ClassNotFoundException, IOException, 
StorageEngineException {
     EnvironmentUtils.envSetUp();
diff --git 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/expression/util/ExpressionOptimizer.java
 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/expression/util/ExpressionOptimizer.java
index ea2a24f..b056286 100644
--- 
a/tsfile/src/main/java/org/apache/iotdb/tsfile/read/expression/util/ExpressionOptimizer.java
+++ 
b/tsfile/src/main/java/org/apache/iotdb/tsfile/read/expression/util/ExpressionOptimizer.java
@@ -70,21 +70,25 @@ public class ExpressionOptimizer {
             (GlobalTimeExpression) right, left, selectedSeries, relation);
       } else if (left.getType() != ExpressionType.GLOBAL_TIME
           && right.getType() != ExpressionType.GLOBAL_TIME) {
-        IExpression regularLeft = optimize(left, selectedSeries);
-        IExpression regularRight = optimize(right, selectedSeries);
-        IBinaryExpression midRet = null;
-        if (relation == ExpressionType.AND) {
-          midRet = BinaryExpression.and(regularLeft, regularRight);
-        } else if (relation == ExpressionType.OR) {
-          midRet = BinaryExpression.or(regularLeft, regularRight);
-        } else {
-          throw new UnsupportedOperationException("unsupported IExpression 
type: " + relation);
-        }
-        if (midRet.getLeft().getType() == ExpressionType.GLOBAL_TIME
-            || midRet.getRight().getType() == ExpressionType.GLOBAL_TIME) {
-          return optimize(midRet, selectedSeries);
-        } else {
-          return midRet;
+        try {
+          IExpression regularLeft = optimize(left, selectedSeries);
+          IExpression regularRight = optimize(right, selectedSeries);
+          IBinaryExpression midRet = null;
+          if (relation == ExpressionType.AND) {
+            midRet = BinaryExpression.and(regularLeft, regularRight);
+          } else if (relation == ExpressionType.OR) {
+            midRet = BinaryExpression.or(regularLeft, regularRight);
+          } else {
+            throw new UnsupportedOperationException("unsupported IExpression 
type: " + relation);
+          }
+          if (midRet.getLeft().getType() == ExpressionType.GLOBAL_TIME
+              || midRet.getRight().getType() == ExpressionType.GLOBAL_TIME) {
+            return optimize(midRet, selectedSeries);
+          } else {
+            return midRet;
+          }
+        } catch (StackOverflowError stackOverflowError) {
+          throw new QueryFilterOptimizationException("StackOverflowError is 
encountered.");
         }
       }
     }

Reply via email to