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 f869a95cfef [IOTDB-6119] Add ConfigNode leader service check (#10985)
f869a95cfef is described below

commit f869a95cfef2f11f6cf9048ce084242c0e5ae5db
Author: Potato <[email protected]>
AuthorDate: Wed Aug 30 12:02:42 2023 +0800

    [IOTDB-6119] Add ConfigNode leader service check (#10985)
---
 .../statemachine/ConfigRegionStateMachine.java     | 56 +++++++++++-----------
 .../manager/consensus/ConsensusManager.java        | 25 ++++++----
 .../iotdb/confignode/manager/cq/CQManager.java     |  9 ----
 .../org/apache/iotdb/consensus/IStateMachine.java  |  5 ++
 .../apache/iotdb/consensus/config/RatisConfig.java | 48 +++++++++++++++++--
 .../ratis/ApplicationStateMachineProxy.java        |  5 ++
 .../apache/iotdb/consensus/ratis/utils/Utils.java  |  3 ++
 .../simple/SimpleConsensusServerImpl.java          |  1 +
 .../iotdb/db/protocol/client/ConfigNodeClient.java |  3 +-
 9 files changed, 106 insertions(+), 49 deletions(-)

diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/statemachine/ConfigRegionStateMachine.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/statemachine/ConfigRegionStateMachine.java
index fdec21c1c00..d097c7dadec 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/statemachine/ConfigRegionStateMachine.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/statemachine/ConfigRegionStateMachine.java
@@ -198,33 +198,7 @@ public class ConfigRegionStateMachine
     // We get currentNodeId here because the currentNodeId
     // couldn't initialize earlier than the ConfigRegionStateMachine
     int currentNodeId = 
ConfigNodeDescriptor.getInstance().getConf().getConfigNodeId();
-
-    if (currentNodeId == newLeaderId) {
-      LOGGER.info(
-          "Current node [nodeId: {}, ip:port: {}] becomes Leader",
-          newLeaderId,
-          currentNodeTEndPoint);
-
-      // Always start load services first
-      configManager.getLoadManager().startLoadServices();
-
-      // Start leader scheduling services
-      configManager.getProcedureManager().shiftExecutor(true);
-      configManager.getRetryFailedTasksThread().startRetryFailedTasksService();
-      configManager.getPartitionManager().startRegionCleaner();
-
-      // we do cq recovery async for two reasons:
-      // 1. For performance: cq recovery may be time-consuming, we use another 
thread to do it in
-      // make notifyLeaderChanged not blocked by it
-      // 2. For correctness: in cq recovery processing, it will use 
ConsensusManager which may be
-      // initialized after notifyLeaderChanged finished
-      threadPool.submit(() -> configManager.getCQManager().startCQScheduler());
-
-      threadPool.submit(
-          () -> 
configManager.getPipeManager().getPipeRuntimeCoordinator().startPipeMetaSync());
-      threadPool.submit(
-          () -> 
configManager.getPipeManager().getPipeRuntimeCoordinator().startPipeHeartbeat());
-    } else {
+    if (currentNodeId != newLeaderId) {
       LOGGER.info(
           "Current node [nodeId:{}, ip:port: {}] is not longer the leader, "
               + "the new leader is [nodeId:{}]",
@@ -244,6 +218,34 @@ public class ConfigRegionStateMachine
     }
   }
 
+  @Override
+  public void notifyLeaderReady() {
+    LOGGER.info(
+        "Current node [nodeId: {}, ip:port: {}] becomes Leader",
+        ConfigNodeDescriptor.getInstance().getConf().getConfigNodeId(),
+        currentNodeTEndPoint);
+
+    // Always start load services first
+    configManager.getLoadManager().startLoadServices();
+
+    // Start leader scheduling services
+    configManager.getProcedureManager().shiftExecutor(true);
+    configManager.getRetryFailedTasksThread().startRetryFailedTasksService();
+    configManager.getPartitionManager().startRegionCleaner();
+
+    // we do cq recovery async for two reasons:
+    // 1. For performance: cq recovery may be time-consuming, we use another 
thread to do it in
+    // make notifyLeaderChanged not blocked by it
+    // 2. For correctness: in cq recovery processing, it will use 
ConsensusManager which may be
+    // initialized after notifyLeaderChanged finished
+    threadPool.submit(() -> configManager.getCQManager().startCQScheduler());
+
+    threadPool.submit(
+        () -> 
configManager.getPipeManager().getPipeRuntimeCoordinator().startPipeMetaSync());
+    threadPool.submit(
+        () -> 
configManager.getPipeManager().getPipeRuntimeCoordinator().startPipeHeartbeat());
+  }
+
   @Override
   public void start() {
     if 
(ConsensusFactory.SIMPLE_CONSENSUS.equals(CONF.getConfigNodeConsensusProtocolClass()))
 {
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/consensus/ConsensusManager.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/consensus/ConsensusManager.java
index 16e90ee0712..744ef7718ea 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/consensus/ConsensusManager.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/consensus/ConsensusManager.java
@@ -336,6 +336,10 @@ public class ConsensusManager {
     return consensusImpl.isLeader(DEFAULT_CONSENSUS_GROUP_ID);
   }
 
+  public boolean isLeaderReady() {
+    return consensusImpl.isLeaderReady(DEFAULT_CONSENSUS_GROUP_ID);
+  }
+
   /** @return ConfigNode-leader's location if leader exists, null otherwise. */
   public TConfigNodeLocation getLeader() {
     for (int retry = 0; retry < 50; retry++) {
@@ -366,25 +370,28 @@ public class ConsensusManager {
   /**
    * Confirm the current ConfigNode's leadership.
    *
-   * @return SUCCESS_STATUS if the current ConfigNode is leader, 
NEED_REDIRECTION otherwise
+   * @return SUCCESS_STATUS if the current ConfigNode is leader and has been 
ready yet,
+   *     NEED_REDIRECTION otherwise
    */
   public TSStatus confirmLeader() {
     TSStatus result = new TSStatus();
-
-    if (isLeader()) {
-      return result.setCode(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+    if (isLeaderReady()) {
+      result.setCode(TSStatusCode.SUCCESS_STATUS.getStatusCode());
     } else {
       result.setCode(TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode());
-      result.setMessage(
-          "The current ConfigNode is not leader, please redirect to a new 
ConfigNode.");
-
+      if (isLeader()) {
+        result.setMessage(
+            "The current ConfigNode is leader but not ready yet, please try 
again later.");
+      } else {
+        result.setMessage(
+            "The current ConfigNode is not leader, please redirect to a new 
ConfigNode.");
+      }
       TConfigNodeLocation leaderLocation = getLeader();
       if (leaderLocation != null) {
         result.setRedirectNode(leaderLocation.getInternalEndPoint());
       }
-
-      return result;
     }
+    return result;
   }
 
   public ConsensusGroupId getConsensusGroupId() {
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/cq/CQManager.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/cq/CQManager.java
index cb7f9063965..02c6e8df242 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/cq/CQManager.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/cq/CQManager.java
@@ -114,15 +114,6 @@ public class CQManager {
   }
 
   public void startCQScheduler() {
-    try {
-      /*
-        TODO: remove this after fixing IOTDB-6085
-        sleep here because IOTDB-6085: NullPointerException in readAsync when 
Ratis leader is changing
-      */
-      Thread.sleep(1000);
-    } catch (InterruptedException e) {
-      Thread.currentThread().interrupt();
-    }
     lock.writeLock().lock();
     try {
       // 1. shutdown previous cq schedule thread pool
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/IStateMachine.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/IStateMachine.java
index 056c897afff..9b9dd46a6da 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/IStateMachine.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/IStateMachine.java
@@ -168,6 +168,11 @@ public interface IStateMachine {
     default void notifyConfigurationChanged(long term, long index, List<Peer> 
newConfiguration) {
       // do nothing default
     }
+
+    /** Notify the {@link IStateMachine} that this server becomes ready after 
changed to leader. */
+    default void notifyLeaderReady() {
+      // do nothing default
+    }
   }
 
   /**
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/RatisConfig.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/RatisConfig.java
index 08bc80d9ed7..c1b97bfceea 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/RatisConfig.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/config/RatisConfig.java
@@ -41,6 +41,7 @@ public class RatisConfig {
   private final Impl impl;
   private final LeaderLogAppender leaderLogAppender;
   private final Read read;
+  private final Utils utils;
 
   private RatisConfig(
       Rpc rpc,
@@ -52,7 +53,8 @@ public class RatisConfig {
       Client client,
       Impl impl,
       LeaderLogAppender leaderLogAppender,
-      Read read) {
+      Read read,
+      Utils utils) {
     this.rpc = rpc;
     this.leaderElection = leaderElection;
     this.snapshot = snapshot;
@@ -63,6 +65,7 @@ public class RatisConfig {
     this.impl = impl;
     this.leaderLogAppender = leaderLogAppender;
     this.read = read;
+    this.utils = utils;
   }
 
   public Rpc getRpc() {
@@ -105,6 +108,10 @@ public class RatisConfig {
     return read;
   }
 
+  public Utils getUtils() {
+    return utils;
+  }
+
   public static Builder newBuilder() {
     return new Builder();
   }
@@ -117,11 +124,11 @@ public class RatisConfig {
     private ThreadPool threadPool;
     private Log log;
     private Grpc grpc;
-
     private Client client;
     private Impl impl;
     private LeaderLogAppender leaderLogAppender;
     private Read read;
+    private Utils utils;
 
     public RatisConfig build() {
       return new RatisConfig(
@@ -135,7 +142,8 @@ public class RatisConfig {
           Optional.ofNullable(impl).orElseGet(() -> Impl.newBuilder().build()),
           Optional.ofNullable(leaderLogAppender)
               .orElseGet(() -> LeaderLogAppender.newBuilder().build()),
-          Optional.ofNullable(read).orElseGet(() -> 
Read.newBuilder().build()));
+          Optional.ofNullable(read).orElseGet(() -> Read.newBuilder().build()),
+          Optional.ofNullable(utils).orElseGet(() -> 
Utils.newBuilder().build()));
     }
 
     public Builder setRpc(Rpc rpc) {
@@ -187,6 +195,10 @@ public class RatisConfig {
       this.read = read;
       return this;
     }
+
+    public void setUtils(Utils utils) {
+      this.utils = utils;
+    }
   }
 
   /** server rpc timeout related. */
@@ -1104,4 +1116,34 @@ public class RatisConfig {
       }
     }
   }
+
+  public static class Utils {
+
+    private final int sleepDeviationThresholdMs;
+
+    private Utils(int sleepDeviationThresholdMs) {
+      this.sleepDeviationThresholdMs = sleepDeviationThresholdMs;
+    }
+
+    public int getSleepDeviationThresholdMs() {
+      return sleepDeviationThresholdMs;
+    }
+
+    public static Utils.Builder newBuilder() {
+      return new Utils.Builder();
+    }
+
+    public static class Builder {
+
+      private int sleepDeviationThresholdMs = 4 * 1000;
+
+      public Utils build() {
+        return new Utils(sleepDeviationThresholdMs);
+      }
+
+      public void setSleepDeviationThresholdMs(int sleepDeviationThresholdMs) {
+        this.sleepDeviationThresholdMs = sleepDeviationThresholdMs;
+      }
+    }
+  }
 }
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/ApplicationStateMachineProxy.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/ApplicationStateMachineProxy.java
index 828d134d231..d7444690a89 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/ApplicationStateMachineProxy.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/ApplicationStateMachineProxy.java
@@ -318,6 +318,11 @@ public class ApplicationStateMachineProxy extends 
BaseStateMachine {
             Utils.fromRaftPeerIdToNodeId(newLeaderId));
   }
 
+  @Override
+  public void notifyLeaderReady() {
+    applicationStateMachine.event().notifyLeaderReady();
+  }
+
   @Override
   public void notifyConfigurationChanged(
       long term, long index, RaftConfigurationProto newRaftConfiguration) {
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/utils/Utils.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/utils/Utils.java
index 14400b7d145..3026cb9cd95 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/utils/Utils.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/ratis/utils/Utils.java
@@ -289,5 +289,8 @@ public class Utils {
             : RaftServerConfigKeys.Read.Option.LINEARIZABLE;
     RaftServerConfigKeys.Read.setOption(properties, option);
     RaftServerConfigKeys.Read.setTimeout(properties, 
config.getRead().getReadTimeout());
+
+    RaftServerConfigKeys.setSleepDeviationThreshold(
+        properties, config.getUtils().getSleepDeviationThresholdMs());
   }
 }
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/simple/SimpleConsensusServerImpl.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/simple/SimpleConsensusServerImpl.java
index b4d8fe87040..758d4f31081 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/simple/SimpleConsensusServerImpl.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/simple/SimpleConsensusServerImpl.java
@@ -50,6 +50,7 @@ public class SimpleConsensusServerImpl implements 
IStateMachine {
     stateMachine.start();
     // Notify itself as the leader
     stateMachine.event().notifyLeaderChanged(peer.getGroupId(), 
peer.getNodeId());
+    stateMachine.event().notifyLeaderReady();
   }
 
   @Override
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/client/ConfigNodeClient.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/client/ConfigNodeClient.java
index d9390ab1976..e731c7dfb59 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/client/ConfigNodeClient.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/client/ConfigNodeClient.java
@@ -307,7 +307,8 @@ public class ConfigNodeClient implements 
IConfigNodeRPCService.Iface, ThriftClie
         configLeader = null;
       }
       logger.warn(
-          "Failed to connect to ConfigNode {} from DataNode {}, because the 
current node is not leader, try next node",
+          "Failed to connect to ConfigNode {} from DataNode {}, because the 
current node is not "
+              + "leader or not ready yet, will try again later",
           configNode,
           config.getAddressAndPort());
       return true;

Reply via email to