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 430cfcd  [To rel/0.12]Fix thrift out of sequence in cluster module 
(#3580)
430cfcd is described below

commit 430cfcd32af615ce887ffd2d03543c6ab8c6a478
Author: Potato <[email protected]>
AuthorDate: Wed Jul 28 15:12:33 2021 +0800

    [To rel/0.12]Fix thrift out of sequence in cluster module (#3580)
---
 .../query/last/ClusterLastQueryExecutor.java        | 21 ++++++---------------
 .../cluster/server/member/DataGroupMember.java      |  8 +-------
 .../cluster/server/member/MetaGroupMember.java      |  6 +++---
 .../iotdb/cluster/server/member/RaftMember.java     |  8 +-------
 .../db/engine/storagegroup/TsFileProcessor.java     |  5 -----
 5 files changed, 11 insertions(+), 37 deletions(-)

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 a30db6c..0ca6a62 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,7 +30,6 @@ 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;
@@ -245,6 +244,7 @@ public class ClusterLastQueryExecutor extends 
LastQueryExecutor {
                 .getClientProvider()
                 .getAsyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS());
       } catch (IOException e) {
+        logger.warn("can not get client for node= {}", node);
         return null;
       }
       buffer =
@@ -259,12 +259,10 @@ public class ClusterLastQueryExecutor extends 
LastQueryExecutor {
     }
 
     private ByteBuffer lastSync(Node node, QueryContext context) throws 
TException {
-      SyncDataClient client = null;
-      try {
-        client =
-            metaGroupMember
-                .getClientProvider()
-                .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS());
+      try (SyncDataClient client =
+          metaGroupMember
+              .getClientProvider()
+              .getSyncDataClient(node, 
RaftServer.getReadOperationTimeoutMS())) {
         return client.last(
             new LastQueryRequest(
                 PartialPath.toStringList(seriesPaths),
@@ -274,15 +272,8 @@ public class ClusterLastQueryExecutor extends 
LastQueryExecutor {
                 group.getHeader(),
                 client.getNode()));
       } catch (IOException e) {
+        logger.warn("can not get client for node= {}", node);
         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/server/member/DataGroupMember.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
index 89d933a..6f6936a 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,7 +696,6 @@ 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);
@@ -719,11 +718,6 @@ 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);
@@ -731,7 +725,7 @@ public class DataGroupMember extends RaftMember {
         return status;
       }
 
-      startTime = 
Timer.Statistic.DATA_GROUP_MEMBER_WAIT_LEADER.getOperationStartTime();
+      long 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 1cc4ad9..768afee 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 = 1;
+  private static final int REFRESH_CLIENT_SEC = 5;
 
   /** how many times is a data record replicated, also the number of nodes in 
a data group */
   private static final int REPLICATION_NUM =
@@ -507,7 +507,7 @@ public class MetaGroupMember extends RaftMember {
       client.refreshConnection(req);
     } catch (IOException ignored) {
     } catch (TException e) {
-      logger.info("encounter refreshing client timeout, throw broken 
connection", 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 {
@@ -529,7 +529,7 @@ public class MetaGroupMember extends RaftMember {
     try {
       client.refreshConnection(new RefreshReuqest(), new 
GenericHandler<>(receiver, null));
     } catch (TException e) {
-      logger.info("encounter refreshing client timeout, throw broken 
connection", e);
+      logger.warn("encounter refreshing client timeout, throw broken 
connection", e);
     }
   }
 
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 3078013..99aeb15 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,7 +1317,6 @@ 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));
@@ -1337,12 +1336,7 @@ public abstract class RaftMember {
       TSStatus status;
       if (e.getCause() instanceof SocketTimeoutException) {
         status = StatusUtils.TIME_OUT;
-        logger.warn(
-            MSG_FORWARD_TIMEOUT + ": {}ms",
-            name,
-            plan,
-            receiver,
-            System.currentTimeMillis() - startTime);
+        logger.warn(MSG_FORWARD_TIMEOUT, name, plan, receiver);
       } else {
         logger.error(MSG_FORWARD_ERROR, name, plan, receiver, e);
         status = StatusUtils.getStatus(StatusUtils.INTERNAL_ERROR, 
e.getMessage());
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 94c63f0..0eca8a2 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,12 +192,7 @@ 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(

Reply via email to