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

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

commit 7294e355410cb982e2be2ccda42824463e3181b2
Author: Beyyes <[email protected]>
AuthorDate: Thu Sep 8 15:00:52 2022 +0800

    perfect addNewRegionPeer, removeRegionPeer logic
---
 .../consensus/request/ConfigPhysicalPlan.java      |  4 --
 .../iotdb/confignode/manager/ConsensusManager.java |  3 +-
 .../iotdb/confignode/manager/ProcedureManager.java |  8 +--
 .../procedure/env/DataNodeRemoveHandler.java       | 67 +++++++++++++++-------
 .../procedure/impl/RegionMigrateProcedure.java     | 44 ++++++++++----
 .../iotdb/db/service/RegionMigrateService.java     | 26 +++++----
 .../impl/DataNodeInternalRPCServiceImpl.java       | 34 +++++------
 7 files changed, 114 insertions(+), 72 deletions(-)

diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java
index f2e4edf976..be411c8e36 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java
@@ -242,10 +242,6 @@ public abstract class ConfigPhysicalPlan implements 
IConsensusRequest {
           throw new IOException("unknown PhysicalPlan type: " + typeNum);
       }
       req.deserializeImpl(buffer);
-      LOGGER.info(
-          "invoking create method in ConfigPhysicalPlan, planType: {}, req: 
{}",
-          ConfigPhysicalPlanType.values()[typeNum],
-          req instanceof UpdateProcedurePlan ? ((UpdateProcedurePlan) 
req).getProcedure() : null);
       return req;
     }
 
diff --git 
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConsensusManager.java
 
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConsensusManager.java
index 1322a498c8..75ce264ba5 100644
--- 
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConsensusManager.java
+++ 
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConsensusManager.java
@@ -47,7 +47,6 @@ import java.util.ArrayList;
 import java.util.Collections;
 import java.util.List;
 import java.util.concurrent.TimeUnit;
-import java.util.stream.Collectors;
 
 /** ConsensusManager maintains consensus class, request will redirect to 
consensus layer */
 public class ConsensusManager {
@@ -107,7 +106,7 @@ public class ConsensusManager {
       createPeerForConsensusGroup(
           Collections.singletonList(
               new TConfigNodeLocation(
-                      seedConfigNodeId,
+                  seedConfigNodeId,
                   new TEndPoint(CONF.getInternalAddress(), 
CONF.getInternalPort()),
                   new TEndPoint(CONF.getInternalAddress(), 
CONF.getConsensusPort()))));
     }
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 6dfe496fad..a92d67cd30 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
@@ -260,7 +260,8 @@ public class ProcedureManager {
   }
 
   public void reportRegionMigrateResult(TRegionMigrateResultReportReq req) {
-    LOGGER.info("Receive DataNode region: {} migrate result: {}", 
req.getRegionId(), req);
+    LOGGER.info(
+        "Receive DataNode region migrate result,regionId: {},req: {}", 
req.getRegionId(), req);
 
     this.executor
         .getProcedures()
@@ -271,11 +272,6 @@ public class ProcedureManager {
                 RegionMigrateProcedure regionMigrateProcedure = 
(RegionMigrateProcedure) procedure;
                 if 
(regionMigrateProcedure.getConsensusGroupId().equals(req.getRegionId())) {
                   regionMigrateProcedure.notifyTheRegionMigrateFinished(req);
-                } else {
-                  LOGGER.warn(
-                      "DataNode report region: {} is not equals ConfigNode 
send region: {}",
-                      req.getRegionId(),
-                      regionMigrateProcedure.getConsensusGroupId());
                 }
               }
             });
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 51e0b3363a..cea2fdbdb7 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
@@ -87,7 +87,7 @@ public class DataNodeRemoveHandler {
   public TSStatus broadcastDisableDataNode(TDataNodeLocation disabledDataNode) 
{
     LOGGER.info(
         "DataNodeRemoveService start send disable the Data Node to cluster, 
{}",
-            getIdWithRpcEndpoint(disabledDataNode));
+        getIdWithRpcEndpoint(disabledDataNode));
     TSStatus status = new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
     List<TEndPoint> otherOnlineDataNodes =
         
configManager.getNodeManager().filterDataNodeThroughStatus(NodeStatus.Running).stream()
@@ -108,7 +108,7 @@ public class DataNodeRemoveHandler {
     }
     LOGGER.info(
         "DataNodeRemoveService finished send disable the Data Node to cluster, 
{}",
-            getIdWithRpcEndpoint(disabledDataNode));
+        getIdWithRpcEndpoint(disabledDataNode));
     status.setMessage("Succeed disable the Data Node from cluster");
     return status;
   }
