This is an automated email from the ASF dual-hosted git repository.
chaow pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 1fd2e19 [IOTDB-1248] forward the pull schema request to the leader
(#2902)
1fd2e19 is described below
commit 1fd2e19392c6e9cb3e4bccc29a581ab766e5eaa5
Author: wangchao316 <[email protected]>
AuthorDate: Fri Mar 26 16:34:27 2021 +0800
[IOTDB-1248] forward the pull schema request to the leader (#2902)
---
.../cluster/server/service/DataAsyncService.java | 81 ++++++++++--------
.../cluster/server/service/DataSyncService.java | 96 +++++++++++-----------
.../cluster/server/member/DataGroupMemberTest.java | 64 ++++++++++++++-
3 files changed, 160 insertions(+), 81 deletions(-)
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataAsyncService.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataAsyncService.java
index 49c5a1e..a592599 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataAsyncService.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataAsyncService.java
@@ -38,6 +38,7 @@ import org.apache.iotdb.cluster.rpc.thrift.PullSnapshotResp;
import org.apache.iotdb.cluster.rpc.thrift.SendSnapshotRequest;
import org.apache.iotdb.cluster.rpc.thrift.SingleSeriesQueryRequest;
import org.apache.iotdb.cluster.rpc.thrift.TSDataService;
+import org.apache.iotdb.cluster.server.NodeCharacter;
import org.apache.iotdb.cluster.server.member.DataGroupMember;
import org.apache.iotdb.db.exception.StorageEngineException;
import org.apache.iotdb.db.exception.metadata.IllegalPathException;
@@ -113,27 +114,35 @@ public class DataAsyncService extends BaseAsyncService
implements TSDataService.
}
}
+ /**
+ * forward the request to the leader
+ *
+ * @param request pull schema request
+ * @param resultHandler result handler
+ */
@Override
public void pullTimeSeriesSchema(
PullSchemaRequest request, AsyncMethodCallback<PullSchemaResp>
resultHandler) {
- try {
- resultHandler.onComplete(
-
dataGroupMember.getLocalQueryExecutor().queryTimeSeriesSchema(request));
- } catch (CheckConsistencyException e) {
- // if this node cannot synchronize with the leader with in a given time,
forward the
- // request to the leader
- AsyncDataClient leaderClient = getLeaderClient();
- if (leaderClient == null) {
- resultHandler.onError(new
LeaderUnknownException(dataGroupMember.getAllNodes()));
- return;
- }
+ if (dataGroupMember.getCharacter() == NodeCharacter.LEADER) {
try {
- leaderClient.pullTimeSeriesSchema(request, resultHandler);
- } catch (TException e1) {
- resultHandler.onError(e1);
+ resultHandler.onComplete(
+
dataGroupMember.getLocalQueryExecutor().queryTimeSeriesSchema(request));
+ return;
+ } catch (CheckConsistencyException | MetadataException e) {
+ resultHandler.onError(e);
}
- } catch (MetadataException e) {
- resultHandler.onError(e);
+ }
+
+ // forward the request to the leader
+ AsyncDataClient leaderClient = getLeaderClient();
+ if (leaderClient == null) {
+ resultHandler.onError(new
LeaderUnknownException(dataGroupMember.getAllNodes()));
+ return;
+ }
+ try {
+ leaderClient.pullTimeSeriesSchema(request, resultHandler);
+ } catch (TException e1) {
+ resultHandler.onError(e1);
}
}
@@ -142,27 +151,35 @@ public class DataAsyncService extends BaseAsyncService
implements TSDataService.
return (AsyncDataClient)
dataGroupMember.getAsyncClient(dataGroupMember.getLeader());
}
+ /**
+ * forward the request to the leader
+ *
+ * @param request pull schema request
+ * @param resultHandler result handler
+ */
@Override
public void pullMeasurementSchema(
PullSchemaRequest request, AsyncMethodCallback<PullSchemaResp>
resultHandler) {
- try {
- resultHandler.onComplete(
-
dataGroupMember.getLocalQueryExecutor().queryMeasurementSchema(request));
- } catch (CheckConsistencyException e) {
- // if this node cannot synchronize with the leader with in a given time,
forward the
- // request to the leader
- AsyncDataClient leaderClient = getLeaderClient();
- if (leaderClient == null) {
- resultHandler.onError(new
LeaderUnknownException(dataGroupMember.getAllNodes()));
- return;
- }
+ if (dataGroupMember.getCharacter() == NodeCharacter.LEADER) {
try {
- leaderClient.pullMeasurementSchema(request, resultHandler);
- } catch (TException e1) {
- resultHandler.onError(e1);
+ resultHandler.onComplete(
+
dataGroupMember.getLocalQueryExecutor().queryMeasurementSchema(request));
+ return;
+ } catch (CheckConsistencyException | MetadataException e) {
+ resultHandler.onError(e);
}
- } catch (IllegalPathException e) {
- resultHandler.onError(e);
+ }
+
+ // forward the request to the leader
+ AsyncDataClient leaderClient = getLeaderClient();
+ if (leaderClient == null) {
+ resultHandler.onError(new
LeaderUnknownException(dataGroupMember.getAllNodes()));
+ return;
+ }
+ try {
+ leaderClient.pullMeasurementSchema(request, resultHandler);
+ } catch (TException e1) {
+ resultHandler.onError(e1);
}
}
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataSyncService.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataSyncService.java
index 94ac92f..062dcc7 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataSyncService.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataSyncService.java
@@ -38,6 +38,7 @@ import org.apache.iotdb.cluster.rpc.thrift.PullSnapshotResp;
import org.apache.iotdb.cluster.rpc.thrift.SendSnapshotRequest;
import org.apache.iotdb.cluster.rpc.thrift.SingleSeriesQueryRequest;
import org.apache.iotdb.cluster.rpc.thrift.TSDataService;
+import org.apache.iotdb.cluster.server.NodeCharacter;
import org.apache.iotdb.cluster.server.member.DataGroupMember;
import org.apache.iotdb.cluster.utils.ClientUtils;
import org.apache.iotdb.db.exception.StorageEngineException;
@@ -120,73 +121,76 @@ public class DataSyncService extends BaseSyncService
implements TSDataService.If
}
/**
- * return the schema, whose measurement Id is the series full path.
+ * forward the request to the leader return the schema, whose measurement Id
is the series full
+ * path.
*
* @param request the pull request
- * @return response
+ * @return response pull schema resp
* @throws TException remind of thrift
*/
@Override
public PullSchemaResp pullTimeSeriesSchema(PullSchemaRequest request) throws
TException {
- try {
- return
dataGroupMember.getLocalQueryExecutor().queryTimeSeriesSchema(request);
- } catch (CheckConsistencyException e) {
- // if this node cannot synchronize with the leader with in a given time,
forward the
- // request to the leader
- dataGroupMember.waitLeader();
- SyncDataClient client =
- (SyncDataClient)
dataGroupMember.getSyncClient(dataGroupMember.getLeader());
- if (client == null) {
- throw new TException(new
LeaderUnknownException(dataGroupMember.getAllNodes()));
- }
- PullSchemaResp pullSchemaResp;
+ if (dataGroupMember.getCharacter() == NodeCharacter.LEADER) {
try {
- pullSchemaResp = client.pullTimeSeriesSchema(request);
- } catch (TException te) {
- client.getInputProtocol().getTransport().close();
- throw te;
- } finally {
- ClientUtils.putBackSyncClient(client);
+ return
dataGroupMember.getLocalQueryExecutor().queryTimeSeriesSchema(request);
+ } catch (CheckConsistencyException | MetadataException e) {
+ throw new TException(e);
}
- return pullSchemaResp;
- } catch (MetadataException e) {
- throw new TException(e);
}
+
+ // forward the request to the leader
+ dataGroupMember.waitLeader();
+ SyncDataClient client =
+ (SyncDataClient)
dataGroupMember.getSyncClient(dataGroupMember.getLeader());
+ if (client == null) {
+ throw new TException(new
LeaderUnknownException(dataGroupMember.getAllNodes()));
+ }
+ PullSchemaResp pullSchemaResp;
+ try {
+ pullSchemaResp = client.pullTimeSeriesSchema(request);
+ } catch (TException te) {
+ client.getInputProtocol().getTransport().close();
+ throw te;
+ } finally {
+ ClientUtils.putBackSyncClient(client);
+ }
+ return pullSchemaResp;
}
/**
- * return the schema, whose measurement Id is the series name.
+ * forward the request to the leader return the schema, whose measurement Id
is the series name.
*
* @param request the pull request
- * @return response
+ * @return response pull schema resp
* @throws TException remind of thrift
*/
@Override
public PullSchemaResp pullMeasurementSchema(PullSchemaRequest request)
throws TException {
- try {
- return
dataGroupMember.getLocalQueryExecutor().queryMeasurementSchema(request);
- } catch (CheckConsistencyException e) {
- // if this node cannot synchronize with the leader with in a given time,
forward the
- // request to the leader
- dataGroupMember.waitLeader();
- SyncDataClient client =
- (SyncDataClient)
dataGroupMember.getSyncClient(dataGroupMember.getLeader());
- if (client == null) {
- throw new TException(new
LeaderUnknownException(dataGroupMember.getAllNodes()));
- }
- PullSchemaResp pullSchemaResp;
+ if (dataGroupMember.getCharacter() == NodeCharacter.LEADER) {
try {
- pullSchemaResp = client.pullMeasurementSchema(request);
- } catch (TException te) {
- client.getInputProtocol().getTransport().close();
- throw te;
- } finally {
- ClientUtils.putBackSyncClient(client);
+ return
dataGroupMember.getLocalQueryExecutor().queryMeasurementSchema(request);
+ } catch (CheckConsistencyException | IllegalPathException e) {
+ throw new TException(e);
}
- return pullSchemaResp;
- } catch (IllegalPathException e) {
- throw new TException(e);
}
+
+ // forward the request to the leader
+ dataGroupMember.waitLeader();
+ SyncDataClient client =
+ (SyncDataClient)
dataGroupMember.getSyncClient(dataGroupMember.getLeader());
+ if (client == null) {
+ throw new TException(new
LeaderUnknownException(dataGroupMember.getAllNodes()));
+ }
+ PullSchemaResp pullSchemaResp;
+ try {
+ pullSchemaResp = client.pullMeasurementSchema(request);
+ } catch (TException te) {
+ client.getInputProtocol().getTransport().close();
+ throw te;
+ } finally {
+ ClientUtils.putBackSyncClient(client);
+ }
+ return pullSchemaResp;
}
@Override
diff --git
a/cluster/src/test/java/org/apache/iotdb/cluster/server/member/DataGroupMemberTest.java
b/cluster/src/test/java/org/apache/iotdb/cluster/server/member/DataGroupMemberTest.java
index a53a32e..77e72b2 100644
---
a/cluster/src/test/java/org/apache/iotdb/cluster/server/member/DataGroupMemberTest.java
+++
b/cluster/src/test/java/org/apache/iotdb/cluster/server/member/DataGroupMemberTest.java
@@ -44,6 +44,7 @@ import org.apache.iotdb.cluster.rpc.thrift.GetAllPathsResult;
import org.apache.iotdb.cluster.rpc.thrift.GroupByRequest;
import org.apache.iotdb.cluster.rpc.thrift.Node;
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.RaftService.AsyncClient;
@@ -55,6 +56,7 @@ import org.apache.iotdb.cluster.server.Response;
import org.apache.iotdb.cluster.server.handlers.caller.GenericHandler;
import
org.apache.iotdb.cluster.server.handlers.caller.PullMeasurementSchemaHandler;
import org.apache.iotdb.cluster.server.handlers.caller.PullSnapshotHandler;
+import
org.apache.iotdb.cluster.server.handlers.caller.PullTimeseriesSchemaHandler;
import org.apache.iotdb.cluster.server.service.DataAsyncService;
import org.apache.iotdb.cluster.utils.Constants;
import org.apache.iotdb.db.engine.StorageEngine;
@@ -206,6 +208,22 @@ public class DataGroupMemberTest extends BaseMember {
return new TestAsyncDataClient(node, dataGroupMemberMap) {
@Override
+ public void pullMeasurementSchema(
+ PullSchemaRequest request,
AsyncMethodCallback<PullSchemaResp> resultHandler) {
+
dataGroupMemberMap.get(request.getHeader()).setCharacter(NodeCharacter.LEADER);
+ new
DataAsyncService(dataGroupMemberMap.get(request.getHeader()))
+ .pullMeasurementSchema(request, resultHandler);
+ }
+
+ @Override
+ public void pullTimeSeriesSchema(
+ PullSchemaRequest request,
AsyncMethodCallback<PullSchemaResp> resultHandler) {
+
dataGroupMemberMap.get(request.getHeader()).setCharacter(NodeCharacter.LEADER);
+ new
DataAsyncService(dataGroupMemberMap.get(request.getHeader()))
+ .pullTimeSeriesSchema(request, resultHandler);
+ }
+
+ @Override
public void pullSnapshot(
PullSnapshotRequest request,
AsyncMethodCallback<PullSnapshotResp> resultHandler) {
@@ -598,20 +616,60 @@ public class DataGroupMemberTest extends BaseMember {
}
@Test
- public void testPullTimeseries() {
- System.out.println("Start testPullTimeseries()");
+ public void testPullTimeseriesSchema() {
+ System.out.println("Start testPullTimeseriesSchema()");
+ int prevTimeOut = RaftServer.getConnectionTimeoutInMS();
+ int prevMaxWait = RaftServer.getSyncLeaderMaxWaitMs();
+ RaftServer.setConnectionTimeoutInMS(20);
+ RaftServer.setSyncLeaderMaxWaitMs(200);
+ try {
+ // sync with leader is temporarily disabled, the request should be
forward to the leader
+ dataGroupMember.setLeader(TestUtils.getNode(0));
+ dataGroupMember.setCharacter(NodeCharacter.FOLLOWER);
+ enableSyncLeader = false;
+
+ PullSchemaRequest request = new PullSchemaRequest();
+
request.setPrefixPaths(Collections.singletonList(TestUtils.getTestSg(0)));
+ request.setHeader(TestUtils.getNode(0));
+ AtomicReference<List<TimeseriesSchema>> result = new AtomicReference<>();
+ PullTimeseriesSchemaHandler handler =
+ new PullTimeseriesSchemaHandler(TestUtils.getNode(1),
request.getPrefixPaths(), result);
+ new DataAsyncService(dataGroupMember).pullTimeSeriesSchema(request,
handler);
+ for (int i = 0; i < 10; i++) {
+ assertTrue(result.get().contains(TestUtils.getTestTimeSeriesSchema(0,
i)));
+ }
+
+ // the member is a leader itself
+ dataGroupMember.setCharacter(NodeCharacter.LEADER);
+ result.set(null);
+ handler =
+ new PullTimeseriesSchemaHandler(TestUtils.getNode(1),
request.getPrefixPaths(), result);
+ new DataAsyncService(dataGroupMember).pullTimeSeriesSchema(request,
handler);
+ for (int i = 0; i < 10; i++) {
+ assertTrue(result.get().contains(TestUtils.getTestTimeSeriesSchema(0,
i)));
+ }
+ } finally {
+ RaftServer.setConnectionTimeoutInMS(prevTimeOut);
+ RaftServer.setSyncLeaderMaxWaitMs(prevMaxWait);
+ }
+ }
+
+ @Test
+ public void testPullMeasurementSchema() {
+ System.out.println("Start testPullMeasurementSchema()");
int prevTimeOut = RaftServer.getConnectionTimeoutInMS();
int prevMaxWait = RaftServer.getSyncLeaderMaxWaitMs();
RaftServer.setConnectionTimeoutInMS(20);
RaftServer.setSyncLeaderMaxWaitMs(200);
try {
// sync with leader is temporarily disabled, the request should be
forward to the leader
- dataGroupMember.setLeader(TestUtils.getNode(1));
+ dataGroupMember.setLeader(TestUtils.getNode(0));
dataGroupMember.setCharacter(NodeCharacter.FOLLOWER);
enableSyncLeader = false;
PullSchemaRequest request = new PullSchemaRequest();
request.setPrefixPaths(Collections.singletonList(TestUtils.getTestSg(0)));
+ request.setHeader(TestUtils.getNode(0));
AtomicReference<List<MeasurementSchema>> result = new
AtomicReference<>();
PullMeasurementSchemaHandler handler =
new PullMeasurementSchemaHandler(TestUtils.getNode(1),
request.getPrefixPaths(), result);