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