@@ -156,16 +156,19 @@ public class DataNodeRemoveHandler {
         filterDataNodeWithOtherRegionReplica(regionId, destDataNode);
     if (!selectedDataNode.isPresent()) {
       LOGGER.warn(
-          "There are no other DataNodes could be selected to perform the add 
peer process, please check RegionGroup: {} by SQL: show regions",
+          "There are no other DataNodes could be selected to perform the add 
peer process, "
+              + "please check RegionGroup: {} by SQL: show regions",
           regionId);
       status = new TSStatus(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode());
       status.setMessage(
-          "There are no other DataNodes could be selected to perform the add 
peer process, please check by SQL: show regions");
+          "There are no other DataNodes could be selected to perform the add 
peer process, "
+              + "please check by SQL: show regions");
       return status;
     }
 
-    // Send addRegionPeer request to the selected DataNode
-    TMaintainPeerReq maintainPeerReq = new TMaintainPeerReq(regionId, 
selectedDataNode.get());
+    // Send addRegionPeer request to the selected DataNode,
+    // destDataNode is where the new RegionReplica is created
+    TMaintainPeerReq maintainPeerReq = new TMaintainPeerReq(regionId, 
destDataNode);
     status =
         SyncDataNodeClientPool.getInstance()
             .sendSyncRequestToDataNodeWithRetry(
@@ -192,33 +195,51 @@ public class DataNodeRemoveHandler {
   public TSStatus removeRegionPeer(TDataNodeLocation originalDataNode, 
TConsensusGroupId regionId) {
     TSStatus status;
 
+    TDataNodeLocation rpcClientDataNode = null;
+
     // Here we pick the DataNode who contains one of the RegionReplica of the 
specified
     // ConsensusGroup except the origin one
     // in order to notify the new ConsensusGroup that the origin peer should 
secede now
     Optional<TDataNodeLocation> selectedDataNode =
         filterDataNodeWithOtherRegionReplica(regionId, originalDataNode);
-    if (!selectedDataNode.isPresent()) {
-      LOGGER.warn(
-          "There are no other DataNodes could be selected to perform the 
remove peer process, please check RegionGroup: {} by SQL: show regions",
-          regionId);
-      status = new TSStatus(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode());
-      status.setMessage(
-          "There are no other DataNodes could be selected to perform the 
remove peer process, please check by SQL: show regions");
-      return status;
+    if (selectedDataNode.isPresent()) {
+      rpcClientDataNode = selectedDataNode.get();
+    } else {
+      // if the originalDataNode to be removed is alive
+      // make the originalDataNode as rpcClientDataNode
+      List<TDataNodeConfiguration> aliveDataNodes =
+          
configManager.getNodeManager().filterDataNodeThroughStatus(NodeStatus.Running);
+      if (aliveDataNodes.stream().anyMatch(node -> 
node.getLocation().equals(originalDataNode))) {
+        LOGGER.info(
+            "Choose the originalDataNode to execute removeRegionPeer, node: 
{}", originalDataNode);
+        rpcClientDataNode = originalDataNode;
+      } else {
+        LOGGER.warn(
+            "There are no other DataNodes could be selected to perform the 
remove peer process, "
+                + "and originalDataNode is not alive"
+                + "please check RegionGroup: {} by SQL: show regions",
+            regionId);
+        status = new 
TSStatus(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode());
+        status.setMessage(
+            "There are no other DataNodes could be selected to perform the 
remove peer process, "
+                + "and originalDataNode is not alive"
+                + "please check by SQL: show regions");
+        return status;
+      }
     }
 
-    // Send addRegionPeer request to the selected DataNode
-    TMaintainPeerReq maintainPeerReq = new TMaintainPeerReq(regionId, 
selectedDataNode.get());
+    // Send removeRegionPeer request to the rpcClientDataNode
+    TMaintainPeerReq maintainPeerReq = new TMaintainPeerReq(regionId, 
rpcClientDataNode);
     status =
         SyncDataNodeClientPool.getInstance()
             .sendSyncRequestToDataNodeWithRetry(
-                selectedDataNode.get().getInternalEndPoint(),
+                rpcClientDataNode.getInternalEndPoint(),
                 maintainPeerReq,
                 DataNodeRequestType.REMOVE_REGION_PEER);
     LOGGER.info(
-        "Send region {} remove peer to {}, wait it finished",
+        "Send action removeRegionPeer, wait it finished, regionId: {}, 
dataNode: {}",
         regionId,
-        selectedDataNode.get().getInternalEndPoint());
+        rpcClientDataNode.getInternalEndPoint());
     return status;
   }
 
@@ -340,7 +361,8 @@ public class DataNodeRemoveHandler {
                 req,
                 DataNodeRequestType.CREATE_NEW_REGION_PEER);
 
-    LOGGER.info("Send action createNewRegionPeer, regionId: {}, dataNode: {}", 
regionId, destDataNode);
+    LOGGER.info(
+        "Send action createNewRegionPeer, regionId: {}, dataNode: {}", 
regionId, destDataNode);
     if (isFailed(status)) {
       LOGGER.error(
           "Send action createNewRegionPeer, regionId: {}, dataNode: {}, 
result: {}",
@@ -495,7 +517,8 @@ public class DataNodeRemoveHandler {
   }
 
   private String getIdWithRpcEndpoint(TDataNodeLocation location) {
-    return String.format("dataNodeId: %s, clientRpcEndPoint: %s",
-            location.getDataNodeId(), location.getClientRpcEndPoint());
+    return String.format(
+        "dataNodeId: %s, clientRpcEndPoint: %s",
+        location.getDataNodeId(), location.getClientRpcEndPoint());
   }
 }
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 e00a4753b8..14dba0a639 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
@@ -40,11 +40,13 @@ import java.io.DataOutputStream;
 import java.io.IOException;
 import java.nio.ByteBuffer;
 
+import static org.apache.iotdb.rpc.TSStatusCode.SUCCESS_STATUS;
+
 /** region migrate procedure */
 public class RegionMigrateProcedure
     extends StateMachineProcedure<ConfigNodeProcedureEnv, 
RegionTransitionState> {
   private static final Logger LOG = 
LoggerFactory.getLogger(RegionMigrateProcedure.class);
-  private static final int retryThreshold = 5;
+  private static final int RETRY_THRESHOLD = 5;
 
   /** Wait region migrate finished */
   private final Object regionMigrateLock = new Object();
@@ -55,6 +57,10 @@ public class RegionMigrateProcedure
 
   private TDataNodeLocation destDataNode;
 
+  private boolean migrateSuccess = true;
+
+  private String migrateResult = "";
+
   public RegionMigrateProcedure() {
     super();
   }
@@ -86,7 +92,7 @@ public class RegionMigrateProcedure
           break;
         case ADD_REGION_PEER:
           tsStatus = 
env.getDataNodeRemoveHandler().addRegionPeer(destDataNode, consensusGroupId);
-          if (tsStatus.getCode() == 
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+          if (tsStatus.getCode() == SUCCESS_STATUS.getStatusCode()) {
             waitForOneMigrationStepFinished(consensusGroupId);
             LOG.info("Wait for region add peer finished, regionId: {}", 
consensusGroupId);
           } else {
@@ -101,7 +107,7 @@ public class RegionMigrateProcedure
         case REMOVE_REGION_PEER:
           tsStatus =
               
env.getDataNodeRemoveHandler().removeRegionPeer(originalDataNode, 
consensusGroupId);
-          if (tsStatus.getCode() == 
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+          if (tsStatus.getCode() == SUCCESS_STATUS.getStatusCode()) {
             waitForOneMigrationStepFinished(consensusGroupId);
             LOG.info("Wait for region {} remove peer finished", 
consensusGroupId);
           } else {
@@ -113,7 +119,7 @@ public class RegionMigrateProcedure
           tsStatus =
               env.getDataNodeRemoveHandler()
                   .deleteOldRegionPeer(originalDataNode, consensusGroupId);
-          if (tsStatus.getCode() == 
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+          if (tsStatus.getCode() == SUCCESS_STATUS.getStatusCode()) {
             waitForOneMigrationStepFinished(consensusGroupId);
             LOG.info("Wait for region {}  remove consensus group finished", 
consensusGroupId);
           }
@@ -131,9 +137,14 @@ public class RegionMigrateProcedure
         setFailure(new ProcedureException("Region migrate failed " + state));
       } else {
         LOG.error(
-            "Retrievable error trying to region migrate {}, state {}", 
originalDataNode, state, e);
-        if (getCycles() > retryThreshold) {
-          setFailure(new ProcedureException("State stuck at " + state));
+            "Failed state is not support rollback, retrievable error trying to 
region migrate {}, state {}",
+            originalDataNode,
+            state,
+            e);
+        if (getCycles() > RETRY_THRESHOLD) {
+          setFailure(
+              new ProcedureException(
+                  "Procedure retried failed exceed 5 times, state stuck at " + 
state));
         }
       }
     }
@@ -170,7 +181,7 @@ public class RegionMigrateProcedure
   protected void releaseLock(ConfigNodeProcedureEnv configNodeProcedureEnv) {
     configNodeProcedureEnv.getSchedulerLock().lock();
     try {
-      LOG.info("{} release lock.", getProcId());
+      LOG.info("procedureId {} release lock.", getProcId());
       if (configNodeProcedureEnv.getRegionMigrateLock().releaseLock(this)) {
         configNodeProcedureEnv
             .getRegionMigrateLock()
@@ -230,12 +241,18 @@ public class RegionMigrateProcedure
     return false;
   }
 
-  public TSStatus waitForOneMigrationStepFinished(TConsensusGroupId 
consensusGroupId) {
-    TSStatus status = new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+  public TSStatus waitForOneMigrationStepFinished(TConsensusGroupId 
consensusGroupId)
+      throws Exception {
+    TSStatus status = new TSStatus(SUCCESS_STATUS.getStatusCode());
     synchronized (regionMigrateLock) {
       try {
         // TODO set timeOut?
         regionMigrateLock.wait();
+
+        if (!migrateSuccess) {
+          throw new RuntimeException(
+              String.format("Region migrate failed, regionId: %s", 
consensusGroupId));
+        }
       } catch (InterruptedException e) {
         LOG.error("region migrate {} interrupt", consensusGroupId, e);
         Thread.currentThread().interrupt();
@@ -250,6 +267,13 @@ public class RegionMigrateProcedure
   public void notifyTheRegionMigrateFinished(TRegionMigrateResultReportReq 
req) {
     // TODO the req is used in roll back
     synchronized (regionMigrateLock) {
+      TSStatus migrateStatus = req.getMigrateResult();
+      // migrate failed
+      if (migrateStatus.getCode() != SUCCESS_STATUS.getStatusCode()) {
+        LOG.info("Region migrate failed, {}", req);
+        migrateSuccess = false;
+        migrateResult = migrateStatus.getMessage();
+      }
       regionMigrateLock.notify();
     }
     LOG.info(
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 809da5b01e..bbc1e65930 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
@@ -67,7 +67,7 @@ public class RegionMigrateService implements IService {
   /**
    * submit AddRegionPeerTask
    *
-   * @param req TMigrateRegionReq
+   * @param req TMaintainPeerReq
    * @return if the submit task succeed
    */
   public synchronized boolean submitAddRegionPeerTask(TMaintainPeerReq req) {
@@ -77,7 +77,7 @@ public class RegionMigrateService implements IService {
       regionMigratePool.submit(new AddRegionPeerTask(req.getRegionId(), 
req.getDestNode()));
     } catch (Exception e) {
       LOGGER.error(
-          "Submit add region peer task error for Region: {} on DataNode: {}.",
+          "Submit addRegionPeer task error for Region: {} on DataNode: {}.",
           req.getRegionId(),
           req.getDestNode().getInternalEndPoint().getIp(),
           e);
@@ -99,7 +99,7 @@ public class RegionMigrateService implements IService {
       regionMigratePool.submit(new RemoveRegionPeerTask(req.getRegionId(), 
req.getDestNode()));
     } catch (Exception e) {
       LOGGER.error(
-          "Submit remove region peer task error for Region: {} on DataNode: 
{}.",
+          "Submit removeRegionPeer task error for Region: {} on DataNode: {}.",
           req.getRegionId(),
           req.getDestNode().getInternalEndPoint().getIp(),
           e);
@@ -114,7 +114,7 @@ public class RegionMigrateService implements IService {
    * @param req TMigrateRegionReq
    * @return submit task succeed?
    */
-  public synchronized boolean 
submitRemoveRegionConsensusGroupTask(TMaintainPeerReq req) {
+  public synchronized boolean submitDeleteOldRegionPeerTask(TMaintainPeerReq 
req) {
 
     boolean submitSucceed = true;
     try {
@@ -122,7 +122,7 @@ public class RegionMigrateService implements IService {
           new RemoveRegionConsensusGroupTask(req.getRegionId(), 
req.getDestNode()));
     } catch (Exception e) {
       LOGGER.error(
-          "Submit remove region peer task error for Region: {} on DataNode: 
{}.",
+          "Submit deleteOldRegionPeerTask error for Region: {} on DataNode: 
{}.",
           req.getRegionId(),
           req.getDestNode().getInternalEndPoint().getIp(),
           e);
@@ -279,7 +279,9 @@ public class RegionMigrateService implements IService {
           taskLogger.error(
               "Add new peer {} for region {} error, retry times: {}", 
newPeerNode, regionId, i, e);
           status.setCode(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode());
-          status.setMessage(String.format("Add peer for region error, peerId: 
%s, regionId: %s, errorMessage: %s",
+          status.setMessage(
+              String.format(
+                  "Add peer for region error, peerId: %s, regionId: %s, 
errorMessage: %s",
                   newPeerNode, regionId, e.getMessage()));
         }
         if (addPeerSucceed && resp != null && resp.isSuccess()) {
@@ -288,9 +290,11 @@ 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: {}", newPeerNode, 
regionId, resp);
         status.setCode(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode());
-        status.setMessage(String.format("Add peer for region error, peerId: 
%s, regionId: %s, resp: %s",
+        status.setMessage(
+            String.format(
+                "Add peer for region error, peerId: %s, regionId: %s, resp: 
%s",
                 newPeerNode, regionId, resp));
         return status;
       }
@@ -331,10 +335,10 @@ public class RegionMigrateService implements IService {
   private static class RemoveRegionPeerTask implements Runnable {
     private static final Logger taskLogger = 
LoggerFactory.getLogger(RemoveRegionPeerTask.class);
 
-    // The RegionGroup that shall perform the add peer process
+    // The RegionGroup that shall perform the remove peer process
     private final TConsensusGroupId tRegionId;
 
-    // The DataNode that selected to perform the add peer process
+    // The DataNode that selected to perform the remove peer process
     private final TDataNodeLocation selectedDataNode;
 
     public RemoveRegionPeerTask(TConsensusGroupId tRegionId, TDataNodeLocation 
selectedDataNode) {
@@ -367,7 +371,7 @@ public class RegionMigrateService implements IService {
       ConsensusGroupId regionId = 
ConsensusGroupId.Factory.createFromTConsensusGroupId(tRegionId);
       TSStatus status = new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
       TEndPoint oldPeerNode = getConsensusEndPoint(selectedDataNode, regionId);
-      taskLogger.info("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++) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
index 3394d5e2d9..06ca5aa730 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
@@ -644,14 +644,13 @@ public class DataNodeInternalRPCServiceImpl implements 
IDataNodeRPCService.Iface
         
ConsensusGroupId.Factory.createFromTConsensusGroupId(req.getRegionId());
     List<Peer> peers =
         req.getRegionLocations().stream()
-            .map(n -> getConsensusEndPoint(n, regionId))
-            .map(node -> new Peer(regionId, node))
+            .map(location -> new Peer(regionId, getConsensusEndPoint(location, 
regionId)))
             .collect(Collectors.toList());
     TSStatus status = createNewRegion(regionId, req.getStorageGroup(), 
req.getTtl());
     if (!isSucceed(status)) {
       return status;
     }
-    return addConsensusGroup(regionId, peers);
+    return createNewRegionPeer(regionId, peers);
   }
 
   @Override
@@ -662,13 +661,13 @@ public class DataNodeInternalRPCServiceImpl implements 
IDataNodeRPCService.Iface
     TSStatus status = new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
     if (submitSucceed) {
       LOGGER.info(
-          "Successfully submit a add region peer task for region: {} on 
DataNode: {}",
+          "Successfully submit addRegionPeer task for region: {} on DataNode: 
{}",
           regionId,
           selectedDataNodeIP);
       return status;
     }
     status.setCode(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode());
-    status.setMessage("submit add region peer task failed, region: " + 
regionId);
+    status.setMessage("Submit addRegionPeer task failed, region: " + regionId);
     return status;
   }
 
@@ -680,13 +679,13 @@ public class DataNodeInternalRPCServiceImpl implements 
IDataNodeRPCService.Iface
     TSStatus status = new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
     if (submitSucceed) {
       LOGGER.info(
-          "Successfully to submit a remove region peer task for region: {} on 
DataNode: {}",
+          "Successfully submit removeRegionPeer task for region: {} on 
DataNode: {}",
           regionId,
           selectedDataNodeIP);
       return status;
     }
     status.setCode(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode());
-    status.setMessage("submit add region peer task failed, region: " + 
regionId);
+    status.setMessage("Submit removeRegionPeer task failed, region: " + 
regionId);
     return status;
   }
 
@@ -694,19 +693,17 @@ public class DataNodeInternalRPCServiceImpl implements 
IDataNodeRPCService.Iface
   public TSStatus deleteOldRegionPeer(TMaintainPeerReq req) throws TException {
     TConsensusGroupId regionId = req.getRegionId();
     String selectedDataNodeIP = 
req.getDestNode().getInternalEndPoint().getIp();
-    boolean submitSucceed =
-        
RegionMigrateService.getInstance().submitRemoveRegionConsensusGroupTask(req);
+    boolean submitSucceed = 
RegionMigrateService.getInstance().submitDeleteOldRegionPeerTask(req);
     TSStatus status = new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
     if (submitSucceed) {
       LOGGER.info(
-          "Successfully to submit a remove region consensus group task for 
region: {} on DataNode: {}",
+          "Successfully submit deleteOldRegionPeer task for region: {} on 
DataNode: {}",
           regionId,
           selectedDataNodeIP);
       return status;
     }
     status.setCode(TSStatusCode.MIGRATE_REGION_ERROR.getStatusCode());
-    status.setMessage(
-        "submit region remove region consensus group task failed, region: " + 
regionId);
+    status.setMessage("Submit deleteOldRegionPeer task failed, region: " + 
regionId);
     return status;
   }
 
@@ -782,8 +779,8 @@ public class DataNodeInternalRPCServiceImpl implements 
IDataNodeRPCService.Iface
     return status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode();
   }
 
-  private TSStatus addConsensusGroup(ConsensusGroupId regionId, List<Peer> 
peers) {
-    LOGGER.info("Start to add consensus group {} to region {}", peers, 
regionId);
+  private TSStatus createNewRegionPeer(ConsensusGroupId regionId, List<Peer> 
peers) {
+    LOGGER.info("Start to createNewRegionPeer {} to region {}", peers, 
regionId);
     TSStatus status = new 
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
     ConsensusGenericResponse resp;
     if (regionId instanceof DataRegionId) {
@@ -793,13 +790,16 @@ public class DataNodeInternalRPCServiceImpl implements 
IDataNodeRPCService.Iface
     }
     if (!resp.isSuccess()) {
       LOGGER.error(
-          "add peers {} to region {} consensus group error", peers, regionId, 
resp.getException());
+          "CreateNewRegionPeer error, peers: {}, regionId: {}, errorMessage",
+          peers,
+          regionId,
+          resp.getException());
       status.setCode(TSStatusCode.REGION_MIGRATE_FAILED.getStatusCode());
       status.setMessage(resp.getException().getMessage());
       return status;
     }
-    LOGGER.info("succeed to add peers {} to region {} consensus group", peers, 
regionId);
-    status.setMessage("add peers to region consensus group " + regionId + 
"succeed");
+    LOGGER.info("Succeed to createNewRegionPeer {} for region {}", peers, 
regionId);
+    status.setMessage("createNewRegionPeer succeed, regionId: " + regionId);
     return status;
   }
 

Reply via email to