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

caogaofei pushed a commit to branch beyyes/master1
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/beyyes/master1 by this push:
     new 20556d3de1 perfect remove datanode process
20556d3de1 is described below

commit 20556d3de1cac3ee53adfc07d4c6359de6c3210b
Author: Beyyes <[email protected]>
AuthorDate: Sun Nov 6 23:02:06 2022 +0800

    perfect remove datanode process
---
 .../iotdb/confignode/persistence/node/NodeInfo.java |  5 +++--
 .../persistence/partition/PartitionInfo.java        |  7 ++++---
 .../partition/StorageGroupPartitionTable.java       | 21 ++++++++++++++++-----
 .../procedure/env/DataNodeRemoveHandler.java        | 10 +++++-----
 .../impl/node/RemoveDataNodeProcedure.java          |  1 +
 .../impl/statemachine/RegionMigrateProcedure.java   |  7 ++++---
 .../iotdb/db/service/RegionMigrateService.java      | 18 ++++++++++--------
 7 files changed, 43 insertions(+), 26 deletions(-)

diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/node/NodeInfo.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/node/NodeInfo.java
index d03fffa7f9..2526abbad6 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/node/NodeInfo.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/node/NodeInfo.java
@@ -168,13 +168,14 @@ public class NodeInfo implements SnapshotProcessor {
         "{}, There are {} data node in cluster before executed 
remove-datanode.sh",
         REMOVE_DATANODE_PROCESS,
         registeredDataNodes.size());
