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

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

commit f4e6435d6249aab60d70d3156e6bcaf6bf0d3828
Author: LebronAl <[email protected]>
AuthorDate: Sun Jul 11 21:16:39 2021 +0800

    init
---
 .../iotdb/cluster/client/DataClientProvider.java   |  9 +++++
 .../iotdb/cluster/client/sync/SyncClientPool.java  | 14 +++++++
 .../iotdb/cluster/server/DataClusterServer.java    |  7 ++++
 .../iotdb/cluster/server/MetaClusterServer.java    |  7 ++++
 .../cluster/server/member/MetaGroupMember.java     | 46 ++++++++++++++++++++++
 .../cluster/server/service/BaseAsyncService.java   |  4 ++
 .../cluster/server/service/BaseSyncService.java    |  4 ++
 .../org/apache/iotdb/db/metadata/MManager.java     |  3 ++
 .../apache/iotdb/session/IoTDBSessionSimpleIT.java |  3 +-
 thrift-cluster/src/main/thrift/cluster.thrift      |  4 ++
 10 files changed, 100 insertions(+), 1 deletion(-)

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 8b954ec..b800454 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
@@ -92,4 +92,13 @@ public class DataClientProvider {
     client.setTimeout(timeout);
     return client;
   }
+
+  public SyncDataClient getSyncDataClientForRefresh(Node node, int timeout) 
throws TException {
+    SyncDataClient client = (SyncDataClient) 
getDataSyncClientPool().getClientForRefresh(node);
+    if (client == null) {
+      throw new TException("can not get client for node=" + 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 2c279c0..77ee1e7 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
@@ -61,6 +61,20 @@ public class SyncClientPool {
     return getClient(node, true);
   }
 
+  public Client getClientForRefresh(Node node) {
+    ClusterNode clusterNode = new ClusterNode(node);
+
+    // As clientCaches is ConcurrentHashMap, computeIfAbsent is thread safety.
+    Deque<Client> clientStack = clientCaches.computeIfAbsent(clusterNode, n -> 
new ArrayDeque<>());
+    synchronized (clientStack) {
+      if (clientStack.isEmpty()) {
+        return null;
+      } else {
+        return clientStack.poll();
+      }
+    }
+  }
+
   /**
    * Get a client of the given node from the cache if one is available, or 
create a new one.
    *
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/DataClusterServer.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/DataClusterServer.java
index e4c81f8..1a0ea16 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/DataClusterServer.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/DataClusterServer.java
@@ -32,6 +32,7 @@ import 
org.apache.iotdb.cluster.partition.slot.SlotPartitionTable;
 import org.apache.iotdb.cluster.rpc.thrift.AppendEntriesRequest;
 import org.apache.iotdb.cluster.rpc.thrift.AppendEntryRequest;
 import org.apache.iotdb.cluster.rpc.thrift.ElectionRequest;
+import org.apache.iotdb.cluster.rpc.thrift.EmptyReuqest;
 import org.apache.iotdb.cluster.rpc.thrift.ExecutNonQueryReq;
 import org.apache.iotdb.cluster.rpc.thrift.GetAggrResultRequest;
 import org.apache.iotdb.cluster.rpc.thrift.GetAllPathsResult;
@@ -921,6 +922,12 @@ public class DataClusterServer extends RaftServer
   }
 
   @Override
+  public void RefreshClient(EmptyReuqest request, AsyncMethodCallback<Void> 
resultHandler) {}
+
+  @Override
+  public void RefreshClient(EmptyReuqest request) {}
+
+  @Override
   public RequestCommitIndexResponse requestCommitIndex(Node header) throws 
TException {
     return getDataSyncService(header).requestCommitIndex(header);
   }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/MetaClusterServer.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/MetaClusterServer.java
index abb8020..10f8622 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/MetaClusterServer.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/MetaClusterServer.java
@@ -29,6 +29,7 @@ import 
org.apache.iotdb.cluster.rpc.thrift.AppendEntriesRequest;
 import org.apache.iotdb.cluster.rpc.thrift.AppendEntryRequest;
 import org.apache.iotdb.cluster.rpc.thrift.CheckStatusResponse;
 import org.apache.iotdb.cluster.rpc.thrift.ElectionRequest;
+import org.apache.iotdb.cluster.rpc.thrift.EmptyReuqest;
 import org.apache.iotdb.cluster.rpc.thrift.ExecutNonQueryReq;
 import org.apache.iotdb.cluster.rpc.thrift.HeartBeatRequest;
 import org.apache.iotdb.cluster.rpc.thrift.HeartBeatResponse;
@@ -226,6 +227,9 @@ public class MetaClusterServer extends RaftServer
   }
 
   @Override
+  public void RefreshClient(EmptyReuqest request, AsyncMethodCallback<Void> 
resultHandler) {}
+
+  @Override
   public void requestCommitIndex(
       Node header, AsyncMethodCallback<RequestCommitIndexResponse> 
resultHandler) {
     asyncService.requestCommitIndex(header, resultHandler);
@@ -334,6 +338,9 @@ public class MetaClusterServer extends RaftServer
   }
 
   @Override
+  public void RefreshClient(EmptyReuqest request) {}
+
+  @Override
   public RequestCommitIndexResponse requestCommitIndex(Node header) throws 
TException {
     return syncService.requestCommitIndex(header);
   }
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 58884e1..63a743c 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
@@ -53,9 +53,11 @@ import org.apache.iotdb.cluster.query.ClusterPlanRouter;
 import org.apache.iotdb.cluster.rpc.thrift.AddNodeResponse;
 import org.apache.iotdb.cluster.rpc.thrift.AppendEntryRequest;
 import org.apache.iotdb.cluster.rpc.thrift.CheckStatusResponse;
+import org.apache.iotdb.cluster.rpc.thrift.EmptyReuqest;
 import org.apache.iotdb.cluster.rpc.thrift.HeartBeatRequest;
 import org.apache.iotdb.cluster.rpc.thrift.HeartBeatResponse;
 import org.apache.iotdb.cluster.rpc.thrift.Node;
+import org.apache.iotdb.cluster.rpc.thrift.RaftService;
 import org.apache.iotdb.cluster.rpc.thrift.SendSnapshotRequest;
 import org.apache.iotdb.cluster.rpc.thrift.StartUpStatus;
 import org.apache.iotdb.cluster.rpc.thrift.TSMetaService;
@@ -96,6 +98,7 @@ import org.apache.iotdb.service.rpc.thrift.EndPoint;
 import org.apache.iotdb.service.rpc.thrift.TSStatus;
 import org.apache.iotdb.tsfile.read.filter.basic.Filter;
 
+import com.google.common.util.concurrent.ThreadFactoryBuilder;
 import org.apache.thrift.TException;
 import org.apache.thrift.protocol.TProtocolFactory;
 import org.apache.thrift.transport.TTransportException;
@@ -163,6 +166,8 @@ public class MetaGroupMember extends RaftMember {
    */
   private static final int REPORT_INTERVAL_SEC = 10;
 
+  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 =
       ClusterDescriptor.getInstance().getConfig().getReplicationNum();
@@ -211,6 +216,8 @@ public class MetaGroupMember extends RaftMember {
 
   private DataClientProvider dataClientProvider;
 
+  private ScheduledExecutorService dataClientRefresher;
+
   /**
    * a single thread pool, every "REPORT_INTERVAL_SEC" seconds, "reportThread" 
will print the status
    * of all raft members in this node
@@ -331,6 +338,9 @@ public class MetaGroupMember extends RaftMember {
         Executors.newSingleThreadScheduledExecutor(n -> new Thread(n, 
"NodeReportThread"));
     hardLinkCleanerThread =
         Executors.newSingleThreadScheduledExecutor(n -> new Thread(n, 
"HardLinkCleaner"));
+    dataClientRefresher =
+        Executors.newSingleThreadScheduledExecutor(
+            new ThreadFactoryBuilder().setNameFormat(getName() + 
"-dataClientRefresher%d").build());
   }
 
   /**
@@ -349,6 +359,15 @@ public class MetaGroupMember extends RaftMember {
     if (clientServer != null) {
       clientServer.stop();
     }
+    if (dataClientRefresher != null) {
+      dataClientRefresher.shutdownNow();
+      try {
+        dataClientRefresher.awaitTermination(10, TimeUnit.SECONDS);
+      } catch (InterruptedException e) {
+        Thread.currentThread().interrupt();
+        logger.error("Unexpected interruption when waiting for reportThread to 
end", e);
+      }
+    }
     if (reportThread != null) {
       reportThread.shutdownNow();
       try {
@@ -460,6 +479,33 @@ public class MetaGroupMember extends RaftMember {
         CLEAN_HARDLINK_INTERVAL_SEC,
         CLEAN_HARDLINK_INTERVAL_SEC,
         TimeUnit.SECONDS);
+    dataClientRefresher.scheduleAtFixedRate(
+        this::RefreshAllClient, REFRESH_CLIENT_SEC, REFRESH_CLIENT_SEC, 
TimeUnit.SECONDS);
+  }
+
+  private void RefreshAllClient() {
+    for (Node receiver : allNodes) {
+      if (!receiver.equals(thisNode)) {
+        RaftService.Client client = null;
+        try {
+          client =
+              getClientProvider()
+                  .getSyncDataClientForRefresh(receiver, 
RaftServer.getWriteOperationTimeoutMS());
+        } catch (TException e) {
+          return;
+        }
+        try {
+          EmptyReuqest req = new EmptyReuqest();
+          client.RefreshClient(req);
+        } 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);
+        }
+      }
+    }
   }
 
   private void generateNodeReport() {
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/BaseAsyncService.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/BaseAsyncService.java
index 8673078..fe5decc 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/BaseAsyncService.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/BaseAsyncService.java
@@ -24,6 +24,7 @@ import 
org.apache.iotdb.cluster.exception.UnknownLogTypeException;
 import org.apache.iotdb.cluster.rpc.thrift.AppendEntriesRequest;
 import org.apache.iotdb.cluster.rpc.thrift.AppendEntryRequest;
 import org.apache.iotdb.cluster.rpc.thrift.ElectionRequest;
+import org.apache.iotdb.cluster.rpc.thrift.EmptyReuqest;
 import org.apache.iotdb.cluster.rpc.thrift.ExecutNonQueryReq;
 import org.apache.iotdb.cluster.rpc.thrift.HeartBeatRequest;
 import org.apache.iotdb.cluster.rpc.thrift.HeartBeatResponse;
@@ -145,6 +146,9 @@ public abstract class BaseAsyncService implements 
RaftService.AsyncIface {
   }
 
   @Override
+  public void RefreshClient(EmptyReuqest request, AsyncMethodCallback<Void> 
resultHandler) {}
+
+  @Override
   public void executeNonQueryPlan(
       ExecutNonQueryReq request, AsyncMethodCallback<TSStatus> resultHandler) {
     if (member.getCharacter() != NodeCharacter.LEADER) {
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/BaseSyncService.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/BaseSyncService.java
index ce200ab..aef2057 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/BaseSyncService.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/BaseSyncService.java
@@ -24,6 +24,7 @@ import 
org.apache.iotdb.cluster.exception.UnknownLogTypeException;
 import org.apache.iotdb.cluster.rpc.thrift.AppendEntriesRequest;
 import org.apache.iotdb.cluster.rpc.thrift.AppendEntryRequest;
 import org.apache.iotdb.cluster.rpc.thrift.ElectionRequest;
+import org.apache.iotdb.cluster.rpc.thrift.EmptyReuqest;
 import org.apache.iotdb.cluster.rpc.thrift.ExecutNonQueryReq;
 import org.apache.iotdb.cluster.rpc.thrift.HeartBeatRequest;
 import org.apache.iotdb.cluster.rpc.thrift.HeartBeatResponse;
@@ -152,6 +153,9 @@ public abstract class BaseSyncService implements 
RaftService.Iface {
   }
 
   @Override
+  public void RefreshClient(EmptyReuqest request) {}
+
+  @Override
   public TSStatus executeNonQueryPlan(ExecutNonQueryReq request) throws 
TException {
     if (member.getCharacter() != NodeCharacter.LEADER) {
       // forward the plan to the leader
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 072ecdd..6eda19a 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
@@ -984,6 +984,9 @@ public class MManager {
       throws MetadataException {
     MNode deviceMNode = getDeviceNode(device);
     MeasurementMNode measurementMNode = (MeasurementMNode) 
deviceMNode.getChild(measurement);
+    if (measurementMNode == null) {
+      return getSeriesSchema(device.concatNode(measurement));
+    }
     return measurementMNode.getSchema();
   }
 
diff --git 
a/session/src/test/java/org/apache/iotdb/session/IoTDBSessionSimpleIT.java 
b/session/src/test/java/org/apache/iotdb/session/IoTDBSessionSimpleIT.java
index a458a89..a5dc3d0 100644
--- a/session/src/test/java/org/apache/iotdb/session/IoTDBSessionSimpleIT.java
+++ b/session/src/test/java/org/apache/iotdb/session/IoTDBSessionSimpleIT.java
@@ -136,7 +136,8 @@ public class IoTDBSessionSimpleIT {
       tablet.reset();
     }
 
-    SessionDataSet dataSet = session.executeQueryStatement("select count(*) 
from root");
+    SessionDataSet dataSet =
+        session.executeQueryStatement("select count(*) from root align by 
device");
     while (dataSet.hasNext()) {
       RowRecord rowRecord = dataSet.next();
       Assert.assertEquals(15L, rowRecord.getFields().get(0).getLongV());
diff --git a/thrift-cluster/src/main/thrift/cluster.thrift 
b/thrift-cluster/src/main/thrift/cluster.thrift
index f23130e..3bfd56b 100644
--- a/thrift-cluster/src/main/thrift/cluster.thrift
+++ b/thrift-cluster/src/main/thrift/cluster.thrift
@@ -261,6 +261,8 @@ struct GetAllPathsResult {
   2: optional list<string> aliasList
 }
 
+struct EmptyReuqest {}
+
 
 service RaftService {
   /**
@@ -313,6 +315,8 @@ service RaftService {
   **/
   rpc.TSStatus executeNonQueryPlan(1:ExecutNonQueryReq request)
 
+  void RefreshClient(1:EmptyReuqest request)
+
   /**
   * Ask the leader for its commit index, used to check whether the node has 
caught up with the
   * leader.

Reply via email to