This is an automated email from the ASF dual-hosted git repository.
haonan 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 fbd0c6a [To rel/0.12][IOTDB-1488] Fix metaMember's forwarding
clientPool timeout in cluster module (#3538)
fbd0c6a is described below
commit fbd0c6a9364de566fec4b7c79e5cbddc165aa8a6
Author: Potato <[email protected]>
AuthorDate: Tue Jul 13 10:53:01 2021 +0800
[To rel/0.12][IOTDB-1488] Fix metaMember's forwarding clientPool timeout in
cluster module (#3538)
---
.../iotdb/cluster/client/DataClientProvider.java | 45 ++++++++++++--
.../cluster/client/async/AsyncClientPool.java | 22 +++++++
.../iotdb/cluster/client/sync/SyncClientPool.java | 21 +++++++
.../iotdb/cluster/server/DataClusterServer.java | 9 +++
.../iotdb/cluster/server/MetaClusterServer.java | 9 +++
.../cluster/server/member/MetaGroupMember.java | 72 ++++++++++++++++++++++
.../cluster/server/service/BaseAsyncService.java | 4 ++
.../cluster/server/service/BaseSyncService.java | 4 ++
.../org/apache/iotdb/db/metadata/MManager.java | 3 +
thrift-cluster/src/main/thrift/cluster.thrift | 6 +-
10 files changed, 190 insertions(+), 5 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 8b954ec..106705f 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
@@ -43,6 +43,8 @@ public class DataClientProvider {
private SyncClientPool dataSyncClientPool;
+ private static final String GET_CLIENT_FAILED_MSG = "can not get client for
node=";
+
public DataClientProvider(TProtocolFactory factory) {
if (!ClusterDescriptor.getInstance().getConfig().isUseAsyncServer()) {
dataSyncClientPool = new SyncClientPool(new
SyncDataClient.FactorySync(factory));
@@ -60,7 +62,7 @@ public class DataClientProvider {
}
/**
- * Get a thrift client that will connect to "node" using the data port.
+ * Get a thrift client from the head of deque that will connect to "node"
using the data port.
*
* @param node the node to be connected
* @param timeout timeout threshold of connection
@@ -68,7 +70,23 @@ public class DataClientProvider {
public AsyncDataClient getAsyncDataClient(Node node, int timeout) throws
IOException {
AsyncDataClient client = (AsyncDataClient)
getDataAsyncClientPool().getClient(node);
if (client == null) {
- throw new IOException("can not get client for node=" + node);
+ throw new IOException(GET_CLIENT_FAILED_MSG + node);
+ }
+ client.setTimeout(timeout);
+ return client;
+ }
+
+ /**
+ * Get a thrift client from the tail of deque that will connect to "node"
using the data port for
+ * refresh.
+ *
+ * @param node the node to be connected
+ * @param timeout timeout threshold of connection
+ */
+ public AsyncDataClient getAsyncDataClientForRefresh(Node node, int timeout)
throws IOException {
+ AsyncDataClient client = (AsyncDataClient)
getDataAsyncClientPool().getClientForRefresh(node);
+ if (client == null) {
+ throw new IOException(GET_CLIENT_FAILED_MSG + node);
}
client.setTimeout(timeout);
return client;
@@ -79,7 +97,7 @@ public class DataClientProvider {
* org.apache.iotdb.cluster.utils.ClientUtils#putBackSyncClient(Client)} to
put the client back
* into the client pool, otherwise there is a risk of client leakage.
*
- * <p>Get a thrift client that will connect to "node" using the data port.
+ * <p>Get a thrift client from the head of deque that will connect to "node"
using the data port.
*
* @param node the node to be connected
* @param timeout timeout threshold of connection
@@ -87,7 +105,26 @@ public class DataClientProvider {
public SyncDataClient getSyncDataClient(Node node, int timeout) throws
TException {
SyncDataClient client = (SyncDataClient)
getDataSyncClientPool().getClient(node);
if (client == null) {
- throw new TException("can not get client for node=" + node);
+ throw new TException(GET_CLIENT_FAILED_MSG + node);
+ }
+ client.setTimeout(timeout);
+ return client;
+ }
+
+ /**
+ * IMPORTANT!!! After calling this function, the caller should make sure to
call {@link
+ * org.apache.iotdb.cluster.utils.ClientUtils#putBackSyncClient(Client)} to
put the client back
+ * into the client pool, otherwise there is a risk of client leakage.
+ *
+ * <p>Get a thrift client from the tail of deque that will connect to "node"
using the data port.
+ *
+ * @param node the node to be connected
+ * @param timeout timeout threshold of connection
+ */
+ public SyncDataClient getSyncDataClientForRefresh(Node node, int timeout)
throws TException {
+ SyncDataClient client = (SyncDataClient)
getDataSyncClientPool().getClientForRefresh(node);
+ if (client == null) {
+ throw new TException(GET_CLIENT_FAILED_MSG + node);
}
client.setTimeout(timeout);
return client;
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/client/async/AsyncClientPool.java
b/cluster/src/main/java/org/apache/iotdb/cluster/client/async/AsyncClientPool.java
index 4ba3fe9..ddabff7 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/client/async/AsyncClientPool.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/client/async/AsyncClientPool.java
@@ -53,6 +53,28 @@ public class AsyncClientPool {
}
/**
+ * Get a client of the given node from the cache if one is available, or
null.
+ *
+ * <p>IMPORTANT!!! The caller should check whether the return value is null
or not!
+ *
+ * @param node the node want to connect
+ * @return if the node can connect, return the client, otherwise null
+ */
+ public AsyncClient getClientForRefresh(Node node) {
+ ClusterNode clusterNode = new ClusterNode(node);
+ // As clientCaches is ConcurrentHashMap, computeIfAbsent is thread safety.
+ Deque<AsyncClient> clientStack =
+ clientCaches.computeIfAbsent(clusterNode, n -> new ArrayDeque<>());
+ synchronized (clientStack) {
+ if (clientStack.isEmpty()) {
+ return null;
+ } else {
+ return clientStack.pollLast();
+ }
+ }
+ }
+
+ /**
* See getClient(Node node, boolean activatedOnly)
*
* @param node
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..67dbae5 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
@@ -52,6 +52,27 @@ public class SyncClientPool {
}
/**
+ * Get a client of the given node from the cache if one is available, or
null.
+ *
+ * <p>IMPORTANT!!! The caller should check whether the return value is null
or not!
+ *
+ * @param node the node want to connect
+ * @return if the node can connect, return the client, otherwise null
+ */
+ 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.pollLast();
+ }
+ }
+ }
+
+ /**
* See getClient(Node node, boolean activatedOnly)
*
* @param node the node want to connect
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..858459b 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
@@ -46,6 +46,7 @@ import org.apache.iotdb.cluster.rpc.thrift.PullSchemaRequest;
import org.apache.iotdb.cluster.rpc.thrift.PullSchemaResp;
import org.apache.iotdb.cluster.rpc.thrift.PullSnapshotRequest;
import org.apache.iotdb.cluster.rpc.thrift.PullSnapshotResp;
+import org.apache.iotdb.cluster.rpc.thrift.RefreshReuqest;
import org.apache.iotdb.cluster.rpc.thrift.RequestCommitIndexResponse;
import org.apache.iotdb.cluster.rpc.thrift.SendSnapshotRequest;
import org.apache.iotdb.cluster.rpc.thrift.SingleSeriesQueryRequest;
@@ -308,6 +309,11 @@ public class DataClusterServer extends RaftServer
}
@Override
+ public void refreshConnection(RefreshReuqest request,
AsyncMethodCallback<Void> resultHandler) {
+ resultHandler.onComplete(null);
+ }
+
+ @Override
public void requestCommitIndex(
Node header, AsyncMethodCallback<RequestCommitIndexResponse>
resultHandler) {
DataAsyncService service = getDataAsyncService(header, resultHandler,
"Request commit index");
@@ -921,6 +927,9 @@ public class DataClusterServer extends RaftServer
}
@Override
+ public void refreshConnection(RefreshReuqest 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..8efd7f5 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
@@ -33,6 +33,7 @@ import org.apache.iotdb.cluster.rpc.thrift.ExecutNonQueryReq;
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.RefreshReuqest;
import org.apache.iotdb.cluster.rpc.thrift.RequestCommitIndexResponse;
import org.apache.iotdb.cluster.rpc.thrift.SendSnapshotRequest;
import org.apache.iotdb.cluster.rpc.thrift.StartUpStatus;
@@ -226,6 +227,11 @@ public class MetaClusterServer extends RaftServer
}
@Override
+ public void refreshConnection(RefreshReuqest request,
AsyncMethodCallback<Void> resultHandler) {
+ resultHandler.onComplete(null);
+ }
+
+ @Override
public void requestCommitIndex(
Node header, AsyncMethodCallback<RequestCommitIndexResponse>
resultHandler) {
asyncService.requestCommitIndex(header, resultHandler);
@@ -334,6 +340,9 @@ public class MetaClusterServer extends RaftServer
}
@Override
+ public void refreshConnection(RefreshReuqest 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..cede9a3 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
@@ -56,6 +56,8 @@ import
org.apache.iotdb.cluster.rpc.thrift.CheckStatusResponse;
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.RefreshReuqest;
import org.apache.iotdb.cluster.rpc.thrift.SendSnapshotRequest;
import org.apache.iotdb.cluster.rpc.thrift.StartUpStatus;
import org.apache.iotdb.cluster.rpc.thrift.TSMetaService;
@@ -163,6 +165,12 @@ public class MetaGroupMember extends RaftMember {
*/
private static final int REPORT_INTERVAL_SEC = 10;
+ /**
+ * 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;
+
/** 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 +219,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 +341,8 @@ public class MetaGroupMember extends RaftMember {
Executors.newSingleThreadScheduledExecutor(n -> new Thread(n,
"NodeReportThread"));
hardLinkCleanerThread =
Executors.newSingleThreadScheduledExecutor(n -> new Thread(n,
"HardLinkCleaner"));
+ dataClientRefresher =
+ Executors.newSingleThreadScheduledExecutor(n -> new Thread(n,
"DataClientRefresher"));
}
/**
@@ -349,6 +361,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 +481,57 @@ public class MetaGroupMember extends RaftMember {
CLEAN_HARDLINK_INTERVAL_SEC,
CLEAN_HARDLINK_INTERVAL_SEC,
TimeUnit.SECONDS);
+ dataClientRefresher.scheduleAtFixedRate(
+ this::refreshClientOnce, REFRESH_CLIENT_SEC, REFRESH_CLIENT_SEC,
TimeUnit.SECONDS);
+ }
+
+ private void refreshClientOnce() {
+ for (Node receiver : allNodes) {
+ if (!receiver.equals(thisNode)) {
+ if (ClusterDescriptor.getInstance().getConfig().isUseAsyncServer()) {
+ refreshClientOnceAsync(receiver);
+ } else {
+ refreshClientOnceSync(receiver);
+ }
+ }
+ }
+ }
+
+ private void refreshClientOnceSync(Node receiver) {
+ RaftService.Client client;
+ try {
+ client =
+ getClientProvider()
+ .getSyncDataClientForRefresh(receiver,
RaftServer.getWriteOperationTimeoutMS());
+ } catch (TException e) {
+ return;
+ }
+ try {
+ RefreshReuqest req = new RefreshReuqest();
+ client.refreshConnection(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 refreshClientOnceAsync(Node receiver) {
+ RaftService.AsyncClient client;
+ try {
+ client =
+ getClientProvider()
+ .getAsyncDataClientForRefresh(receiver,
RaftServer.getWriteOperationTimeoutMS());
+ } catch (IOException e) {
+ return;
+ }
+ try {
+ client.refreshConnection(new RefreshReuqest(), new
GenericHandler<>(receiver, null));
+ } catch (TException e) {
+ logger.warn("encounter refreshing client timeout, throw broken
connection", e);
+ }
}
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..aa6c9c7 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
@@ -30,6 +30,7 @@ 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.RaftService.AsyncClient;
+import org.apache.iotdb.cluster.rpc.thrift.RefreshReuqest;
import org.apache.iotdb.cluster.rpc.thrift.RequestCommitIndexResponse;
import org.apache.iotdb.cluster.server.NodeCharacter;
import org.apache.iotdb.cluster.server.member.RaftMember;
@@ -145,6 +146,9 @@ public abstract class BaseAsyncService implements
RaftService.AsyncIface {
}
@Override
+ public void refreshConnection(RefreshReuqest 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..c9a3016 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
@@ -30,6 +30,7 @@ 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.RaftService.Client;
+import org.apache.iotdb.cluster.rpc.thrift.RefreshReuqest;
import org.apache.iotdb.cluster.rpc.thrift.RequestCommitIndexResponse;
import org.apache.iotdb.cluster.server.NodeCharacter;
import org.apache.iotdb.cluster.server.member.RaftMember;
@@ -152,6 +153,9 @@ public abstract class BaseSyncService implements
RaftService.Iface {
}
@Override
+ public void refreshConnection(RefreshReuqest 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/thrift-cluster/src/main/thrift/cluster.thrift
b/thrift-cluster/src/main/thrift/cluster.thrift
index f23130e..e9332a5 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 RefreshReuqest {}
+
service RaftService {
/**
@@ -336,7 +338,9 @@ service RaftService {
* When a follower finds that it already has a file in a snapshot locally, it
calls this
* interface to notify the leader to remove the associated hardlink.
**/
- void removeHardLink(1: string hardLinkPath)
+ void removeHardLink(1:string hardLinkPath)
+
+ void refreshConnection(1:RefreshReuqest request)
}