This is an automated email from the ASF dual-hosted git repository.

chaow pushed a commit to branch optimize_sync_meta
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit de6d61547caaad392163bb688a293452dba0e722
Author: chaow <[email protected]>
AuthorDate: Fri Apr 9 12:10:59 2021 +0800

    optimize sync leader for meta
---
 .../iotdb/cluster/server/DataClusterServer.java    |  6 ++-
 .../iotdb/cluster/server/MetaClusterServer.java    |  6 ++-
 .../iotdb/cluster/server/member/RaftMember.java    | 47 ++++++++++++++++++----
 .../cluster/server/service/BaseAsyncService.java   | 19 +++++++--
 .../cluster/server/service/BaseSyncService.java    | 23 ++++++++---
 .../cluster/server/member/DataGroupMemberTest.java |  5 ++-
 .../cluster/server/member/RaftMemberTest.java      |  9 +++--
 thrift-cluster/src/main/thrift/cluster.thrift      |  8 +++-
 8 files changed, 96 insertions(+), 27 deletions(-)

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 a54de23..e4c81f8 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.RequestCommitIndexResponse;
 import org.apache.iotdb.cluster.rpc.thrift.SendSnapshotRequest;
 import org.apache.iotdb.cluster.rpc.thrift.SingleSeriesQueryRequest;
 import org.apache.iotdb.cluster.rpc.thrift.TSDataService;
@@ -307,7 +308,8 @@ public class DataClusterServer extends RaftServer
   }
 
   @Override
