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: "