This is an automated email from the ASF dual-hosted git repository.
lta pushed a commit to branch cluster_scalability
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/cluster_scalability by this
push:
new 5576d78 This commit fix following issues: 1. fix a issue of
watiFollowerToSync too long 2. change wait time of syncLeader 3. fix
nullPointer of result in PullSnapshotTask 4. Remove local data of
workingTsFileProcessor when remove partition
5576d78 is described below
commit 5576d785d79f2e115cda9a0d47d640d1cb7dee1e
Author: lta <[email protected]>
AuthorDate: Thu Mar 4 11:24:14 2021 +0800
This commit fix following issues:
1. fix a issue of watiFollowerToSync too long
2. change wait time of syncLeader
3. fix nullPointer of result in PullSnapshotTask
4. Remove local data of workingTsFileProcessor when remove partition
---
.../cluster/log/snapshot/PullSnapshotTask.java | 17 ++++--
.../iotdb/cluster/server/DataClusterServer.java | 17 ++++--
.../cluster/server/PullSnapshotHintService.java | 6 ++
.../apache/iotdb/cluster/server/RaftServer.java | 2 +-
.../cluster/server/heartbeat/HeartbeatThread.java | 3 +-
.../cluster/server/member/DataGroupMember.java | 67 ++++++++++++++--------
.../iotdb/cluster/server/member/RaftMember.java | 41 +++----------
.../engine/storagegroup/StorageGroupProcessor.java | 1 +
.../org/apache/iotdb/db/utils/CommonUtils.java | 4 +-
9 files changed, 86 insertions(+), 72 deletions(-)
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/log/snapshot/PullSnapshotTask.java
b/cluster/src/main/java/org/apache/iotdb/cluster/log/snapshot/PullSnapshotTask.java
index 1dc3247..fc9b968 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/log/snapshot/PullSnapshotTask.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/log/snapshot/PullSnapshotTask.java
@@ -89,8 +89,9 @@ public class PullSnapshotTask<T extends Snapshot> implements
Callable<Void> {
throws InterruptedException, TException {
Node node = descriptor.getPreviousHolders().get(nodeIndex);
if (logger.isDebugEnabled()) {
- logger.debug("Pulling {} snapshots from {} of {}",
descriptor.getSlots().size(), node,
- descriptor.getPreviousHolders().getHeader());
+ logger.debug("Pulling slot {} and other {} snapshots from {} of {} for
{}",
+ descriptor.getSlots().get(0), descriptor.getSlots().size() - 1, node,
+ descriptor.getPreviousHolders().getHeader(), newMember.getName());
}
Map<Integer, T> result = pullSnapshot(node);
@@ -115,9 +116,11 @@ public class PullSnapshotTask<T extends Snapshot>
implements Callable<Void> {
descriptor.getPreviousHolders().get(nodeIndex));
}
try {
- Snapshot snapshot = result.values().iterator().next();
- SnapshotInstaller installer = snapshot.getDefaultInstaller(newMember);
- installer.install(result);
+ if (result.size() > 0) {
+ Snapshot snapshot = result.values().iterator().next();
+ SnapshotInstaller installer =
snapshot.getDefaultInstaller(newMember);
+ installer.install(result);
+ }
// inform the previous holders that one member has successfully pulled
snapshot
newMember.registerPullSnapshotHint(descriptor);
return true;
@@ -179,6 +182,10 @@ public class PullSnapshotTask<T extends Snapshot>
implements Callable<Void> {
nodeIndex = (nodeIndex + 1) % descriptor.getPreviousHolders().size();
finished = pullSnapshot(nodeIndex);
if (!finished) {
+ if (logger.isDebugEnabled()) {
+ logger.debug("Cannot pull slot {} from {}, retry",
descriptor.getSlots(),
+ descriptor.getPreviousHolders().get(nodeIndex));
+ }
Thread
.sleep(
ClusterDescriptor.getInstance().getConfig().getPullSnapshotRetryIntervalMs());
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/DataClusterServer.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/DataClusterServer.java
index 2a0dc2a..e95df88 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/DataClusterServer.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/DataClusterServer.java
@@ -554,8 +554,8 @@ public class DataClusterServer extends RaftServer
implements TSDataService.Async
boolean shouldLeave = dataGroupMember.addNode(node, result);
if (shouldLeave) {
logger.info("This node does not belong to {} any more",
dataGroupMember.getAllNodes());
+ removeMember(entry.getKey(), entry.getValue(), false);
entryIterator.remove();
- removeMember(entry.getKey(), entry.getValue());
}
}
@@ -576,11 +576,16 @@ public class DataClusterServer extends RaftServer
implements TSDataService.Async
* @param header
* @param dataGroupMember
*/
- private void removeMember(RaftNode header, DataGroupMember dataGroupMember) {
-
dataGroupMember.getStopStatus().setSyncSuccess(dataGroupMember.syncLeader());
+ private void removeMember(RaftNode header, DataGroupMember dataGroupMember,
boolean waitFollowersToSync) {
+ if (dataGroupMember.syncLeader()) {
+ dataGroupMember.setHasSyncedLeaderBeforeRemoved(true);
+ }
dataGroupMember.setReadOnly();
- dataGroupMember.waitFollowersToSync();
- dataGroupMember.stop();
+ if (waitFollowersToSync && dataGroupMember.getCharacter() ==
NodeCharacter.LEADER) {
+ dataGroupMember.getAppendLogThreadPool().submit(() ->
dataGroupMember.waitFollowersToSync());
+ } else {
+ dataGroupMember.stop();
+ }
stoppedMemberManager.put(header, dataGroupMember);
logger.info("Data group member has removed, header {}, group is {}.",
header,
dataGroupMember.getAllNodes());
@@ -658,7 +663,7 @@ public class DataClusterServer extends RaftServer
implements TSDataService.Async
DataGroupMember dataGroupMember = entry.getValue();
if (dataGroupMember.getHeader().equals(node) || node.equals(thisNode))
{
entryIterator.remove();
- removeMember(entry.getKey(), dataGroupMember);
+ removeMember(entry.getKey(), dataGroupMember,
dataGroupMember.getHeader().equals(node));
} else {
// the group should be updated and pull new slots from the removed
node
dataGroupMember.removeNode(node, removalResult);
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/PullSnapshotHintService.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/PullSnapshotHintService.java
index 0cc1452..d3d611f 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/PullSnapshotHintService.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/PullSnapshotHintService.java
@@ -83,6 +83,12 @@ public class PullSnapshotHintService {
for (Iterator<Node> iter = hint.receivers.iterator(); iter.hasNext(); ) {
Node receiver = iter.next();
try {
+ if (logger.isDebugEnabled()) {
+ logger.debug(
+ "{}: start to send hint to target group {}, receiver {}, slot
is {} and other {}",
+ member.getName(), hint.receivers, receiver, hint.slots.get(0),
+ hint.slots.size() - 1);
+ }
boolean result = sendHint(receiver, hint);
if (result) {
iter.remove();
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/RaftServer.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/RaftServer.java
index 2e925c1..4a061ca 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/server/RaftServer.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/server/RaftServer.java
@@ -61,7 +61,7 @@ public abstract class RaftServer implements
RaftService.AsyncIface, RaftService.
ClusterDescriptor.getInstance().getConfig().getReadOperationTimeoutMS();
private static int writeOperationTimeoutMS =
ClusterDescriptor.getInstance().getConfig().getWriteOperationTimeoutMS();
- private static int syncLeaderMaxWaitMs = 20 * 1000;
+ private static int syncLeaderMaxWaitMs = 10 * 1000;
private static long heartBeatIntervalMs = 1000L;
ClusterConfig config = ClusterDescriptor.getInstance().getConfig();
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
index 99bf3f5..82eb020 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java
@@ -147,7 +147,8 @@ public class HeartbeatThread implements Runnable {
@SuppressWarnings("java:S2445")
private void sendHeartbeats(Collection<Node> nodes) {
if (logger.isDebugEnabled()) {
- logger.debug("{}: Send heartbeat to {} followers", memberName,
nodes.size() - 1);
+ logger.debug("{}: Send heartbeat to {} followers, commit log index =
{}", memberName,
+ nodes.size() - 1, request.getCommitLogIndex());
}
synchronized (nodes) {
// avoid concurrent modification
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
index a9f39be..fbf2dd4 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
@@ -307,15 +307,16 @@ public class DataGroupMember extends RaftMember {
Set<Integer> lostSlots = ((SlotNodeAdditionResult) result).getLostSlots()
.getOrDefault(new RaftNode(getHeader(), getRaftGroupId()),
Collections.emptySet());
for (Integer lostSlot : lostSlots) {
- slotManager.setToSending(lostSlot);
+ slotManager.setToSending(lostSlot, false);
}
+ slotManager.save();
synchronized (allNodes) {
if (allNodes.contains(node) && allNodes.size() >
config.getReplicationNum()) {
// remove the last node because the group size is fixed to replication
number
Node removedNode = allNodes.remove(allNodes.size() - 1);
peerMap.remove(removedNode);
- if (removedNode.equals(leader.get())) {
+ if (removedNode.equals(leader.get()) && !removedNode.equals(thisNode))
{
// if the leader is removed, also start an election immediately
synchronized (term) {
setCharacter(NodeCharacter.ELECTOR);
@@ -420,11 +421,11 @@ public class DataGroupMember extends RaftMember {
boolean canGetSnapshot;
/**
* There are two conditions that can get snapshot:
- * 1. The raft member is stopped and sync status is successful which means
it has synced leader successfully before stop.
+ * 1. The raft member is stopped and it has synced leader successfully
before stop.
* 2. The raft member is not stopped and syncing leader is successful.
*/
- if (stopStatus.stop) {
- canGetSnapshot = stopStatus.syncSuccess;
+ if (isHasSyncedLeaderBeforeRemoved()) {
+ canGetSnapshot = true;
} else {
canGetSnapshot = syncLeader();
}
@@ -523,6 +524,8 @@ public class DataGroupMember extends RaftMember {
private void pullFileSnapshot(PullSnapshotTaskDescriptor descriptor, File
snapshotSave) {
// If this node is the member of previous holder, it's unnecessary to pull
data again
if (descriptor.getPreviousHolders().contains(thisNode)) {
+ logger.info("{}: {} and other {} don't need to pull because there
already has such data locally", name,
+ descriptor.getSlots().get(0), descriptor.getSlots().size() - 1);
// inform the previous holders that one member has successfully pulled
snapshot directly
registerPullSnapshotHint(descriptor);
return;
@@ -785,7 +788,7 @@ public class DataGroupMember extends RaftMember {
syncLeader();
synchronized (allNodes) {
- if (allNodes.contains(removedNode) && allNodes.size() >
config.getReplicationNum()) {
+ if (allNodes.contains(removedNode)) {
// update the group if the deleted node was in it
allNodes.remove(removedNode);
peerMap.remove(removedNode);
@@ -810,27 +813,36 @@ public class DataGroupMember extends RaftMember {
}
}
+ /**
+ * When the header of a partition group is removed, it needs to wait all
followers to sync data because
+ * there has no new leader.
+ */
public void waitFollowersToSync() {
- if (character != NodeCharacter.LEADER) {
- return;
- }
- for (Map.Entry<Node, Peer> entry: peerMap.entrySet()) {
- Node node = entry.getKey();
- if (node.equals(thisNode)) {
- continue;
- }
- Peer peer = entry.getValue();
- while (peer.getMatchIndex() < logManager.getCommitLogIndex()) {
- try {
- Thread.sleep(10);
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- logger.warn("{}: Unexpected interruption when waiting follower {} to
sync, raft id is {}",
- name, node, getRaftGroupId());
+ try {
+ for (Map.Entry<Node, Peer> entry : peerMap.entrySet()) {
+ Node node = entry.getKey();
+ if (node.equals(thisNode)) {
+ continue;
+ }
+ Peer peer = entry.getValue();
+ while (peer.getMatchIndex() < logManager.getCommitLogIndex()) {
+ try {
+ Thread.sleep(10);
+ if (character != NodeCharacter.LEADER) {
+ return;
+ }
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ logger
+ .warn("{}: Unexpected interruption when waiting follower {} to
sync, raft id is {}",
+ name, node, getRaftGroupId());
+ }
}
+ logger.info("{}: Follower {} has synced with leader, raft id is {}",
name, node,
+ getRaftGroupId());
}
- logger.info("{}: Follower {} has synced with leader, raft id is {}",
name, node,
- getRaftGroupId());
+ } finally {
+ stop();
}
}
@@ -872,6 +884,13 @@ public class DataGroupMember extends RaftMember {
}
public boolean onSnapshotInstalled(List<Integer> slots) {
+ if (!isHasSyncedLeaderBeforeRemoved()) {
+
getMetaGroupMember().waitUtil(getMetaGroupMember().getPartitionTable().getLastMetaLogIndex());
+ }
+ if (logger.isDebugEnabled()) {
+ logger.debug("{} received one replication snapshot installed of slot {}
and other {} slots",
+ name, slots.get(0), slots.size() - 1);
+ }
List<Integer> removableSlots = new ArrayList<>();
for (Integer slot : slots) {
int sentReplicaNum = slotManager.sentOneReplication(slot, false);
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
index 83eec55..eb2e6cf 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
@@ -85,7 +85,6 @@ import org.apache.iotdb.cluster.utils.ClientUtils;
import org.apache.iotdb.cluster.utils.IOUtils;
import org.apache.iotdb.cluster.utils.PlanSerializer;
import org.apache.iotdb.cluster.utils.StatusUtils;
-import org.apache.iotdb.cluster.utils.nodetool.function.Status;
import org.apache.iotdb.db.exception.BatchProcessException;
import org.apache.iotdb.db.exception.IoTDBException;
import org.apache.iotdb.db.exception.metadata.IllegalPathException;
@@ -233,6 +232,7 @@ public abstract class RaftMember {
* a thread pool that is used to do commit log tasks asynchronous in
heartbeat thread
*/
private ExecutorService commitLogPool;
+
/**
* logDispatcher buff the logs orderly according to their log indexes and
send them sequentially,
* which avoids the followers receiving out-of-order logs, forcing them to
wait for previous
@@ -240,7 +240,7 @@ public abstract class RaftMember {
*/
private LogDispatcher logDispatcher;
- protected StopStatus stopStatus;
+ private boolean hasSyncedLeaderBeforeRemoved = false;
protected RaftMember() {
}
@@ -253,7 +253,6 @@ public abstract class RaftMember {
this.asyncHeartbeatClientPool = asyncHeartbeatPool;
this.syncHeartbeatClientPool = syncHeartbeatPool;
this.asyncSendLogClientPool = asyncClientPool;
- this.stopStatus = new StopStatus();
}
protected RaftMember(String name, AsyncClientPool asyncPool, SyncClientPool
syncPool,
@@ -265,7 +264,6 @@ public abstract class RaftMember {
this.asyncHeartbeatClientPool = asyncHeartbeatPool;
this.syncHeartbeatClientPool = syncHeartbeatPool;
this.asyncSendLogClientPool = asyncSendLogClientPool;
- this.stopStatus = new StopStatus();
}
/**
@@ -364,7 +362,6 @@ public abstract class RaftMember {
catchUpService = null;
heartBeatService = null;
appendLogThreadPool = null;
- stopStatus.setStop(true);
logger.info("Member {} stopped", name);
}
@@ -783,7 +780,7 @@ public abstract class RaftMember {
* Wait until the leader of this node becomes known or time out.
*/
public void waitLeader() {
- if (stopStatus.isStop()) {
+ if (hasSyncedLeaderBeforeRemoved) {
return;
}
long startTime = System.currentTimeMillis();
@@ -1023,7 +1020,7 @@ public abstract class RaftMember {
}
synchronized (commitIdResult) {
client.requestCommitIndex(getHeader(), getRaftGroupId(), new
GenericHandler<>(leader.get(), commitIdResult));
- commitIdResult.wait(RaftServer.getSyncLeaderMaxWaitMs());
+ commitIdResult.wait(RaftServer.getReadOperationTimeoutMS());
}
return commitIdResult.get();
}
@@ -1576,8 +1573,6 @@ public abstract class RaftMember {
* Send the given log to all the followers and decide the result by how many
followers return a
* success.
*
- * @param requiredQuorum the number of votes needed to make the log valid,
when requiredQuorum <=
- * 0, half of the cluster size will be used.
* @return an AppendLogResult
*/
protected AppendLogResult sendLogToFollowers(Log log) {
@@ -1899,31 +1894,11 @@ public abstract class RaftMember {
OK, TIME_OUT, LEADERSHIP_STALE
}
- public class StopStatus {
-
- boolean stop;
-
- boolean syncSuccess;
-
- public boolean isStop() {
- return stop;
- }
-
- public void setStop(boolean stop) {
- this.stop = stop;
- }
-
- public boolean isSyncSuccess() {
- return syncSuccess;
- }
-
- public void setSyncSuccess(boolean syncSuccess) {
- this.syncSuccess = syncSuccess;
- }
+ public boolean isHasSyncedLeaderBeforeRemoved() {
+ return hasSyncedLeaderBeforeRemoved;
}
- public StopStatus getStopStatus() {
- return stopStatus;
+ public void setHasSyncedLeaderBeforeRemoved(boolean
hasSyncedLeaderAfterRemoved) {
+ this.hasSyncedLeaderBeforeRemoved = hasSyncedLeaderAfterRemoved;
}
-
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
index 641f113..2349228 100755
---
a/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
+++
b/server/src/main/java/org/apache/iotdb/db/engine/storagegroup/StorageGroupProcessor.java
@@ -2644,6 +2644,7 @@ public class StorageGroupProcessor {
if (filter.satisfy(logicalStorageGroupName, partitionId)) {
processor.syncClose();
iterator.remove();
+ processor.getTsFileResource().remove();
tsFileManagement.remove(processor.getTsFileResource(), sequence);
updateLatestFlushTimeToPartition(partitionId, Long.MIN_VALUE);
logger.debug("{} is removed during deleting partitions",
diff --git a/server/src/main/java/org/apache/iotdb/db/utils/CommonUtils.java
b/server/src/main/java/org/apache/iotdb/db/utils/CommonUtils.java
index 30b3952..3451ce6 100644
--- a/server/src/main/java/org/apache/iotdb/db/utils/CommonUtils.java
+++ b/server/src/main/java/org/apache/iotdb/db/utils/CommonUtils.java
@@ -202,8 +202,8 @@ public class CommonUtils {
}
private static void badUse(Exception e) {
- System.out.println("memory-tool: " + e.getMessage());
- System.out.println("See 'memory-tool help' or 'memory-tool help
<command>'.");
+ System.out.println("node-tool: " + e.getMessage());
+ System.out.println("See 'node-tool help' or 'node-tool help <command>'.");
}
private static void err(Throwable e) {