-  public void requestCommitIndex(Node header, AsyncMethodCallback<Long> 
resultHandler) {
+  public void requestCommitIndex(
+      Node header, AsyncMethodCallback<RequestCommitIndexResponse> 
resultHandler) {
     DataAsyncService service = getDataAsyncService(header, resultHandler, 
"Request commit index");
     if (service != null) {
       service.requestCommitIndex(header, resultHandler);
@@ -919,7 +921,7 @@ public class DataClusterServer extends RaftServer
   }
 
   @Override
-  public long requestCommitIndex(Node header) throws TException {
+  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 02d53b3..12e286f 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.RequestCommitIndexResponse;
 import org.apache.iotdb.cluster.rpc.thrift.SendSnapshotRequest;
 import org.apache.iotdb.cluster.rpc.thrift.StartUpStatus;
 import org.apache.iotdb.cluster.rpc.thrift.TNodeStatus;
@@ -224,7 +225,8 @@ public class MetaClusterServer extends RaftServer
   }
 
   @Override
-  public void requestCommitIndex(Node header, AsyncMethodCallback<Long> 
resultHandler) {
+  public void requestCommitIndex(
+      Node header, AsyncMethodCallback<RequestCommitIndexResponse> 
resultHandler) {
     asyncService.requestCommitIndex(header, resultHandler);
   }
 
@@ -331,7 +333,7 @@ public class MetaClusterServer extends RaftServer
   }
 
   @Override
-  public long requestCommitIndex(Node header) throws TException {
+  public RequestCommitIndexResponse requestCommitIndex(Node header) throws 
TException {
     return syncService.requestCommitIndex(header);
   }
 
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 571424a..5bcd65c 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
@@ -47,6 +47,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.AsyncClient;
 import org.apache.iotdb.cluster.rpc.thrift.RaftService.Client;
+import org.apache.iotdb.cluster.rpc.thrift.RequestCommitIndexResponse;
 import org.apache.iotdb.cluster.server.NodeCharacter;
 import org.apache.iotdb.cluster.server.RaftServer;
 import org.apache.iotdb.cluster.server.Response;
@@ -872,8 +873,34 @@ public abstract class RaftMember {
   protected boolean waitUntilCatchUp(CheckConsistency checkConsistency)
       throws CheckConsistencyException {
     long leaderCommitId = Long.MIN_VALUE;
+    RequestCommitIndexResponse response;
     try {
-      leaderCommitId = config.isUseAsyncServer() ? requestCommitIdAsync() : 
requestCommitIdSync();
+      response = config.isUseAsyncServer() ? requestCommitIdAsync() : 
requestCommitIdSync();
+      leaderCommitId = response.getCommitLogIndex();
+      if (term.get() == response.getTerm() && logManager.getCommitLogIndex() < 
leaderCommitId) {
+        // there are more local logs that can be committed, commit them in a 
ThreadPool so the
+        // heartbeat response will not be blocked
+        CommitLogTask commitLogTask =
+            new CommitLogTask(logManager, leaderCommitId, 
response.getCommitLogTerm());
+        commitLogTask.registerCallback(new CommitLogCallback(this));
+        // if the log is not consistent, the commitment will be blocked until 
the leader makes the
+        // node catch up
+        if (commitLogPool != null && !commitLogPool.isShutdown()) {
+          commitLogPool.submit(commitLogTask);
+        }
+
+        logger.debug(
+            "{}: Inconsistent log found, leaderCommit: {}-{}, localCommit: 
{}-{}, "
+                + "localLast: {}-{}",
+            name,
+            response.getCommitLogIndex(),
+            response.getCommitLogTerm(),
+            logManager.getCommitLogIndex(),
+            logManager.getCommitLogTerm(),
+            logManager.getLastLogIndex(),
+            logManager.getLastLogTerm());
+      }
+
       return syncLocalApply(leaderCommitId);
     } catch (TException e) {
       logger.error(MSG_NO_LEADER_COMMIT_INDEX, name, leader.get(), e);
@@ -1057,9 +1084,12 @@ public abstract class RaftMember {
   }
 
   @SuppressWarnings("java:S2274") // enable timeout
-  protected long requestCommitIdAsync() throws TException, 
InterruptedException {
+  protected RequestCommitIndexResponse requestCommitIdAsync()
+      throws TException, InterruptedException {
     // use Long.MAX_VALUE to indicate a timeout
-    AtomicReference<Long> commitIdResult = new 
AtomicReference<>(Long.MAX_VALUE);
+    RequestCommitIndexResponse response =
+        new RequestCommitIndexResponse(Long.MAX_VALUE, Long.MAX_VALUE, 
Long.MAX_VALUE);
+    AtomicReference<RequestCommitIndexResponse> commitIdResult = new 
AtomicReference<>(response);
     AsyncClient client = getAsyncClient(leader.get());
     if (client == null) {
       // cannot connect to the leader
@@ -1073,24 +1103,25 @@ public abstract class RaftMember {
     return commitIdResult.get();
   }
 
-  private long requestCommitIdSync() throws TException {
+  private RequestCommitIndexResponse requestCommitIdSync() throws TException {
     Client client = getSyncClient(leader.get());
+    RequestCommitIndexResponse response;
     if (client == null) {
       // cannot connect to the leader
       logger.warn(MSG_NO_LEADER_IN_SYNC, name);
       // use Long.MAX_VALUE to indicate a timeouts
-      return Long.MAX_VALUE;
+      response = new RequestCommitIndexResponse(Long.MAX_VALUE, 
Long.MAX_VALUE, Long.MAX_VALUE);
+      return response;
     }
-    long commitIndex;
     try {
-      commitIndex = client.requestCommitIndex(getHeader());
+      response = client.requestCommitIndex(getHeader());
     } catch (TException e) {
       client.getInputProtocol().getTransport().close();
       throw e;
     } finally {
       ClientUtils.putBackSyncClient(client);
     }
-    return commitIndex;
+    return response;
   }
 
   /**
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 07dbdea..8673078 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.RequestCommitIndexResponse;
 import org.apache.iotdb.cluster.server.NodeCharacter;
 import org.apache.iotdb.cluster.server.member.RaftMember;
 import org.apache.iotdb.cluster.utils.IOUtils;
@@ -85,10 +86,22 @@ public abstract class BaseAsyncService implements 
RaftService.AsyncIface {
   }
 
   @Override
-  public void requestCommitIndex(Node header, AsyncMethodCallback<Long> 
resultHandler) {
-    long commitIndex = member.getCommitIndex();
+  public void requestCommitIndex(
+      Node header, AsyncMethodCallback<RequestCommitIndexResponse> 
resultHandler) {
+    long commitIndex;
+    long commitTerm;
+    long curTerm;
+    synchronized (member.getTerm()) {
+      commitIndex = member.getLogManager().getCommitLogIndex();
+      commitTerm = member.getLogManager().getCommitLogTerm();
+      curTerm = member.getTerm().get();
+    }
+
+    RequestCommitIndexResponse response =
+        new RequestCommitIndexResponse(curTerm, commitIndex, commitTerm);
+
     if (commitIndex != Long.MIN_VALUE) {
-      resultHandler.onComplete(commitIndex);
+      resultHandler.onComplete(response);
       return;
     }
 
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 697f54e..ce200ab 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.RequestCommitIndexResponse;
 import org.apache.iotdb.cluster.server.NodeCharacter;
 import org.apache.iotdb.cluster.server.member.RaftMember;
 import org.apache.iotdb.cluster.utils.ClientUtils;
@@ -93,10 +94,22 @@ public abstract class BaseSyncService implements 
RaftService.Iface {
   }
 
   @Override
-  public long requestCommitIndex(Node header) throws TException {
-    long commitIndex = member.getCommitIndex();
+  public RequestCommitIndexResponse requestCommitIndex(Node header) throws 
TException {
+
+    long commitIndex;
+    long commitTerm;
+    long curTerm;
+    synchronized (member.getTerm()) {
+      commitIndex = member.getLogManager().getCommitLogIndex();
+      commitTerm = member.getLogManager().getCommitLogTerm();
+      curTerm = member.getTerm().get();
+    }
+
+    RequestCommitIndexResponse response =
+        new RequestCommitIndexResponse(curTerm, commitIndex, commitTerm);
+
     if (commitIndex != Long.MIN_VALUE) {
-      return commitIndex;
+      return response;
     }
 
     member.waitLeader();
@@ -105,14 +118,14 @@ public abstract class BaseSyncService implements 
RaftService.Iface {
       throw new TException(new LeaderUnknownException(member.getAllNodes()));
     }
     try {
-      commitIndex = client.requestCommitIndex(header);
+      response = client.requestCommitIndex(header);
     } catch (TException e) {
       client.getInputProtocol().getTransport().close();
       throw e;
     } finally {
       ClientUtils.putBackSyncClient(client);
     }
-    return commitIndex;
+    return response;
   }
 
   @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 77e72b2..f6d5e5d 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
@@ -48,6 +48,7 @@ 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;
+import org.apache.iotdb.cluster.rpc.thrift.RequestCommitIndexResponse;
 import org.apache.iotdb.cluster.rpc.thrift.SendSnapshotRequest;
 import org.apache.iotdb.cluster.rpc.thrift.SingleSeriesQueryRequest;
 import org.apache.iotdb.cluster.server.NodeCharacter;
@@ -244,11 +245,11 @@ public class DataGroupMemberTest extends BaseMember {
 
                 @Override
                 public void requestCommitIndex(
-                    Node header, AsyncMethodCallback<Long> resultHandler) {
+                    Node header, 
AsyncMethodCallback<RequestCommitIndexResponse> resultHandler) {
                   new Thread(
                           () -> {
                             if (enableSyncLeader) {
-                              resultHandler.onComplete(-1L);
+                              resultHandler.onComplete(new 
RequestCommitIndexResponse());
                             } else {
                               resultHandler.onError(new TestException());
                             }
diff --git 
a/cluster/src/test/java/org/apache/iotdb/cluster/server/member/RaftMemberTest.java
 
b/cluster/src/test/java/org/apache/iotdb/cluster/server/member/RaftMemberTest.java
index 494f845..694c4aa 100644
--- 
a/cluster/src/test/java/org/apache/iotdb/cluster/server/member/RaftMemberTest.java
+++ 
b/cluster/src/test/java/org/apache/iotdb/cluster/server/member/RaftMemberTest.java
@@ -29,6 +29,7 @@ import 
org.apache.iotdb.cluster.log.manage.PartitionedSnapshotLogManager;
 import org.apache.iotdb.cluster.rpc.thrift.AppendEntryRequest;
 import org.apache.iotdb.cluster.rpc.thrift.Node;
 import org.apache.iotdb.cluster.rpc.thrift.RaftService;
+import org.apache.iotdb.cluster.rpc.thrift.RequestCommitIndexResponse;
 import org.apache.iotdb.cluster.server.NodeCharacter;
 import org.apache.iotdb.cluster.server.Response;
 
@@ -178,8 +179,8 @@ public class RaftMemberTest extends BaseMember {
           }
 
           @Override
-          protected long requestCommitIdAsync() {
-            return 5;
+          protected RequestCommitIndexResponse requestCommitIdAsync() {
+            return new RequestCommitIndexResponse(5, 5, 5);
           }
 
           @Override
@@ -215,8 +216,8 @@ public class RaftMemberTest extends BaseMember {
           }
 
           @Override
-          protected long requestCommitIdAsync() {
-            return 1000L;
+          protected RequestCommitIndexResponse requestCommitIdAsync() {
+            return new RequestCommitIndexResponse(1000, 1000, 1000);
           }
 
           @Override
diff --git a/thrift-cluster/src/main/thrift/cluster.thrift 
b/thrift-cluster/src/main/thrift/cluster.thrift
index c8edbe3..f23130e 100644
--- a/thrift-cluster/src/main/thrift/cluster.thrift
+++ b/thrift-cluster/src/main/thrift/cluster.thrift
@@ -58,6 +58,12 @@ struct HeartBeatResponse {
   7: optional Node header
 }
 
+struct RequestCommitIndexResponse {
+  1: required long term // leader's meta log
+  2: required long commitLogIndex  // leader's meta log
+  3: required long commitLogTerm
+}
+
 // node -> node
 struct ElectionRequest {
   1: required long term
@@ -311,7 +317,7 @@ service RaftService {
   * Ask the leader for its commit index, used to check whether the node has 
caught up with the
   * leader.
   **/
-  long requestCommitIndex(1:Node header)
+  RequestCommitIndexResponse requestCommitIndex(1:Node header)
 
 
   /**

Reply via email to