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) {

Reply via email to