+
+    dataNodeInfoReadWriteLock.writeLock().lock();
     try {
-      dataNodeInfoReadWriteLock.writeLock().lock();
       req.getDataNodeLocations()
           .forEach(
               removeDataNodes -> {
                 registeredDataNodes.remove(removeDataNodes.getDataNodeId());
-                LOGGER.info("removed the datanode {} from cluster", 
removeDataNodes);
+                LOGGER.info("Removed the datanode {} from cluster", 
removeDataNodes);
               });
     } finally {
       dataNodeInfoReadWriteLock.writeLock().unlock();
diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/PartitionInfo.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/PartitionInfo.java
index 6d2ca47e58..84115faae2 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/PartitionInfo.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/PartitionInfo.java
@@ -480,9 +480,10 @@ public class PartitionInfo implements SnapshotProcessor {
     TConsensusGroupId regionId = req.getRegionId();
     TDataNodeLocation oldNode = req.getOldNode();
     TDataNodeLocation newNode = req.getNewNode();
-    storageGroupPartitionTables
-        .values().stream().filter(sgPartitionTable -> 
sgPartitionTable.containRegion(regionId))
-        .forEach(sgPartitionTable -> 
sgPartitionTable.updateRegionLocation(regionId, oldNode, newNode));
+    storageGroupPartitionTables.values().stream()
+        .filter(sgPartitionTable -> sgPartitionTable.containRegion(regionId))
+        .forEach(
+            sgPartitionTable -> 
sgPartitionTable.updateRegionLocation(regionId, oldNode, newNode));
 
     return status;
   }
diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/StorageGroupPartitionTable.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/StorageGroupPartitionTable.java
index dd2e859e7e..b763fded02 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/StorageGroupPartitionTable.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/StorageGroupPartitionTable.java
@@ -432,12 +432,18 @@ public class StorageGroupPartitionTable {
   private void addRegionNewLocation(TConsensusGroupId regionId, 
TDataNodeLocation node) {
     RegionGroup regionGroup = regionGroupMap.get(regionId);
     if (regionGroup == null) {
-      LOGGER.warn("Cannot find RegionGroup for region {} when 
addRegionNewLocation in {}", regionId, storageGroupName);
+      LOGGER.warn(
+          "Cannot find RegionGroup for region {} when addRegionNewLocation in 
{}",
+          regionId,
+          storageGroupName);
       return;
     }
     if (regionGroup.getReplicaSet().getDataNodeLocations().contains(node)) {
-      LOGGER.info("Node is already in region locations when 
addRegionNewLocation in {}, node: {}, region: {}",
-              storageGroupName, node, regionId);
+      LOGGER.info(
+          "Node is already in region locations when addRegionNewLocation in 
{}, node: {}, region: {}",
+          storageGroupName,
+          node,
+          regionId);
       return;
     }
     regionGroup.getReplicaSet().getDataNodeLocations().add(node);
@@ -446,13 +452,18 @@ public class StorageGroupPartitionTable {
   private void removeRegionOldLocation(TConsensusGroupId regionId, 
TDataNodeLocation node) {
     RegionGroup regionGroup = regionGroupMap.get(regionId);
     if (regionGroup == null) {
-      LOGGER.warn("Cannot find RegionGroup for region {} when 
removeRegionOldLocation in {}", regionId, storageGroupName);
+      LOGGER.warn(
+          "Cannot find RegionGroup for region {} when removeRegionOldLocation 
in {}",
+          regionId,
+          storageGroupName);
       return;
     }
     if (!regionGroup.getReplicaSet().getDataNodeLocations().contains(node)) {
       LOGGER.info(
           "Node is not in region locations when removeRegionOldLocation in {}, 
no need to remove it, node: {}, region: {}",
-          storageGroupName, node, regionId);
+          storageGroupName,
+          node,
+          regionId);
       return;
     }
     regionGroup.getReplicaSet().getDataNodeLocations().remove(node);
diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/DataNodeRemoveHandler.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/DataNodeRemoveHandler.java
index 8a48624046..2e29166188 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/DataNodeRemoveHandler.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/DataNodeRemoveHandler.java
@@ -249,11 +249,11 @@ public class DataNodeRemoveHandler {
                 maintainPeerReq,
                 DataNodeRequestType.ADD_REGION_PEER);
     LOGGER.info(
-        "{}, Send action addRegionPeer finished, regionId: {}, rpcDataNode: 
{}, destDataNode: {}",
+        "{}, Send action addRegionPeer finished, regionId: {}, rpcDataNode: 
{},  destDataNode: {}",
         REMOVE_DATANODE_PROCESS,
         regionId,
         getIdWithRpcEndpoint(selectedDataNode.get()),
-        destDataNode);
+        getIdWithRpcEndpoint(destDataNode));
     return status;
   }
 
@@ -365,8 +365,8 @@ public class DataNodeRemoveHandler {
         "UpdateRegionLocationCache finished, region:{}, result:{}, old:{}, 
new:{}",
         regionId,
         status,
-            getIdWithRpcEndpoint(originalDataNode),
-            getIdWithRpcEndpoint(destDataNode));
+        getIdWithRpcEndpoint(originalDataNode),
+        getIdWithRpcEndpoint(destDataNode));
 
     // Broadcast the latest RegionRouteMap when Region migration finished
     configManager.getLoadManager().broadcastLatestRegionRouteMap();
@@ -412,7 +412,7 @@ public class DataNodeRemoveHandler {
    * @param dataNode old data node
    */
   public void stopDataNode(TDataNodeLocation dataNode) {
-    LOGGER.info("{}, Begin to stop Data Node {}", REMOVE_DATANODE_PROCESS, 
dataNode);
+    LOGGER.info("{}, Begin to stop DataNode {}", REMOVE_DATANODE_PROCESS, 
dataNode);
     
AsyncDataNodeClientPool.getInstance().resetClient(dataNode.getInternalEndPoint());
     TSStatus status =
         SyncDataNodeClientPool.getInstance()
diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/node/RemoveDataNodeProcedure.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/node/RemoveDataNodeProcedure.java
index fe59a56d97..6b7272c36f 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/node/RemoveDataNodeProcedure.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/node/RemoveDataNodeProcedure.java
@@ -95,6 +95,7 @@ public class RemoveDataNodeProcedure extends 
AbstractNodeProcedure<RemoveDataNod
           setNextState(RemoveDataNodeState.STOP_DATA_NODE);
           break;
         case STOP_DATA_NODE:
+          // TODO if region migrate is failed, don't execute STOP_DATA_NODE
           
env.getDataNodeRemoveHandler().removeDataNodePersistence(disableDataNodeLocation);
           env.getDataNodeRemoveHandler().stopDataNode(disableDataNodeLocation);
           return Flow.NO_MORE_STATE;
diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/statemachine/RegionMigrateProcedure.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/statemachine/RegionMigrateProcedure.java
index 22509193a7..da40d93400 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/statemachine/RegionMigrateProcedure.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/statemachine/RegionMigrateProcedure.java
@@ -40,6 +40,7 @@ import java.io.IOException;
 import java.nio.ByteBuffer;
 
 import static 
org.apache.iotdb.confignode.conf.ConfigNodeConstant.REMOVE_DATANODE_PROCESS;
+import static 
org.apache.iotdb.confignode.procedure.env.DataNodeRemoveHandler.getIdWithRpcEndpoint;
 import static org.apache.iotdb.rpc.TSStatusCode.SUCCESS_STATUS;
 
 /** region migrate procedure */
@@ -145,7 +146,7 @@ public class RegionMigrateProcedure
             "{}, Failed state [{}] is not support rollback, originalDataNode: 
{}",
             REMOVE_DATANODE_PROCESS,
             state,
-            originalDataNode);
+                getIdWithRpcEndpoint(originalDataNode));
         if (getCycles() > RETRY_THRESHOLD) {
           setFailure(
               new ProcedureException(
@@ -283,7 +284,7 @@ public class RegionMigrateProcedure
   public void notifyTheRegionMigrateFinished(TRegionMigrateResultReportReq 
req) {
 
     LOG.info(
-        "{}, ConfigNode received DataNode reported region migrate result: {}",
+        "{}, ConfigNode received region migrate result reported by DataNode: 
{}",
         REMOVE_DATANODE_PROCESS,
         req);
 
@@ -293,7 +294,7 @@ public class RegionMigrateProcedure
       // migrate failed
       if (migrateStatus.getCode() != SUCCESS_STATUS.getStatusCode()) {
         LOG.info(
-            "{}, Region migrate executed failed in DataNode, migrateStatus: 
{}",
+            "{}, Region migrate failed in DataNode, migrateStatus: {}",
             REMOVE_DATANODE_PROCESS,
             migrateStatus);
         migrateSuccess = false;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/service/RegionMigrateService.java 
b/server/src/main/java/org/apache/iotdb/db/service/RegionMigrateService.java
index a1f73fdb14..80d83d0711 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/RegionMigrateService.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/RegionMigrateService.java
@@ -52,7 +52,9 @@ import java.util.Map;
 public class RegionMigrateService implements IService {
   private static final Logger LOGGER = 
LoggerFactory.getLogger(RegionMigrateService.class);
 
-  private static final int RETRY = 5;
+  public static final String REMOVE_DATANODE_PROCESS = 
"[REMOVE_DATANODE_PROCESS]";
+
+  private static final int MAX_RETRY_NUM = 5;
 
   private static final int SLEEP_MILLIS = 5000;
 
@@ -172,7 +174,7 @@ public class RegionMigrateService implements IService {
     @Override
     public void start() {
       if (this.pool != null) {
-        poolLogger.info("Data Node region migrate pool start");
+        poolLogger.info("DataNode region migrate pool start");
       }
     }
 
@@ -266,7 +268,7 @@ public class RegionMigrateService implements IService {
       TEndPoint newPeerNode = getConsensusEndPoint(selectedDataNode, regionId);
       taskLogger.info("Start to add peer {} for region {}", newPeerNode, 
tRegionId);
       boolean addPeerSucceed = true;
-      for (int i = 0; i < RETRY; i++) {
+      for (int i = 0; i < MAX_RETRY_NUM; i++) {
         try {
           if (!addPeerSucceed) {
             Thread.sleep(SLEEP_MILLIS);
@@ -278,7 +280,7 @@ public class RegionMigrateService implements IService {
         } catch (Throwable e) {
           addPeerSucceed = false;
           taskLogger.error(
-              "Add new peer {} for region {} error, retry times: {}", 
newPeerNode, regionId, i, e);
+              "{}, Add new peer {} for region {} error, retry times: {}", 
REMOVE_DATANODE_PROCESS, newPeerNode, regionId, i, e);
           status.setCode(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode());
           status.setMessage(
               String.format(
@@ -291,7 +293,7 @@ public class RegionMigrateService implements IService {
       }
       if (!addPeerSucceed || resp == null || !resp.isSuccess()) {
         taskLogger.error(
-            "Add new peer {} for region {} failed, resp: {}", newPeerNode, 
regionId, resp);
+            "{}, Add new peer {} for region {} failed, resp: {}", 
REMOVE_DATANODE_PROCESS, newPeerNode, regionId, resp);
         status.setCode(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode());
         status.setMessage(
             String.format(
@@ -300,7 +302,7 @@ public class RegionMigrateService implements IService {
         return status;
       }
 
-      taskLogger.info("Succeed to add peer {} for region {}", newPeerNode, 
regionId);
+      taskLogger.info("{}, Succeed to add peer {} for region {}", 
REMOVE_DATANODE_PROCESS, newPeerNode, regionId);
       status.setCode(TSStatusCode.SUCCESS_STATUS.getStatusCode());
       status.setMessage("add peer " + newPeerNode + " for region " + regionId 
+ " succeed");
       return status;
@@ -375,7 +377,7 @@ public class RegionMigrateService implements IService {
       taskLogger.info("Start to remove peer {} for region {}", oldPeerNode, 
regionId);
       ConsensusGenericResponse resp = null;
       boolean removePeerSucceed = true;
-      for (int i = 0; i < RETRY; i++) {
+      for (int i = 0; i < MAX_RETRY_NUM; i++) {
         try {
           if (!removePeerSucceed) {
             Thread.sleep(SLEEP_MILLIS);
@@ -387,7 +389,7 @@ public class RegionMigrateService implements IService {
         } catch (Throwable e) {
           removePeerSucceed = false;
           taskLogger.error(
-              "remove peer {} for region {} error, retry times: {}", 
oldPeerNode, regionId, i, e);
+              "Remove peer {} for region {} error, retry times: {}", 
oldPeerNode, regionId, i, e);
           status.setCode(TSStatusCode.REGION_MIGRATE_FAILED.getStatusCode());
           status.setMessage(
               "remove peer: "

Reply via email to