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(