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);

Reply via email to