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

tanxinyu pushed a commit to branch fix_thrift_out_of_sequence
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 1252a68e872e8e2b8e3ce4a8bea0e78db9785870
Author: LebronAl <[email protected]>
AuthorDate: Thu Jul 15 19:03:39 2021 +0800

    fix
---
 .../iotdb/cluster/client/DataClientProvider.java   |  9 ++-
 .../iotdb/cluster/coordinator/Coordinator.java     | 17 ++----
 .../apache/iotdb/cluster/metadata/CMManager.java   | 65 +++++++++++++++-------
 .../apache/iotdb/cluster/metadata/MetaPuller.java  | 17 ++++--
 .../iotdb/cluster/query/ClusterPlanExecutor.java   | 35 +++++++++---
 .../cluster/query/aggregate/ClusterAggregator.java |  8 ++-
 .../cluster/query/fill/ClusterPreviousFill.java    | 25 ++++++---
 .../query/groupby/RemoteGroupByExecutor.java       | 21 +++++--
 .../query/last/ClusterLastQueryExecutor.java       | 25 ++++++---
 .../cluster/query/reader/ClusterReaderFactory.java |  8 ++-
 .../iotdb/cluster/query/reader/DataSourceInfo.java | 15 +++--
 .../reader/RemoteSeriesReaderByTimestamp.java      |  1 +
 .../query/reader/RemoteSimpleSeriesReader.java     |  1 +
 .../query/reader/mult/MultDataSourceInfo.java      | 14 +++--
 .../query/reader/mult/RemoteMultSeriesReader.java  | 16 +++---
 .../cluster/server/member/MetaGroupMember.java     | 11 ++--
 16 files changed, 193 insertions(+), 95 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/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..9a3d69d 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,13 @@ 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) {
+              syncDataClient.getInputProtocol().getTransport().close();
+              throw e;
+            }
           }
         }
         if (result != null) {
@@ -978,13 +983,17 @@ 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) {
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
         }
       }
     }
@@ -1212,8 +1221,12 @@ 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) {
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
+        }
       }
     }
 
@@ -1338,8 +1351,12 @@ 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) {
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
+        }
       }
     }
     return paths;
@@ -1792,9 +1809,14 @@ 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) {
+          syncDataClient.getInputProtocol().getTransport().close();
+          throw e;
+        }
       }
     }
     return resultBinary;
@@ -1818,9 +1840,14 @@ 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) {
+          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..a772259 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,17 @@ 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) {
+          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..980ba38 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,12 @@ 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) {
+              syncDataClient.getInputProtocol().getTransport().close();
+              throw e;
+            }
           }
         }
         logger.debug(
@@ -363,8 +367,13 @@ 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) {
+              syncDataClient.getInputProtocol().getTransport().close();
+              throw e;
+            }
           }
         }
         if (paths != null) {
@@ -449,8 +458,13 @@ 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) {
+              syncDataClient.getInputProtocol().getTransport().close();
+              throw e;
+            }
           }
         }
         if (nextChildrenNodes != null) {
@@ -558,8 +572,13 @@ 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) {
+              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..e8ec8e2 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,12 @@ public class ClusterAggregator {
           metaGroupMember
               .getClientProvider()
               .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
-
-        resultBuffers = syncDataClient.getAggrResult(request);
+        try {
+          resultBuffers = syncDataClient.getAggrResult(request);
+        } catch (TException e) {
+          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..9c46845 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 (TException | IOException e) {
       logger.error(
           "{}: Cannot perform previous fill of {} to {}",
           metaGroupMember.getName(),
           arguments.getPath(),
           node,
           e);
+      if (e instanceof TException) {
+        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..30c115c 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,13 @@ 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();
+          }
         }
       }
     } catch (TException e) {
@@ -118,7 +123,7 @@ public class RemoteGroupByExecutor implements 
GroupByExecutor {
   @Override
   public Pair<Long, Object> peekNextNotNullValue(long nextStartTime, long 
nextEndTime)
       throws IOException {
-    ByteBuffer aggrBuffer;
+    ByteBuffer aggrBuffer = null;
     try {
       if (ClusterDescriptor.getInstance().getConfig().isUseAsyncServer()) {
         AsyncDataClient client =
@@ -133,9 +138,13 @@ 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();
+          }
         }
       }
     } 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..bea7c1d 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,29 @@ 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) {
+        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..44767b7 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
@@ -156,10 +156,15 @@ public class DataSourceInfo {
       throws TException {
 
     Long newReaderId;
-    try (SyncDataClient client =
-        this.metaGroupMember
-            .getClientProvider()
-            .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())) {
+    try {
+      SyncDataClient client =
+          this.metaGroupMember
+              .getClientProvider()
+              .getSyncDataClient(node, RaftServer.getReadOperationTimeoutMS())
+    }catch (IOException e){
+      return null;
+    }
+    try () {
 
       if (byTimestamp) {
         newReaderId = client.querySingleSeriesByTimestamp(request);
@@ -201,7 +206,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..2431951 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,7 @@ public class RemoteSeriesReaderByTimestamp implements 
IReaderByTimestamp {
       return curSyncClient.fetchSingleSeriesByTimestamps(
           sourceInfo.getHeader(), sourceInfo.getReaderId(), timestampList);
     } catch (TException e) {
+      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..70a0641 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,7 @@ public class RemoteSimpleSeriesReader implements 
IPointReader {
       curSyncClient = 
sourceInfo.getCurSyncClient(RaftServer.getReadOperationTimeoutMS());
       return curSyncClient.fetchSingleSeries(sourceInfo.getHeader(), 
sourceInfo.getReaderId());
     } catch (TException e) {
+      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..8ffe4ce2 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 = null;
     try (SyncDataClient client =
         this.metaGroupMember
             .getClientProvider()
@@ -189,7 +188,12 @@ 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();
+      }
       return newReaderId;
     }
   }
@@ -212,7 +216,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..31f7caa 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,17 @@ 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) {
+        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/member/MetaGroupMember.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java
index cede9a3..e5f0a0f 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
@@ -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);
       // 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);
+      }
     }
   }
 

Reply via email to