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.
