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

wangchao316 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 4c8edc5661 [IOTDB-3985] Retry removePeer for region bug (#6829)
4c8edc5661 is described below

commit 4c8edc56611991d75b7a995f12da91fbc4e76d95
Author: QiangShaowei <[email protected]>
AuthorDate: Fri Jul 29 11:55:13 2022 +0800

    [IOTDB-3985] Retry removePeer for region bug (#6829)
    
    [IOTDB-3985] Retry removePeer for region bug (#6829)
---
 .../iotdb/confignode/manager/ProcedureManager.java |  2 +-
 .../iotdb/confignode/persistence/NodeInfo.java     |  3 +++
 .../procedure/env/DataNodeRemoveHandler.java       | 10 +++++----
 .../procedure/impl/RegionMigrateProcedure.java     | 12 +++++++++--
 .../iotdb/db/service/RegionMigrateService.java     | 25 +++++++++++-----------
 5 files changed, 32 insertions(+), 20 deletions(-)

diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
index 564a921f7e..afcf242b32 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
@@ -246,7 +246,7 @@ public class ProcedureManager {
               if (procedure instanceof RegionMigrateProcedure) {
                 RegionMigrateProcedure regionMigrateProcedure = 
(RegionMigrateProcedure) procedure;
                 if 
(regionMigrateProcedure.getConsensusGroupId().equals(req.getRegionId())) {
-                  regionMigrateProcedure.notifyTheRegionMigrateFinished();
+                  regionMigrateProcedure.notifyTheRegionMigrateFinished(req);
                 }
               }
             });
diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/NodeInfo.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/NodeInfo.java
index 5d71b5da60..5830c6127a 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/NodeInfo.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/NodeInfo.java
@@ -195,16 +195,19 @@ public class NodeInfo implements SnapshotProcessor {
    * @return TSStatus
    */
   public TSStatus removeDataNode(RemoveDataNodePlan req) {
+    LOGGER.info("there are {} data node in cluster before remove some", 
registeredDataNodes.size());
     try {
       dataNodeInfoReadWriteLock.writeLock().lock();
       req.getDataNodeLocations()
           .forEach(
               removeDataNodes -> {
                 registeredDataNodes.remove(removeDataNodes.getDataNodeId());
+                LOGGER.info("removed the datanode {} from cluster", 
removeDataNodes);
               });
     } finally {
       dataNodeInfoReadWriteLock.writeLock().unlock();
     }
+    LOGGER.info("there are {} data node in cluster after remove some", 
registeredDataNodes.size());
     return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
   }
 
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 4663ed47d3..aa5a0b39e7 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
@@ -175,8 +175,10 @@ public class DataNodeRemoveHandler {
                 originalDataNode.getInternalEndPoint(),
                 migrateRegionReq,
                 DataNodeRequestType.MIGRATE_REGION);
-    LOGGER.debug(
-        "send region {} migrate action to {}, wait it finished", regionId, 
originalDataNode);
+    LOGGER.info(
+        "send region {} migrate action to {}, wait it finished",
+        regionId,
+        originalDataNode.getInternalEndPoint());
     return status;
   }
 
@@ -191,7 +193,7 @@ public class DataNodeRemoveHandler {
       TConsensusGroupId regionId,
       TDataNodeLocation originalDataNode,
       TDataNodeLocation destDataNode) {
-    LOGGER.debug(
+    LOGGER.info(
         "start to update region {} location from {} to {} when it migrate 
succeed",
         regionId,
         originalDataNode.getInternalEndPoint().getIp(),
@@ -258,7 +260,7 @@ public class DataNodeRemoveHandler {
             .addToRegionConsensusGroup(
                 // TODO replace with real ttl
                 regionReplicaNodes, regionId, destDataNode, storageGroup, 
Long.MAX_VALUE);
-    LOGGER.debug("send add region {} consensus group to {}", regionId, 
destDataNode);
+    LOGGER.info("send add region {} consensus group to {}", regionId, 
destDataNode);
     if (isFailed(status)) {
       LOGGER.error(
           "add new node {} to region {} consensus group failed,  result: {}",
diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/RegionMigrateProcedure.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/RegionMigrateProcedure.java
index 9a8c13505f..a2641da308 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/RegionMigrateProcedure.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/RegionMigrateProcedure.java
@@ -30,6 +30,7 @@ import 
org.apache.iotdb.confignode.procedure.exception.ProcedureException;
 import org.apache.iotdb.confignode.procedure.state.ProcedureLockState;
 import org.apache.iotdb.confignode.procedure.state.RegionTransitionState;
 import org.apache.iotdb.confignode.procedure.store.ProcedureFactory;
+import org.apache.iotdb.confignode.rpc.thrift.TRegionMigrateResultReportReq;
 import org.apache.iotdb.rpc.TSStatusCode;
 
 import org.slf4j.Logger;
@@ -91,7 +92,7 @@ public class RegionMigrateProcedure
         case WAIT_FOR_REGION_MIGRATE_FINISHED:
           waitForTheRegionMigrateFinished(consensusGroupId);
           setNextState(RegionTransitionState.UPDATE_REGION_LOCATION_CACHE);
-          LOG.info("Wait for region migrate finished");
+          LOG.info("Wait for region {}  migrate finished", consensusGroupId);
           break;
         case UPDATE_REGION_LOCATION_CACHE:
           env.getDataNodeRemoveHandler()
@@ -218,7 +219,14 @@ public class RegionMigrateProcedure
     return status;
   }
 
-  public void notifyTheRegionMigrateFinished() {
+  /**
+   * DN report region migrate result to CN, and continue
+   *
+   * @param req
+   */
+  public void notifyTheRegionMigrateFinished(TRegionMigrateResultReportReq 
req) {
+    LOG.info("DataNode reported region {} migrate result:{} ", 
req.getRegionId(), req);
+    // TODO the req is used in roll back
     synchronized (regionMigrateLock) {
       regionMigrateLock.notify();
     }
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 d6c304c2d7..0aaae47d64 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
@@ -303,13 +303,11 @@ public class RegionMigrateService implements IService {
       TSStatus status;
       try (ConfigNodeClient client = new ConfigNodeClient()) {
         status = client.reportRegionMigrateResult(req);
-        if (taskLogger.isDebugEnabled()) {
-          taskLogger.debug(
-              "report region {} migrate result {} to Config node succeed, 
result: {}",
-              tRegionId,
-              req,
-              status);
-        }
+        taskLogger.info(
+            "report region {} migrate result {} to Config node succeed, 
result: {}",
+            tRegionId,
+            req,
+            status);
       }
     }
 
@@ -318,7 +316,7 @@ public class RegionMigrateService implements IService {
       TSStatus status = new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
       ConsensusGenericResponse resp = null;
       TEndPoint newPeerNode = getConsensusEndPoint(toNode, regionId);
-      taskLogger.debug("start to add peer {} for region {}", newPeerNode, 
tRegionId);
+      taskLogger.info("start to add peer {} for region {}", newPeerNode, 
tRegionId);
       boolean addPeerSucceed = true;
       for (int i = 0; i < RETRY; i++) {
         try {
@@ -352,7 +350,7 @@ public class RegionMigrateService implements IService {
         return status;
       }
 
-      taskLogger.debug("succeed to add peer {} for region {}", newPeerNode, 
regionId);
+      taskLogger.info("succeed to add peer {} for region {}", newPeerNode, 
regionId);
       status.setCode(TSStatusCode.SUCCESS_STATUS.getStatusCode());
       status.setMessage("add peer " + newPeerNode + " for region " + regionId 
+ " succeed");
       return status;
@@ -382,7 +380,7 @@ public class RegionMigrateService implements IService {
       ConsensusGroupId regionId = 
ConsensusGroupId.Factory.createFromTConsensusGroupId(tRegionId);
       TSStatus status = new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
       TEndPoint oldPeerNode = getConsensusEndPoint(fromNode, regionId);
-      taskLogger.debug("start to remove peer {} for region {}", oldPeerNode, 
regionId);
+      taskLogger.info("start to remove peer {} for region {}", oldPeerNode, 
regionId);
       ConsensusGenericResponse resp = null;
       boolean removePeerSucceed = true;
       for (int i = 0; i < RETRY; i++) {
@@ -391,6 +389,7 @@ public class RegionMigrateService implements IService {
             Thread.sleep(SLEEP_MILLIS);
           }
           resp = removeRegionPeer(regionId, new Peer(regionId, oldPeerNode));
+          removePeerSucceed = true;
         } catch (Throwable e) {
           removePeerSucceed = false;
           taskLogger.error(
@@ -417,7 +416,7 @@ public class RegionMigrateService implements IService {
         return status;
       }
 
-      taskLogger.debug("succeed to remove peer {} for region {}", oldPeerNode, 
regionId);
+      taskLogger.info("succeed to remove peer {} for region {}", oldPeerNode, 
regionId);
       status.setCode(TSStatusCode.SUCCESS_STATUS.getStatusCode());
       status.setMessage("remove peer " + oldPeerNode + " for region " + 
regionId + " succeed");
       return status;
@@ -425,7 +424,7 @@ public class RegionMigrateService implements IService {
 
     private TSStatus removeConsensusGroup() {
       ConsensusGroupId regionId = 
ConsensusGroupId.Factory.createFromTConsensusGroupId(tRegionId);
-      taskLogger.debug("start to remove region {} consensus group", regionId);
+      taskLogger.info("start to remove region {} consensus group", regionId);
       TSStatus status = new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
       ConsensusGenericResponse resp;
       try {
@@ -454,7 +453,7 @@ public class RegionMigrateService implements IService {
                 + resp.getException().getMessage());
         return status;
       }
-      taskLogger.debug("succeed to remove region {} consensus group", 
regionId);
+      taskLogger.info("succeed to remove region {} consensus group", regionId);
       status.setMessage("remove region consensus group " + regionId + 
"succeed");
       return status;
     }

Reply via email to