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; }
