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

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

commit ba4355a030af3952db3239699d1f631d3a9311aa
Author: jt <[email protected]>
AuthorDate: Wed Dec 1 09:15:22 2021 +0800

    add RESPONSE_OUT_OF_WINDOW
---
 .../java/org/apache/iotdb/cluster/expr/ExprMember.java   |  2 +-
 .../java/org/apache/iotdb/cluster/log/LogDispatcher.java | 16 +++++++++++++---
 .../java/org/apache/iotdb/cluster/server/Response.java   |  1 +
 3 files changed, 15 insertions(+), 4 deletions(-)

diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/expr/ExprMember.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/expr/ExprMember.java
index a2c62ed..04e6b97 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/expr/ExprMember.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/expr/ExprMember.java
@@ -285,7 +285,7 @@ public class ExprMember extends MetaGroupMember {
         Statistic.RAFT_WINDOW_LENGTH.add(windowLength);
       } else {
         
Timer.Statistic.RAFT_RECEIVER_APPEND_ENTRY.calOperationCostTimeFromStart(startTime);
-        result.setStatus(Response.RESPONSE_LOG_MISMATCH);
+        result.setStatus(Response.RESPONSE_OUT_OF_WINDOW);
         result.setHeader(getHeader());
         return result;
       }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/log/LogDispatcher.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/log/LogDispatcher.java
index 9558714..beaaa77 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/log/LogDispatcher.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/log/LogDispatcher.java
@@ -27,6 +27,7 @@ import org.apache.iotdb.cluster.rpc.thrift.AppendEntryResult;
 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.server.Response;
 import org.apache.iotdb.cluster.server.handlers.caller.AppendNodeEntryHandler;
 import org.apache.iotdb.cluster.server.member.RaftMember;
 import org.apache.iotdb.cluster.server.monitor.Peer;
@@ -422,11 +423,20 @@ public class LogDispatcher {
               logRequest.newLeaderTerm,
               peer,
               logRequest.quorumSize);
+      // TODO add async interface
+      int retries = 5;
       try {
         long operationStartTime = 
Statistic.RAFT_SENDER_SEND_LOG.getOperationStartTime();
-        AppendEntryResult result = 
client.appendEntry(logRequest.appendEntryRequest);
-        
Statistic.RAFT_SENDER_SEND_LOG.calOperationCostTimeFromStart(operationStartTime);
-        handler.onComplete(result);
+        for (int i = 0; i < retries; i++) {
+          AppendEntryResult result = 
client.appendEntry(logRequest.appendEntryRequest);
+          if (result.status == Response.RESPONSE_OUT_OF_WINDOW) {
+            Thread.sleep(100);
+          } else {
+            
Statistic.RAFT_SENDER_SEND_LOG.calOperationCostTimeFromStart(operationStartTime);
+            handler.onComplete(result);
+            break;
+          }
+        }
       } catch (TException e) {
         client.getInputProtocol().getTransport().close();
         ClientUtils.putBackSyncClient(client);
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/Response.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/Response.java
index e4ab4ef..b520248 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/server/Response.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/server/Response.java
@@ -52,6 +52,7 @@ public class Response {
   public static final long RESPONSE_NODE_IS_NOT_IN_GROUP = -11;
   public static final long RESPONSE_STRONG_ACCEPT = -12;
   public static final long RESPONSE_WEAK_ACCEPT = -13;
+  public static final long RESPONSE_OUT_OF_WINDOW = -14;
   // the request is not executed locally anc should be forwarded
   public static final long RESPONSE_NULL = Long.MIN_VALUE;
 

Reply via email to