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

tanxinyu 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 6c24a24  [IOTDB-2392] Memory control of raft log in cluster  (#4825)
6c24a24 is described below

commit 6c24a242b140e0afcede846ab95bff2ed434572b
Author: Mrquan <[email protected]>
AuthorDate: Sat Jan 15 15:10:29 2022 +0800

    [IOTDB-2392] Memory control of raft log in cluster  (#4825)
    
    * fix the client ip bug in cluster
    
    * fix a typo
    
    * Add an example for Cluster setup on 3 nodes
    
    * memory control for raft log
    
    * memory control for raft log
    
    * memory control for raft log
    
    * memory control for raft log
    
    * memory control for raft log
    
    Co-authored-by: 权思屹 <[email protected]>
---
 .../resources/conf/iotdb-cluster.properties        |  16 ++-
 .../org/apache/iotdb/cluster/ClusterIoTDB.java     |  26 +++-
 .../apache/iotdb/cluster/config/ClusterConfig.java |  40 ++++++
 .../iotdb/cluster/config/ClusterDescriptor.java    |  20 ++-
 .../org/apache/iotdb/cluster/server/Response.java  |   3 +-
 .../server/handlers/caller/LogCatchUpHandler.java  |  11 +-
 .../iotdb/cluster/server/member/RaftMember.java    | 154 +++++++++++++++------
 7 files changed, 212 insertions(+), 58 deletions(-)

diff --git a/cluster/src/assembly/resources/conf/iotdb-cluster.properties 
b/cluster/src/assembly/resources/conf/iotdb-cluster.properties
index 9bcf19a..ac06f6e 100644
--- a/cluster/src/assembly/resources/conf/iotdb-cluster.properties
+++ b/cluster/src/assembly/resources/conf/iotdb-cluster.properties
@@ -113,10 +113,10 @@ multi_raft_factor=1
 # memory footprint
 # max_num_of_logs_in_mem=2000
 
-# maximum memory size of committed logs in memory, when reached, a log 
deletion will be triggered.
-# Increasing the number will reduce the chance to use snapshot in catch-ups, 
but will also increase
-# memory footprint, default is 512MB
-# max_memory_size_for_raft_log=536870912
+# Ratio of write memory allocated for raft log, 0.2 by default
+# Increasing the number will reduce the memory allocated for write process in 
iotdb, but will also
+# increase the memory footprint for raft log, which reduces the chance to use 
snapshot in catch-ups
+# raft_log_memory_proportion=0.2
 
 # deletion check period of the submitted log
 # log_deletion_check_interval_second=-1
@@ -152,6 +152,14 @@ multi_raft_factor=1
 # These indexes are used to index the location of the log on the disk
 # max_raft_log_index_size_in_memory=10000
 
+# If leader finds too many uncommitted raft logs, raft group leader will wait 
for a short period of
+# time, and then append the raft log
+# uncommitted_raft_log_num_for_reject_threshold=500
+
+# If followers find too many committed raft logs have not been applied, 
followers will reject the raft
+# log sent by leader
+# unapplied_raft_log_num_for_reject_threshold=500
+
 # The maximum size of the raft log saved on disk for each file (in bytes) of 
each raft group.
 # The default size is 1GB
 # max_raft_log_persist_data_size_per_file=1073741824
diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
index 008854a..b15635b 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
@@ -57,6 +57,7 @@ import 
org.apache.iotdb.cluster.server.service.MetaSyncService;
 import org.apache.iotdb.cluster.utils.ClusterUtils;
 import org.apache.iotdb.cluster.utils.nodetool.ClusterMonitor;
 import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory;
+import org.apache.iotdb.db.conf.IoTDBConfig;
 import org.apache.iotdb.db.conf.IoTDBConfigCheck;
 import org.apache.iotdb.db.conf.IoTDBConstant;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
@@ -276,22 +277,35 @@ public class ClusterIoTDB implements ClusterIoTDBMBean {
 
   private boolean serverCheckAndInit() throws ConfigurationException, 
IOException {
     IoTDBConfigCheck.getInstance().checkConfig();
+    IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
     // init server's configuration first, because the cluster configuration 
may read settings from
     // the server's configuration.
-    IoTDBDescriptor.getInstance().getConfig().setSyncEnable(false);
+    config.setSyncEnable(false);
     // auto create schema is took over by cluster module, so we disable it in 
the server module.
-    
IoTDBDescriptor.getInstance().getConfig().setAutoCreateSchemaEnabled(false);
+    config.setAutoCreateSchemaEnabled(false);
     // check cluster config
     String checkResult = clusterConfigCheck();
     if (checkResult != null) {
       logger.error(checkResult);
       return false;
     }
+    ClusterConfig clusterConfig = ClusterDescriptor.getInstance().getConfig();
     // if client ip is the default address, set it same with internal ip
-    if 
(IoTDBDescriptor.getInstance().getConfig().getRpcAddress().equals("0.0.0.0")) {
-      IoTDBDescriptor.getInstance()
-          .getConfig()
-          
.setRpcAddress(ClusterDescriptor.getInstance().getConfig().getInternalIp());
+    if (config.getRpcAddress().equals("0.0.0.0")) {
+      config.setRpcAddress(clusterConfig.getInternalIp());
+    }
+    // set the memory allocated for raft log of each raft log manager
+    if (clusterConfig.getReplicationNum() > 1) {
+      clusterConfig.setMaxMemorySizeForRaftLog(
+          (long)
+              (config.getAllocateMemoryForWrite()
+                  * clusterConfig.getRaftLogMemoryProportion()
+                  / clusterConfig.getReplicationNum()));
+      // calculate remaining memory allocated for write process
+      config.setAllocateMemoryForWrite(
+          (long)
+              (config.getAllocateMemoryForWrite()
+                  * (1 - clusterConfig.getRaftLogMemoryProportion())));
     }
     return true;
   }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java
index 84a97a3..ac91f8c 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java
@@ -78,6 +78,9 @@ public class ClusterConfig {
   /** max memory size of committed logs in memory, default 512M */
   private long maxMemorySizeForRaftLog = 536870912;
 
+  /** Ratio of write memory allocated for raft log */
+  private double RaftLogMemoryProportion = 0.2;
+
   /** deletion check period of the submitted log */
   private int logDeleteCheckIntervalSecond = -1;
 
@@ -137,6 +140,18 @@ public class ClusterConfig {
   private int maxRaftLogIndexSizeInMemory = 10000;
 
   /**
+   * If leader finds too many uncommitted raft logs, raft group leader will 
wait for a short period
+   * of time, and then append the raft log
+   */
+  private int UnCommittedRaftLogNumForRejectThreshold = 500;
+
+  /**
+   * If followers find too many committed raft logs have not been applied, 
followers will reject the
+   * raft log sent by leader
+   */
+  private int UnAppliedRaftLogNumForRejectThreshold = 500;
+
+  /**
    * The maximum size of the raft log saved on disk for each file (in bytes) 
of each raft group. The
    * default size is 1GB
    */
@@ -375,6 +390,23 @@ public class ClusterConfig {
     this.maxNumOfLogsInMem = maxNumOfLogsInMem;
   }
 
+  public int getUnCommittedRaftLogNumForRejectThreshold() {
+    return UnCommittedRaftLogNumForRejectThreshold;
+  }
+
+  public void setUnCommittedRaftLogNumForRejectThreshold(
+      int unCommittedRaftLogNumForRejectThreshold) {
+    UnCommittedRaftLogNumForRejectThreshold = 
unCommittedRaftLogNumForRejectThreshold;
+  }
+
+  public int getUnAppliedRaftLogNumForRejectThreshold() {
+    return UnAppliedRaftLogNumForRejectThreshold;
+  }
+
+  public void setUnAppliedRaftLogNumForRejectThreshold(int 
unAppliedRaftLogNumForRejectThreshold) {
+    UnAppliedRaftLogNumForRejectThreshold = 
unAppliedRaftLogNumForRejectThreshold;
+  }
+
   public int getRaftLogBufferSize() {
     return raftLogBufferSize;
   }
@@ -423,6 +455,14 @@ public class ClusterConfig {
     this.maxMemorySizeForRaftLog = maxMemorySizeForRaftLog;
   }
 
+  public double getRaftLogMemoryProportion() {
+    return RaftLogMemoryProportion;
+  }
+
+  public void setRaftLogMemoryProportion(double raftLogMemoryProportion) {
+    RaftLogMemoryProportion = raftLogMemoryProportion;
+  }
+
   public int getMaxRaftLogPersistDataSizePerFile() {
     return maxRaftLogPersistDataSizePerFile;
   }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterDescriptor.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterDescriptor.java
index 0ebac8c..bf5385a 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterDescriptor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterDescriptor.java
@@ -215,11 +215,11 @@ public class ClusterDescriptor {
             properties.getProperty(
                 "max_num_of_logs_in_mem", 
String.valueOf(config.getMaxNumOfLogsInMem()))));
 
-    config.setMaxMemorySizeForRaftLog(
-        Long.parseLong(
+    config.setRaftLogMemoryProportion(
+        Double.parseDouble(
             properties.getProperty(
-                "max_memory_size_for_raft_log",
-                String.valueOf(config.getMaxMemorySizeForRaftLog()))));
+                "raft_log_memory_proportion",
+                String.valueOf(config.getRaftLogMemoryProportion()))));
 
     config.setLogDeleteCheckIntervalSecond(
         Integer.parseInt(
@@ -269,6 +269,18 @@ public class ClusterDescriptor {
                 "max_raft_log_index_size_in_memory",
                 String.valueOf(config.getMaxRaftLogIndexSizeInMemory()))));
 
+    config.setUnCommittedRaftLogNumForRejectThreshold(
+        Integer.parseInt(
+            properties.getProperty(
+                "uncommitted_raft_log_num_for_reject_threshold",
+                
String.valueOf(config.getUnCommittedRaftLogNumForRejectThreshold()))));
+
+    config.setUnAppliedRaftLogNumForRejectThreshold(
+        Integer.parseInt(
+            properties.getProperty(
+                "unapplied_raft_log_num_for_reject_threshold",
+                
String.valueOf(config.getUnAppliedRaftLogNumForRejectThreshold()))));
+
     config.setMaxRaftLogPersistDataSizePerFile(
         Integer.parseInt(
             properties.getProperty(
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 387549d..6b286c1 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,9 +52,10 @@ public class Response {
   public static final long RESPONSE_NODE_IS_NOT_IN_GROUP = -11;
   // the request is not executed locally anc should be forwarded
   public static final long RESPONSE_NULL = Long.MIN_VALUE;
-
   // the meta engine is not ready (except for the partitionTable is ready)
   public static final long RESPONSE_META_NOT_READY = -12;
+  // the cluster is too busy to reject new committed logs
+  public static final long RESPONSE_TOO_BUSY = -13;
 
   private Response() {
     // enum-like class
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/handlers/caller/LogCatchUpHandler.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/handlers/caller/LogCatchUpHandler.java
index ddf3af1..79acb2b 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/handlers/caller/LogCatchUpHandler.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/handlers/caller/LogCatchUpHandler.java
@@ -29,8 +29,7 @@ import org.slf4j.LoggerFactory;
 
 import java.util.concurrent.atomic.AtomicBoolean;
 
-import static org.apache.iotdb.cluster.server.Response.RESPONSE_AGREE;
-import static org.apache.iotdb.cluster.server.Response.RESPONSE_LOG_MISMATCH;
+import static org.apache.iotdb.cluster.server.Response.*;
 
 /**
  * LogCatchUpHandler checks the result of appending a log in a catch-up task 
and decides to abort
@@ -64,6 +63,14 @@ public class LogCatchUpHandler implements 
AsyncMethodCallback<Long> {
         appendSucceed.set(true);
         appendSucceed.notifyAll();
       }
+    } else if (resp == RESPONSE_TOO_BUSY) {
+      // this may occur when the follower has too many logs unapplied, so we 
abort the
+      // catch-up task
+      logger.info("{}: Catchup task rejected by receiver {}", memberName, 
follower);
+      synchronized (appendSucceed) {
+        appendSucceed.set(false);
+        appendSucceed.notifyAll();
+      }
     } else {
       // the follower's term has updated, which means a new leader is elected
       logger.debug("{}: Received a rejection because term is updated to: {}", 
memberName, resp);
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 7901370..7124251 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
@@ -68,6 +68,7 @@ import org.apache.iotdb.cluster.utils.StatusUtils;
 import org.apache.iotdb.db.concurrent.IoTDBThreadPoolFactory;
 import org.apache.iotdb.db.concurrent.IoTThreadFactory;
 import org.apache.iotdb.db.conf.IoTDBConstant;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.exception.BatchProcessException;
 import org.apache.iotdb.db.exception.IoTDBException;
 import org.apache.iotdb.db.exception.metadata.DuplicatedTemplateException;
@@ -991,16 +992,32 @@ public abstract class RaftMember implements 
RaftMemberMBean {
       return StatusUtils.INTERNAL_ERROR;
     }
 
-    // assign term and index to the new log and append it
-    synchronized (logManager) {
-      if (!(plan instanceof LogPlan)) {
-        plan.setIndex(logManager.getLastLogIndex() + 1);
+    long startWaitingTime = System.currentTimeMillis();
+    while (true) {
+      // assign term and index to the new log and append it
+      synchronized (logManager) {
+        if (logManager.getLastLogIndex() - logManager.getCommitLogIndex()
+            <= config.getUnCommittedRaftLogNumForRejectThreshold()) {
+          if (!(plan instanceof LogPlan)) {
+            plan.setIndex(logManager.getLastLogIndex() + 1);
+          }
+          log.setCurrLogTerm(getTerm().get());
+          log.setCurrLogIndex(logManager.getLastLogIndex() + 1);
+          logManager.append(log);
+          break;
+        }
+      }
+      try {
+        TimeUnit.MILLISECONDS.sleep(
+            
IoTDBDescriptor.getInstance().getConfig().getCheckPeriodWhenInsertBlocked());
+        if (System.currentTimeMillis() - startWaitingTime
+            > 
IoTDBDescriptor.getInstance().getConfig().getMaxWaitingTimeWhenInsertBlocked()) 
{
+          return StatusUtils.getStatus(TSStatusCode.WRITE_PROCESS_REJECT);
+        }
+      } catch (InterruptedException e) {
+        Thread.currentThread().interrupt();
       }
-      log.setCurrLogTerm(getTerm().get());
-      log.setCurrLogIndex(logManager.getLastLogIndex() + 1);
-      logManager.append(log);
     }
-
     
Timer.Statistic.RAFT_SENDER_APPEND_LOG.calOperationCostTimeFromStart(startTime);
 
     try {
@@ -1043,25 +1060,40 @@ public abstract class RaftMember implements 
RaftMemberMBean {
               + "or reduce the size of requests you send.");
       return StatusUtils.INTERNAL_ERROR;
     }
-
     long startTime =
         
Statistic.RAFT_SENDER_COMPETE_LOG_MANAGER_BEFORE_APPEND_V2.getOperationStartTime();
-    synchronized (logManager) {
-      
Statistic.RAFT_SENDER_COMPETE_LOG_MANAGER_BEFORE_APPEND_V2.calOperationCostTimeFromStart(
-          startTime);
-
-      if (!(plan instanceof LogPlan)) {
-        plan.setIndex(logManager.getLastLogIndex() + 1);
+    long startWaitingTime = System.currentTimeMillis();
+    while (true) {
+      synchronized (logManager) {
+        if (!IoTDBDescriptor.getInstance().getConfig().isEnableMemControl()
+            || (logManager.getLastLogIndex() - logManager.getCommitLogIndex()
+                <= config.getUnCommittedRaftLogNumForRejectThreshold())) {
+          
Statistic.RAFT_SENDER_COMPETE_LOG_MANAGER_BEFORE_APPEND_V2.calOperationCostTimeFromStart(
+              startTime);
+          if (!(plan instanceof LogPlan)) {
+            plan.setIndex(logManager.getLastLogIndex() + 1);
+          }
+          log.setCurrLogTerm(getTerm().get());
+          log.setCurrLogIndex(logManager.getLastLogIndex() + 1);
+          startTime = 
Timer.Statistic.RAFT_SENDER_APPEND_LOG_V2.getOperationStartTime();
+          logManager.append(log);
+          
Timer.Statistic.RAFT_SENDER_APPEND_LOG_V2.calOperationCostTimeFromStart(startTime);
+          startTime = 
Statistic.RAFT_SENDER_BUILD_LOG_REQUEST.getOperationStartTime();
+          sendLogRequest = buildSendLogRequest(log);
+          
Statistic.RAFT_SENDER_BUILD_LOG_REQUEST.calOperationCostTimeFromStart(startTime);
+          break;
+        }
+      }
+      try {
+        TimeUnit.MILLISECONDS.sleep(
+            
IoTDBDescriptor.getInstance().getConfig().getCheckPeriodWhenInsertBlocked());
+        if (System.currentTimeMillis() - startWaitingTime
+            > 
IoTDBDescriptor.getInstance().getConfig().getMaxWaitingTimeWhenInsertBlocked()) 
{
+          return StatusUtils.getStatus(TSStatusCode.WRITE_PROCESS_REJECT);
+        }
+      } catch (InterruptedException e) {
+        Thread.currentThread().interrupt();
       }
-      log.setCurrLogTerm(getTerm().get());
-      log.setCurrLogIndex(logManager.getLastLogIndex() + 1);
-
-      startTime = 
Timer.Statistic.RAFT_SENDER_APPEND_LOG_V2.getOperationStartTime();
-      logManager.append(log);
-      
Timer.Statistic.RAFT_SENDER_APPEND_LOG_V2.calOperationCostTimeFromStart(startTime);
-      startTime = 
Statistic.RAFT_SENDER_BUILD_LOG_REQUEST.getOperationStartTime();
-      sendLogRequest = buildSendLogRequest(log);
-      
Statistic.RAFT_SENDER_BUILD_LOG_REQUEST.calOperationCostTimeFromStart(startTime);
     }
 
     startTime = Statistic.RAFT_SENDER_OFFER_LOG.getOperationStartTime();
@@ -1940,10 +1972,12 @@ public abstract class RaftMember implements 
RaftMemberMBean {
 
   /**
    * Find the local previous log of "log". If such log is found, discard all 
local logs behind it
-   * and append "log" to it. Otherwise report a log mismatch.
+   * and append "log" to it. Otherwise report a log mismatch. If too many 
committed logs have not
+   * been applied, reject the appendEntry request.
    *
    * @return Response.RESPONSE_AGREE when the log is successfully appended or 
Response
-   *     .RESPONSE_LOG_MISMATCH if the previous log of "log" is not found.
+   *     .RESPONSE_LOG_MISMATCH if the previous log of "log" is not found or 
Response
+   *     .RESPONSE_TOO_BUSY if too many committed logs have not been applied.
    */
   protected long appendEntry(long prevLogIndex, long prevLogTerm, long 
leaderCommit, Log log) {
     long resp = checkPrevLogIndex(prevLogIndex);
@@ -1952,9 +1986,27 @@ public abstract class RaftMember implements 
RaftMemberMBean {
     }
 
     long startTime = 
Timer.Statistic.RAFT_RECEIVER_APPEND_ENTRY.getOperationStartTime();
+    long startWaitingTime = System.currentTimeMillis();
     long success;
-    synchronized (logManager) {
-      success = logManager.maybeAppend(prevLogIndex, prevLogTerm, 
leaderCommit, log);
+    while (true) {
+      synchronized (logManager) {
+        // TODO: Consider memory footprint to execute a precise rejection
+        if ((logManager.getCommitLogIndex() - 
logManager.getMaxHaveAppliedCommitIndex())
+            <= config.getUnAppliedRaftLogNumForRejectThreshold()) {
+          success = logManager.maybeAppend(prevLogIndex, prevLogTerm, 
leaderCommit, log);
+          break;
+        }
+        try {
+          TimeUnit.MILLISECONDS.sleep(
+              
IoTDBDescriptor.getInstance().getConfig().getCheckPeriodWhenInsertBlocked());
+          if (System.currentTimeMillis() - startWaitingTime
+              > 
IoTDBDescriptor.getInstance().getConfig().getMaxWaitingTimeWhenInsertBlocked()) 
{
+            return Response.RESPONSE_TOO_BUSY;
+          }
+        } catch (InterruptedException e) {
+          Thread.currentThread().interrupt();
+        }
+      }
     }
     
Timer.Statistic.RAFT_RECEIVER_APPEND_ENTRY.calOperationCostTimeFromStart(startTime);
     if (success != -1) {
@@ -2010,11 +2062,13 @@ public abstract class RaftMember implements 
RaftMemberMBean {
 
   /**
    * Find the local previous log of "log". If such log is found, discard all 
local logs behind it
-   * and append "log" to it. Otherwise report a log mismatch.
+   * and append "log" to it. Otherwise report a log mismatch. If too many 
committed logs have not
+   * been applied, reject the appendEntry request.
    *
    * @param logs append logs
    * @return Response.RESPONSE_AGREE when the log is successfully appended or 
Response
-   *     .RESPONSE_LOG_MISMATCH if the previous log of "log" is not found.
+   *     .RESPONSE_LOG_MISMATCH if the previous log of "log" is not found 
Response
+   *     .RESPONSE_TOO_BUSY if too many committed logs have not been applied.
    */
   private long appendEntries(
       long prevLogIndex, long prevLogTerm, long leaderCommit, List<Log> logs) {
@@ -2033,18 +2087,36 @@ public abstract class RaftMember implements 
RaftMemberMBean {
       return resp;
     }
 
-    synchronized (logManager) {
-      long startTime = 
Timer.Statistic.RAFT_RECEIVER_APPEND_ENTRY.getOperationStartTime();
-      resp = logManager.maybeAppend(prevLogIndex, prevLogTerm, leaderCommit, 
logs);
-      
Timer.Statistic.RAFT_RECEIVER_APPEND_ENTRY.calOperationCostTimeFromStart(startTime);
-      if (resp != -1) {
-        if (logger.isDebugEnabled()) {
-          logger.debug("{} append a new log list {}, commit to {}", name, 
logs, leaderCommit);
+    long startWaitingTime = System.currentTimeMillis();
+    while (true) {
+      synchronized (logManager) {
+        // TODO: Consider memory footprint to execute a precise rejection
+        if ((logManager.getCommitLogIndex() - 
logManager.getMaxHaveAppliedCommitIndex())
+            <= config.getUnAppliedRaftLogNumForRejectThreshold()) {
+          long startTime = 
Timer.Statistic.RAFT_RECEIVER_APPEND_ENTRY.getOperationStartTime();
+          resp = logManager.maybeAppend(prevLogIndex, prevLogTerm, 
leaderCommit, logs);
+          
Timer.Statistic.RAFT_RECEIVER_APPEND_ENTRY.calOperationCostTimeFromStart(startTime);
+          if (resp != -1) {
+            if (logger.isDebugEnabled()) {
+              logger.debug("{} append a new log list {}, commit to {}", name, 
logs, leaderCommit);
+            }
+            resp = Response.RESPONSE_AGREE;
+          } else {
+            // the incoming log points to an illegal position, reject it
+            resp = Response.RESPONSE_LOG_MISMATCH;
+          }
+          break;
         }
-        resp = Response.RESPONSE_AGREE;
-      } else {
-        // the incoming log points to an illegal position, reject it
-        resp = Response.RESPONSE_LOG_MISMATCH;
+      }
+      try {
+        TimeUnit.MILLISECONDS.sleep(
+            
IoTDBDescriptor.getInstance().getConfig().getCheckPeriodWhenInsertBlocked());
+        if (System.currentTimeMillis() - startWaitingTime
+            > 
IoTDBDescriptor.getInstance().getConfig().getMaxWaitingTimeWhenInsertBlocked()) 
{
+          return Response.RESPONSE_TOO_BUSY;
+        }
+      } catch (InterruptedException e) {
+        Thread.currentThread().interrupt();
       }
     }
     return resp;

Reply via email to