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;
