This is an automated email from the ASF dual-hosted git repository.

hxd pushed a commit to branch cluster-
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/cluster- by this push:
     new 12ad291  [cluster-refactor]split DataGroupServiceImpls into engine and 
impls  (#4062)
12ad291 is described below

commit 12ad29153c06388af9974c4f3dbdf00b559614c7
Author: lisijia <[email protected]>
AuthorDate: Tue Oct 5 14:31:17 2021 +0800

    [cluster-refactor]split DataGroupServiceImpls into engine and impls  (#4062)
    
    * [cluster-refactor]split DataGroupServiceImpls into engine and impls
---
 .../org/apache/iotdb/cluster/ClusterIoTDB.java     |  25 +-
 .../cluster/server/member/MetaGroupMember.java     |   6 +-
 .../cluster/server/service/DataGroupEngine.java    | 516 ++++++++++++++
 ...ceImplsMBean.java => DataGroupEngineMBean.java} |   2 +-
 .../server/service/DataGroupServiceImpls.java      | 739 +++++----------------
 .../cluster/server/member/MetaGroupMemberTest.java |  29 +-
 6 files changed, 738 insertions(+), 579 deletions(-)

diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
index 40883ce..57d9304 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
@@ -47,6 +47,7 @@ import 
org.apache.iotdb.cluster.server.raft.DataRaftHeartBeatService;
 import org.apache.iotdb.cluster.server.raft.DataRaftService;
 import org.apache.iotdb.cluster.server.raft.MetaRaftHeartBeatService;
 import org.apache.iotdb.cluster.server.raft.MetaRaftService;
+import org.apache.iotdb.cluster.server.service.DataGroupEngine;
 import org.apache.iotdb.cluster.server.service.DataGroupServiceImpls;
 import org.apache.iotdb.cluster.server.service.MetaAsyncService;
 import org.apache.iotdb.cluster.server.service.MetaSyncService;
@@ -105,13 +106,13 @@ public class ClusterIoTDB implements ClusterIoTDBMBean {
 
   private MetaGroupMember metaGroupEngine;
 
-  // TODO we can split dataGroupServiceImpls into two parts: the rpc impl and 
the engine.
-  private DataGroupServiceImpls dataGroupEngine;
+  // split DataGroupServiceImpls into engine and impls
+  private DataGroupEngine dataGroupEngine;
 
   private Node thisNode;
   private Coordinator coordinator;
 
-  private IoTDB iotdb = IoTDB.getInstance();
+  private final IoTDB iotdb = IoTDB.getInstance();
 
   // Cluster IoTDB uses a individual registerManager with its parent.
   private RegisterManager registerManager = new RegisterManager();
@@ -157,7 +158,12 @@ public class ClusterIoTDB implements ClusterIoTDBMBean {
     ((CMManager) IoTDB.metaManager).setCoordinator(coordinator);
     MetaPuller.getInstance().init(metaGroupEngine);
 
-    dataGroupEngine = new DataGroupServiceImpls(protocolFactory, 
metaGroupEngine);
+    // from the scope of the DataGroupEngine,it should be singleton pattern
+    // the way of setting MetaGroupMember in DataGroupEngine may need a better 
modification in
+    // future commit.
+    DataGroupEngine.setProtocolFactory(protocolFactory);
+    DataGroupEngine.setMetaGroupMember(metaGroupEngine);
+    dataGroupEngine = DataGroupEngine.getInstance();
     clientManager =
         new ClientManager(
             ClusterDescriptor.getInstance().getConfig().isUseAsyncServer(),
@@ -331,18 +337,19 @@ public class ClusterIoTDB implements ClusterIoTDBMBean {
     registerManager.register(dataGroupEngine);
 
     // rpc service initialize
+    DataGroupServiceImpls dataGroupServiceImpls = new DataGroupServiceImpls();
     if (ClusterDescriptor.getInstance().getConfig().isUseAsyncServer()) {
       MetaAsyncService metaAsyncService = new 
MetaAsyncService(metaGroupEngine);
       
MetaRaftHeartBeatService.getInstance().initAsyncedServiceImpl(metaAsyncService);
       MetaRaftService.getInstance().initAsyncedServiceImpl(metaAsyncService);
-      DataRaftService.getInstance().initAsyncedServiceImpl(dataGroupEngine);
-      
DataRaftHeartBeatService.getInstance().initAsyncedServiceImpl(dataGroupEngine);
+      
DataRaftService.getInstance().initAsyncedServiceImpl(dataGroupServiceImpls);
+      
DataRaftHeartBeatService.getInstance().initAsyncedServiceImpl(dataGroupServiceImpls);
     } else {
       MetaSyncService syncService = new MetaSyncService(metaGroupEngine);
       
MetaRaftHeartBeatService.getInstance().initSyncedServiceImpl(syncService);
       MetaRaftService.getInstance().initSyncedServiceImpl(syncService);
-      DataRaftService.getInstance().initSyncedServiceImpl(dataGroupEngine);
-      
DataRaftHeartBeatService.getInstance().initSyncedServiceImpl(dataGroupEngine);
+      
DataRaftService.getInstance().initSyncedServiceImpl(dataGroupServiceImpls);
+      
DataRaftHeartBeatService.getInstance().initSyncedServiceImpl(dataGroupServiceImpls);
     }
     // start RPC service
     logger.info("start Meta Heartbeat RPC service... ");
@@ -594,7 +601,7 @@ public class ClusterIoTDB implements ClusterIoTDBMBean {
     return registerManager;
   }
 
-  public DataGroupServiceImpls getDataGroupEngine() {
+  public DataGroupEngine getDataGroupEngine() {
     return dataGroupEngine;
   }
 
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java
index ccfb219..45fb3e0 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/MetaGroupMember.java
@@ -69,7 +69,7 @@ import 
org.apache.iotdb.cluster.server.heartbeat.MetaHeartbeatThread;
 import org.apache.iotdb.cluster.server.monitor.NodeReport.MetaMemberReport;
 import org.apache.iotdb.cluster.server.monitor.NodeStatusManager;
 import org.apache.iotdb.cluster.server.monitor.Timer;
-import org.apache.iotdb.cluster.server.service.DataGroupServiceImpls;
+import org.apache.iotdb.cluster.server.service.DataGroupEngine;
 import org.apache.iotdb.cluster.utils.ClusterUtils;
 import org.apache.iotdb.cluster.utils.PartitionUtils;
 import org.apache.iotdb.cluster.utils.StatusUtils;
@@ -272,7 +272,7 @@ public class MetaGroupMember extends RaftMember implements 
IService, MetaGroupMe
     return localDataMember.closePartition(storageGroupName, partitionId, 
isSeq);
   }
 
-  DataGroupServiceImpls getDataGroupEngine() {
+  DataGroupEngine getDataGroupEngine() {
     return ClusterIoTDB.getInstance().getDataGroupEngine();
   }
 
@@ -1662,7 +1662,7 @@ public class MetaGroupMember extends RaftMember 
implements IService, MetaGroupMe
     this.partitionTable = partitionTable;
     router = new ClusterPlanRouter(partitionTable);
     this.coordinator.setRouter(router);
-    DataGroupServiceImpls dClusterServer = getDataGroupEngine();
+    DataGroupEngine dClusterServer = getDataGroupEngine();
     if (dClusterServer != null) {
       dClusterServer.setPartitionTable(partitionTable);
     }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupEngine.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupEngine.java
new file mode 100644
index 0000000..223a44c
--- /dev/null
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupEngine.java
@@ -0,0 +1,516 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.iotdb.cluster.server.service;
+
+import org.apache.iotdb.cluster.ClusterIoTDB;
+import org.apache.iotdb.cluster.config.ClusterDescriptor;
+import org.apache.iotdb.cluster.exception.CheckConsistencyException;
+import org.apache.iotdb.cluster.exception.NoHeaderNodeException;
+import org.apache.iotdb.cluster.exception.NotInSameGroupException;
+import org.apache.iotdb.cluster.exception.PartitionTableUnavailableException;
+import org.apache.iotdb.cluster.log.logtypes.AddNodeLog;
+import org.apache.iotdb.cluster.log.logtypes.RemoveNodeLog;
+import org.apache.iotdb.cluster.partition.NodeAdditionResult;
+import org.apache.iotdb.cluster.partition.NodeRemovalResult;
+import org.apache.iotdb.cluster.partition.PartitionGroup;
+import org.apache.iotdb.cluster.partition.PartitionTable;
+import org.apache.iotdb.cluster.partition.slot.SlotPartitionTable;
+import org.apache.iotdb.cluster.rpc.thrift.Node;
+import org.apache.iotdb.cluster.rpc.thrift.RaftNode;
+import org.apache.iotdb.cluster.server.NodeCharacter;
+import org.apache.iotdb.cluster.server.StoppedMemberManager;
+import org.apache.iotdb.cluster.server.member.DataGroupMember;
+import org.apache.iotdb.cluster.server.member.MetaGroupMember;
+import org.apache.iotdb.cluster.server.monitor.NodeReport.DataMemberReport;
+import org.apache.iotdb.db.exception.StartupException;
+import org.apache.iotdb.db.service.IService;
+import org.apache.iotdb.db.service.ServiceType;
+import org.apache.iotdb.db.utils.TestOnly;
+
+import org.apache.thrift.async.AsyncMethodCallback;
+import org.apache.thrift.protocol.TProtocolFactory;
+import org.apache.thrift.transport.TTransportException;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.ArrayList;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.Map.Entry;
+import java.util.concurrent.ConcurrentHashMap;
+
+public class DataGroupEngine implements IService, DataGroupEngineMBean {
+
+  private static final Logger logger = 
LoggerFactory.getLogger(DataGroupEngine.class);
+  // key: the header of a data group, value: the member representing this node 
in this group and
+  // it is currently at service
+  private static final Map<RaftNode, DataGroupMember> headerGroupMap = new 
ConcurrentHashMap<>();
+  private static final Map<RaftNode, DataAsyncService> asyncServiceMap = new 
ConcurrentHashMap<>();
+  private static final Map<RaftNode, DataSyncService> syncServiceMap = new 
ConcurrentHashMap<>();
+  // key: the header of a data group, value: the member representing this node 
in this group but
+  // it is out of service because another node has joined the group and 
expelled this node, or
+  // the node itself is removed, but it is still stored to provide snapshot 
for other nodes
+  private final StoppedMemberManager stoppedMemberManager;
+  private PartitionTable partitionTable;
+  private DataGroupMember.Factory dataMemberFactory;
+  private static MetaGroupMember metaGroupMember;
+  private final Node thisNode = ClusterIoTDB.getInstance().getThisNode();
+  private static TProtocolFactory protocolFactory;
+
+  private DataGroupEngine() {
+    dataMemberFactory = new DataGroupMember.Factory(protocolFactory, 
metaGroupMember);
+    stoppedMemberManager = new StoppedMemberManager(dataMemberFactory);
+  }
+
+  public static DataGroupEngine getInstance() {
+    if (metaGroupMember == null || protocolFactory == null) {
+      logger.error("MetaGroupMember or protocolFactory init failed.");
+    }
+    return InstanceHolder.Instance;
+  }
+
+  @TestOnly
+  public DataGroupEngine(
+      DataGroupMember.Factory dataMemberFactory, MetaGroupMember 
metaGroupMember) {
+    DataGroupEngine.metaGroupMember = metaGroupMember;
+    this.stoppedMemberManager = new StoppedMemberManager(dataMemberFactory);
+  }
+
+  @Override
+  public void start() throws StartupException {}
+
+  @Override
+  public void stop() {
+    closeLogManagers();
+    for (DataGroupMember member : headerGroupMap.values()) {
+      member.stop();
+    }
+  }
+
+  @Override
+  public ServiceType getID() {
+    return ServiceType.CLUSTER_DATA_ENGINE;
+  }
+
+  public void closeLogManagers() {
+    for (DataGroupMember member : headerGroupMap.values()) {
+      member.closeLogManager();
+    }
+  }
+
+  public <T> DataAsyncService getDataAsyncService(
+      RaftNode header, AsyncMethodCallback<T> resultHandler, Object request) {
+    return asyncServiceMap.computeIfAbsent(
+        header,
+        h -> {
+          DataGroupMember dataMember = getDataMember(header, resultHandler, 
request);
+          return dataMember != null ? new DataAsyncService(dataMember) : null;
+        });
+  }
+
+  public <T> DataSyncService getDataSyncService(RaftNode header) {
+    return syncServiceMap.computeIfAbsent(
+        header,
+        h -> {
+          DataGroupMember dataMember = getDataMember(header, null, null);
+          return dataMember != null ? new DataSyncService(dataMember) : null;
+        });
+  }
+
+  /**
+   * Add a DataGroupMember into this server, if a member with the same header 
exists, the old member
+   * will be stopped and replaced by the new one.
+   *
+   * @param dataGroupMember
+   */
+  public DataGroupMember addDataGroupMember(DataGroupMember dataGroupMember, 
RaftNode header) {
+    synchronized (headerGroupMap) {
+      // TODO this method won't update headerMap if a new dataGroupMember 
comes with the same header
+      // ?
+      if (headerGroupMap.containsKey(header)) {
+        logger.debug("Group {} already exist.", dataGroupMember.getAllNodes());
+        return headerGroupMap.get(header);
+      }
+      stoppedMemberManager.remove(header);
+      headerGroupMap.put(header, dataGroupMember);
+
+      dataGroupMember.start();
+    }
+    logger.info("Add group {} successfully.", dataGroupMember.getName());
+    resetServiceCache(header); // avoid dead-lock
+
+    return dataGroupMember;
+  }
+
+  private void resetServiceCache(RaftNode header) {
+    asyncServiceMap.remove(header);
+    syncServiceMap.remove(header);
+  }
+
+  /**
+   * @param header the header of the group which the local node is in
+   * @param resultHandler can be set to null if the request is an internal 
request
+   * @param request the toString() of this parameter should explain what the 
request is and it is
+   *     only used in logs for tracing
+   * @return
+   */
+  public <T> DataGroupMember getDataMember(
+      RaftNode header, AsyncMethodCallback<T> resultHandler, Object request) {
+    // if the resultHandler is not null, then the request is a external one 
and must be with a
+    // header
+    if (header.getNode() == null) {
+      if (resultHandler != null) {
+        resultHandler.onError(new NoHeaderNodeException());
+      }
+      return null;
+    }
+    DataGroupMember member = stoppedMemberManager.get(header);
+    if (member != null) {
+      return member;
+    }
+
+    // avoid creating two members for a header
+    Exception ex = null;
+    member = headerGroupMap.get(header);
+    if (member != null) {
+      return member;
+    }
+    logger.info("Received a request \"{}\" from unregistered header {}", 
request, header);
+    if (partitionTable != null) {
+      try {
+        member = createNewMember(header);
+      } catch (NotInSameGroupException | CheckConsistencyException e) {
+        ex = e;
+      }
+    } else {
+      logger.info("Partition is not ready, cannot create member");
+      ex = new PartitionTableUnavailableException(thisNode);
+    }
+    if (ex != null && resultHandler != null) {
+      resultHandler.onError(ex);
+    }
+    return member;
+  }
+
+  /**
+   * @param header
+   * @return A DataGroupMember representing this node in the data group of the 
header.
+   * @throws NotInSameGroupException If this node is not in the group of the 
header.
+   */
+  private DataGroupMember createNewMember(RaftNode header)
+      throws NotInSameGroupException, CheckConsistencyException {
+    PartitionGroup partitionGroup;
+    partitionGroup = partitionTable.getHeaderGroup(header);
+    if (partitionGroup == null || !partitionGroup.contains(thisNode)) {
+      // if the partition table is old, this node may have not been moved to 
the new group
+      metaGroupMember.syncLeaderWithConsistencyCheck(true);
+      partitionGroup = partitionTable.getHeaderGroup(header);
+    }
+    DataGroupMember member;
+    synchronized (headerGroupMap) {
+      member = headerGroupMap.get(header);
+      if (member != null) {
+        return member;
+      }
+      if (partitionGroup != null && partitionGroup.contains(thisNode)) {
+        // the two nodes are in the same group, create a new data member
+        member = dataMemberFactory.create(partitionGroup);
+        headerGroupMap.put(header, member);
+        stoppedMemberManager.remove(header);
+        logger.info("Created a member for header {}, group is {}", header, 
partitionGroup);
+        member.start();
+      } else {
+        // the member may have been stopped after syncLeader
+        member = stoppedMemberManager.get(header);
+        if (member != null) {
+          return member;
+        }
+        logger.info(
+            "This node {} does not belong to the group {}, header {}",
+            thisNode,
+            partitionGroup,
+            header);
+        throw new NotInSameGroupException(partitionGroup, thisNode);
+      }
+    }
+    return member;
+  }
+
+  public void preAddNodeForDataGroup(AddNodeLog log, DataGroupMember 
targetDataGroupMember) {
+
+    // Make sure the previous add/remove node log has applied
+    metaGroupMember.syncLocalApply(log.getMetaLogIndex() - 1, false);
+
+    // Check the validity of the partition table
+    if 
(!metaGroupMember.getPartitionTable().deserialize(log.getPartitionTable())) {
+      return;
+    }
+
+    targetDataGroupMember.preAddNode(log.getNewNode());
+  }
+
+  /**
+   * Try adding the node into the group of each DataGroupMember, and if the 
DataGroupMember no
+   * longer stays in that group, also remove and stop it. If the new group 
contains this node, also
+   * create and add a new DataGroupMember for it.
+   *
+   * @param node
+   * @param result
+   */
+  public void addNode(Node node, NodeAdditionResult result) {
+    // If the node executed adding itself to the cluster, it's unnecessary to 
add new groups because
+    // they already exist.
+    if (node.equals(thisNode)) {
+      return;
+    }
+    Iterator<Entry<RaftNode, DataGroupMember>> entryIterator = 
headerGroupMap.entrySet().iterator();
+    synchronized (headerGroupMap) {
+      while (entryIterator.hasNext()) {
+        Entry<RaftNode, DataGroupMember> entry = entryIterator.next();
+        DataGroupMember dataGroupMember = entry.getValue();
+        // the member may be extruded from the group, remove and stop it if so
+        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();
+        }
+      }
+
+      if (logger.isDebugEnabled()) {
+        logger.debug(
+            "Data cluster server: start to handle new groups when adding new 
node {}", node);
+      }
+      for (PartitionGroup newGroup : result.getNewGroupList()) {
+        if (newGroup.contains(thisNode)) {
+          RaftNode header = newGroup.getHeader();
+          logger.info("Adding this node into a new group {}", newGroup);
+          DataGroupMember dataGroupMember = dataMemberFactory.create(newGroup);
+          dataGroupMember = addDataGroupMember(dataGroupMember, header);
+          dataGroupMember.pullNodeAdditionSnapshots(
+              ((SlotPartitionTable) partitionTable).getNodeSlots(header), 
node);
+        }
+      }
+    }
+  }
+
+  /**
+   * When the node joins a cluster, it also creates a new data group and a 
corresponding member
+   * which has no data. This is to make that member pull data from other nodes.
+   */
+  public void pullSnapshots() {
+    for (int raftId = 0;
+        raftId < 
ClusterDescriptor.getInstance().getConfig().getMultiRaftFactor();
+        raftId++) {
+      RaftNode raftNode = new RaftNode(thisNode, raftId);
+      List<Integer> slots = ((SlotPartitionTable) 
partitionTable).getNodeSlots(raftNode);
+      DataGroupMember dataGroupMember = headerGroupMap.get(raftNode);
+      dataGroupMember.pullNodeAdditionSnapshots(slots, thisNode);
+    }
+  }
+
+  /**
+   * Make sure the group will not receive new raft logs
+   *
+   * @param header
+   * @param dataGroupMember
+   */
+  private void removeMember(
+      RaftNode header, DataGroupMember dataGroupMember, boolean removedGroup) {
+    dataGroupMember.setReadOnly();
+    if (!removedGroup) {
+      dataGroupMember.stop();
+    } else {
+      if (dataGroupMember.getCharacter() != NodeCharacter.LEADER) {
+        new Thread(
+                () -> {
+                  try {
+                    dataGroupMember.syncLeader(null);
+                    dataGroupMember.stop();
+                  } catch (CheckConsistencyException e) {
+                    logger.warn("Failed to check consistency.", e);
+                  }
+                })
+            .start();
+      }
+    }
+    stoppedMemberManager.put(header, dataGroupMember);
+    logger.info(
+        "Data group member has removed, header {}, group is {}.",
+        header,
+        dataGroupMember.getAllNodes());
+  }
+
+  /**
+   * Set the partition table as the in-use one and build a DataGroupMember for 
each local group (the
+   * group which the local node is in) and start them.
+   *
+   * @param partitionTable
+   * @throws TTransportException
+   */
+  @SuppressWarnings("java:S1135")
+  public void buildDataGroupMembers(PartitionTable partitionTable) {
+    setPartitionTable(partitionTable);
+    // TODO-Cluster: if there are unchanged members, do not stop and restart 
them
+    // clear previous members if the partition table is reloaded
+    for (DataGroupMember value : headerGroupMap.values()) {
+      value.stop();
+    }
+
+    for (DataGroupMember value : headerGroupMap.values()) {
+      value.setUnchanged(false);
+    }
+
+    List<PartitionGroup> partitionGroups = partitionTable.getLocalGroups();
+    for (PartitionGroup partitionGroup : partitionGroups) {
+      RaftNode header = partitionGroup.getHeader();
+      DataGroupMember prevMember = headerGroupMap.get(header);
+      if (prevMember == null || 
!prevMember.getAllNodes().equals(partitionGroup)) {
+        logger.info("Building member of data group: {}", partitionGroup);
+        // no previous member or member changed
+        DataGroupMember dataGroupMember = 
dataMemberFactory.create(partitionGroup);
+        // the previous member will be replaced here
+        addDataGroupMember(dataGroupMember, header);
+        dataGroupMember.setUnchanged(true);
+      } else {
+        prevMember.setUnchanged(true);
+        prevMember.start();
+        // TODO do we nedd call other functions in addDataGroupMember() ?
+      }
+    }
+
+    // remove out-dated members of this node
+    headerGroupMap.entrySet().removeIf(e -> !e.getValue().isUnchanged());
+
+    logger.info("Data group members are ready");
+  }
+
+  public void preRemoveNodeForDataGroup(RemoveNodeLog log, DataGroupMember 
targetDataGroupMember) {
+
+    // Make sure the previous add/remove node log has applied
+    metaGroupMember.syncLocalApply(log.getMetaLogIndex() - 1, false);
+
+    // Check the validity of the partition table
+    if 
(!metaGroupMember.getPartitionTable().deserialize(log.getPartitionTable())) {
+      return;
+    }
+
+    logger.debug(
+        "Pre removing a node {} from {}",
+        log.getRemovedNode(),
+        targetDataGroupMember.getAllNodes());
+    targetDataGroupMember.preRemoveNode(log.getRemovedNode());
+  }
+
+  /**
+   * Try removing a node from the groups of each DataGroupMember. If the node 
is the header of some
+   * group, set the member to read only so that it can still provide data for 
other nodes that has
+   * not yet pulled its data. Otherwise, just change the node list of the 
member and pull new data.
+   * And create a new DataGroupMember if this node should join a new group 
because of this removal.
+   *
+   * @param node
+   * @param removalResult cluster changes due to the node removal
+   */
+  public void removeNode(Node node, NodeRemovalResult removalResult) {
+    Iterator<Entry<RaftNode, DataGroupMember>> entryIterator = 
headerGroupMap.entrySet().iterator();
+    synchronized (headerGroupMap) {
+      while (entryIterator.hasNext()) {
+        Entry<RaftNode, DataGroupMember> entry = entryIterator.next();
+        DataGroupMember dataGroupMember = entry.getValue();
+        if (dataGroupMember.getHeader().getNode().equals(node) || 
node.equals(thisNode)) {
+          entryIterator.remove();
+          removeMember(
+              entry.getKey(), dataGroupMember, 
dataGroupMember.getHeader().getNode().equals(node));
+        } else {
+          // the group should be updated
+          dataGroupMember.removeNode(node);
+        }
+      }
+
+      if (logger.isDebugEnabled()) {
+        logger.debug(
+            "Data cluster server: start to handle new groups and pulling data 
when removing node {}",
+            node);
+      }
+      // if the removed group contains the local node, the local node should 
join a new group to
+      // preserve the replication number
+      for (PartitionGroup group : partitionTable.getLocalGroups()) {
+        RaftNode header = group.getHeader();
+        if (!headerGroupMap.containsKey(header)) {
+          logger.info("{} should join a new group {}", thisNode, group);
+          DataGroupMember dataGroupMember = dataMemberFactory.create(group);
+          addDataGroupMember(dataGroupMember, header);
+        }
+        // pull new slots from the removed node
+        headerGroupMap.get(header).pullSlots(removalResult);
+      }
+    }
+  }
+
+  public void setPartitionTable(PartitionTable partitionTable) {
+    this.partitionTable = partitionTable;
+  }
+
+  /** @return The reports of every DataGroupMember in this node. */
+  public List<DataMemberReport> genMemberReports() {
+    List<DataMemberReport> dataMemberReports = new ArrayList<>();
+    for (DataGroupMember value : headerGroupMap.values()) {
+
+      dataMemberReports.add(value.genReport());
+    }
+    return dataMemberReports;
+  }
+
+  public Map<RaftNode, DataGroupMember> getHeaderGroupMap() {
+    return headerGroupMap;
+  }
+
+  public static void setProtocolFactory(TProtocolFactory protocolFactory) {
+    DataGroupEngine.protocolFactory = protocolFactory;
+  }
+
+  public static void setMetaGroupMember(MetaGroupMember metaGroupMember) {
+    DataGroupEngine.metaGroupMember = metaGroupMember;
+  }
+
+  static class InstanceHolder {
+    private static final DataGroupEngine Instance = new DataGroupEngine();
+  }
+
+  @Override
+  public String getHeaderGroupMapAsString() {
+    return headerGroupMap.toString();
+  }
+
+  @Override
+  public int getAsyncServiceMapSize() {
+    return asyncServiceMap.size();
+  }
+
+  @Override
+  public int getSyncServiceMapSize() {
+    return syncServiceMap.size();
+  }
+
+  @Override
+  public String getPartitionTable() {
+    return partitionTable.toString();
+  }
+}
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupServiceImplsMBean.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupEngineMBean.java
similarity index 95%
rename from 
cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupServiceImplsMBean.java
rename to 
cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupEngineMBean.java
index 988fc6f..b0fcb84 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupServiceImplsMBean.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupEngineMBean.java
@@ -18,7 +18,7 @@
  */
 package org.apache.iotdb.cluster.server.service;
 
-public interface DataGroupServiceImplsMBean {
+public interface DataGroupEngineMBean {
 
   String getHeaderGroupMapAsString();
 
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupServiceImpls.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupServiceImpls.java
index ac660dd..e46fb7e 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupServiceImpls.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/service/DataGroupServiceImpls.java
@@ -19,19 +19,6 @@
 
 package org.apache.iotdb.cluster.server.service;
 
-import org.apache.iotdb.cluster.ClusterIoTDB;
-import org.apache.iotdb.cluster.config.ClusterDescriptor;
-import org.apache.iotdb.cluster.exception.CheckConsistencyException;
-import org.apache.iotdb.cluster.exception.NoHeaderNodeException;
-import org.apache.iotdb.cluster.exception.NotInSameGroupException;
-import org.apache.iotdb.cluster.exception.PartitionTableUnavailableException;
-import org.apache.iotdb.cluster.log.logtypes.AddNodeLog;
-import org.apache.iotdb.cluster.log.logtypes.RemoveNodeLog;
-import org.apache.iotdb.cluster.partition.NodeAdditionResult;
-import org.apache.iotdb.cluster.partition.NodeRemovalResult;
-import org.apache.iotdb.cluster.partition.PartitionGroup;
-import org.apache.iotdb.cluster.partition.PartitionTable;
-import org.apache.iotdb.cluster.partition.slot.SlotPartitionTable;
 import org.apache.iotdb.cluster.rpc.thrift.AppendEntriesRequest;
 import org.apache.iotdb.cluster.rpc.thrift.AppendEntryRequest;
 import org.apache.iotdb.cluster.rpc.thrift.ElectionRequest;
@@ -54,232 +41,28 @@ import 
org.apache.iotdb.cluster.rpc.thrift.RequestCommitIndexResponse;
 import org.apache.iotdb.cluster.rpc.thrift.SendSnapshotRequest;
 import org.apache.iotdb.cluster.rpc.thrift.SingleSeriesQueryRequest;
 import org.apache.iotdb.cluster.rpc.thrift.TSDataService;
-import org.apache.iotdb.cluster.server.NodeCharacter;
-import org.apache.iotdb.cluster.server.StoppedMemberManager;
-import org.apache.iotdb.cluster.server.member.DataGroupMember;
-import org.apache.iotdb.cluster.server.member.MetaGroupMember;
-import org.apache.iotdb.cluster.server.monitor.NodeReport.DataMemberReport;
 import org.apache.iotdb.cluster.utils.IOUtils;
-import org.apache.iotdb.db.exception.StartupException;
-import org.apache.iotdb.db.service.IService;
-import org.apache.iotdb.db.service.ServiceType;
-import org.apache.iotdb.db.utils.TestOnly;
 import org.apache.iotdb.service.rpc.thrift.TSStatus;
 
 import org.apache.thrift.TException;
 import org.apache.thrift.async.AsyncMethodCallback;
-import org.apache.thrift.protocol.TProtocolFactory;
-import org.apache.thrift.transport.TTransportException;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
 
 import java.io.File;
 import java.io.IOException;
 import java.nio.ByteBuffer;
 import java.nio.file.Files;
-import java.util.ArrayList;
-import java.util.Iterator;
 import java.util.List;
 import java.util.Map;
-import java.util.Map.Entry;
 import java.util.Set;
-import java.util.concurrent.ConcurrentHashMap;
 
-public class DataGroupServiceImpls
-    implements TSDataService.AsyncIface, TSDataService.Iface, IService, 
DataGroupServiceImplsMBean {
-
-  private static final Logger logger = 
LoggerFactory.getLogger(DataGroupServiceImpls.class);
-
-  // key: the header of a data group, value: the member representing this node 
in this group and
-  // it is currently at service
-  private Map<RaftNode, DataGroupMember> headerGroupMap = new 
ConcurrentHashMap<>();
-  private Map<RaftNode, DataAsyncService> asyncServiceMap = new 
ConcurrentHashMap<>();
-  private Map<RaftNode, DataSyncService> syncServiceMap = new 
ConcurrentHashMap<>();
-  // key: the header of a data group, value: the member representing this node 
in this group but
-  // it is out of service because another node has joined the group and 
expelled this node, or
-  // the node itself is removed, but it is still stored to provide snapshot 
for other nodes
-  private StoppedMemberManager stoppedMemberManager;
-  private PartitionTable partitionTable;
-  private DataGroupMember.Factory dataMemberFactory;
-  private MetaGroupMember metaGroupMember;
-
-  private Node thisNode = ClusterIoTDB.getInstance().getThisNode();
-
-  public DataGroupServiceImpls(TProtocolFactory protocolFactory, 
MetaGroupMember metaGroupMember) {
-    dataMemberFactory = new DataGroupMember.Factory(protocolFactory, 
metaGroupMember);
-    this.metaGroupMember = metaGroupMember;
-    stoppedMemberManager = new StoppedMemberManager(dataMemberFactory);
-  }
-
-  @TestOnly
-  public DataGroupServiceImpls(
-      DataGroupMember.Factory dataMemberFactory, MetaGroupMember 
metaGroupMember) {
-    this.metaGroupMember = metaGroupMember;
-    this.stoppedMemberManager = new StoppedMemberManager(dataMemberFactory);
-  }
-
-  @Override
-  public void start() throws StartupException {
-    // seems do nothing
-  }
-
-  //  @Override
-  // TODO
-  public void stop() {
-    closeLogManagers();
-    for (DataGroupMember member : headerGroupMap.values()) {
-      member.stop();
-    }
-  }
-
-  @Override
-  public ServiceType getID() {
-    return ServiceType.CLUSTER_DATA_ENGINE;
-  }
-
-  /**
-   * Add a DataGroupMember into this server, if a member with the same header 
exists, the old member
-   * will be stopped and replaced by the new one.
-   *
-   * @param dataGroupMember
-   */
-  public DataGroupMember addDataGroupMember(DataGroupMember dataGroupMember, 
RaftNode header) {
-    synchronized (headerGroupMap) {
-      if (headerGroupMap.containsKey(header)) {
-        logger.debug("Group {} already exist.", dataGroupMember.getAllNodes());
-        return headerGroupMap.get(header);
-      }
-      stoppedMemberManager.remove(header);
-      headerGroupMap.put(header, dataGroupMember);
-
-      dataGroupMember.start();
-    }
-    logger.info("Add group {} successfully.", dataGroupMember.getName());
-    resetServiceCache(header); // avoid dead-lock
-
-    return dataGroupMember;
-  }
-
-  private void resetServiceCache(RaftNode header) {
-    asyncServiceMap.remove(header);
-    syncServiceMap.remove(header);
-  }
-
-  private <T> DataAsyncService getDataAsyncService(
-      RaftNode header, AsyncMethodCallback<T> resultHandler, Object request) {
-    return asyncServiceMap.computeIfAbsent(
-        header,
-        h -> {
-          DataGroupMember dataMember = getDataMember(header, resultHandler, 
request);
-          return dataMember != null ? new DataAsyncService(dataMember) : null;
-        });
-  }
-
-  private DataSyncService getDataSyncService(RaftNode header) {
-    return syncServiceMap.computeIfAbsent(
-        header,
-        h -> {
-          DataGroupMember dataMember = getDataMember(header, null, null);
-          return dataMember != null ? new DataSyncService(dataMember) : null;
-        });
-  }
-
-  /**
-   * @param header the header of the group which the local node is in
-   * @param resultHandler can be set to null if the request is an internal 
request
-   * @param request the toString() of this parameter should explain what the 
request is and it is
-   *     only used in logs for tracing
-   * @return
-   */
-  public <T> DataGroupMember getDataMember(
-      RaftNode header, AsyncMethodCallback<T> resultHandler, Object request) {
-    // if the resultHandler is not null, then the request is a external one 
and must be with a
-    // header
-    if (header.getNode() == null) {
-      if (resultHandler != null) {
-        resultHandler.onError(new NoHeaderNodeException());
-      }
-      return null;
-    }
-    DataGroupMember member = stoppedMemberManager.get(header);
-    if (member != null) {
-      return member;
-    }
-
-    // avoid creating two members for a header
-    Exception ex = null;
-    member = headerGroupMap.get(header);
-    if (member != null) {
-      return member;
-    }
-    logger.info("Received a request \"{}\" from unregistered header {}", 
request, header);
-    if (partitionTable != null) {
-      try {
-        member = createNewMember(header);
-      } catch (NotInSameGroupException | CheckConsistencyException e) {
-        ex = e;
-      }
-    } else {
-      logger.info("Partition is not ready, cannot create member");
-      ex = new PartitionTableUnavailableException(thisNode);
-    }
-    if (ex != null && resultHandler != null) {
-      resultHandler.onError(ex);
-    }
-    return member;
-  }
-
-  /**
-   * @param header
-   * @return A DataGroupMember representing this node in the data group of the 
header.
-   * @throws NotInSameGroupException If this node is not in the group of the 
header.
-   */
-  private DataGroupMember createNewMember(RaftNode header)
-      throws NotInSameGroupException, CheckConsistencyException {
-    PartitionGroup partitionGroup;
-    partitionGroup = partitionTable.getHeaderGroup(header);
-    if (partitionGroup == null || !partitionGroup.contains(thisNode)) {
-      // if the partition table is old, this node may have not been moved to 
the new group
-      metaGroupMember.syncLeaderWithConsistencyCheck(true);
-      partitionGroup = partitionTable.getHeaderGroup(header);
-    }
-    DataGroupMember member;
-    synchronized (headerGroupMap) {
-      member = headerGroupMap.get(header);
-      if (member != null) {
-        return member;
-      }
-      if (partitionGroup != null && partitionGroup.contains(thisNode)) {
-        // the two nodes are in the same group, create a new data member
-        member = dataMemberFactory.create(partitionGroup);
-        headerGroupMap.put(header, member);
-        stoppedMemberManager.remove(header);
-        logger.info("Created a member for header {}, group is {}", header, 
partitionGroup);
-        member.start();
-      } else {
-        // the member may have been stopped after syncLeader
-        member = stoppedMemberManager.get(header);
-        if (member != null) {
-          return member;
-        }
-        logger.info(
-            "This node {} does not belong to the group {}, header {}",
-            thisNode,
-            partitionGroup,
-            header);
-        throw new NotInSameGroupException(partitionGroup, thisNode);
-      }
-    }
-    return member;
-  }
-
-  // Forward requests. Find the DataGroupMember that is in the group of the 
header of the
-  // request, and forward the request to it. See methods in DataGroupMember 
for details.
+public class DataGroupServiceImpls implements TSDataService.AsyncIface, 
TSDataService.Iface {
 
   @Override
   public void sendHeartbeat(
       HeartBeatRequest request, AsyncMethodCallback<HeartBeatResponse> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, request);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, request);
     if (service != null) {
       service.sendHeartbeat(request, resultHandler);
     }
@@ -287,7 +70,9 @@ public class DataGroupServiceImpls
 
   @Override
   public void startElection(ElectionRequest request, AsyncMethodCallback<Long> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, request);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, request);
     if (service != null) {
       service.startElection(request, resultHandler);
     }
@@ -295,7 +80,9 @@ public class DataGroupServiceImpls
 
   @Override
   public void appendEntries(AppendEntriesRequest request, 
AsyncMethodCallback<Long> resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, request);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, request);
     if (service != null) {
       service.appendEntries(request, resultHandler);
     }
@@ -303,7 +90,9 @@ public class DataGroupServiceImpls
 
   @Override
   public void appendEntry(AppendEntryRequest request, 
AsyncMethodCallback<Long> resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, request);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, request);
     if (service != null) {
       service.appendEntry(request, resultHandler);
     }
@@ -311,7 +100,9 @@ public class DataGroupServiceImpls
 
   @Override
   public void sendSnapshot(SendSnapshotRequest request, 
AsyncMethodCallback<Void> resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, request);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, request);
     if (service != null) {
       service.sendSnapshot(request, resultHandler);
     }
@@ -320,7 +111,9 @@ public class DataGroupServiceImpls
   @Override
   public void pullSnapshot(
       PullSnapshotRequest request, AsyncMethodCallback<PullSnapshotResp> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, request);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, request);
     if (service != null) {
       service.pullSnapshot(request, resultHandler);
     }
@@ -329,7 +122,9 @@ public class DataGroupServiceImpls
   @Override
   public void executeNonQueryPlan(
       ExecutNonQueryReq request, AsyncMethodCallback<TSStatus> resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, request);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, request);
     if (service != null) {
       service.executeNonQueryPlan(request, resultHandler);
     }
@@ -338,7 +133,9 @@ public class DataGroupServiceImpls
   @Override
   public void requestCommitIndex(
       RaftNode header, AsyncMethodCallback<RequestCommitIndexResponse> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(header, resultHandler, 
"Request commit index");
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(header, resultHandler, "Request commit 
index");
     if (service != null) {
       service.requestCommitIndex(header, resultHandler);
     }
@@ -358,8 +155,9 @@ public class DataGroupServiceImpls
   public void querySingleSeries(
       SingleSeriesQueryRequest request, AsyncMethodCallback<Long> 
resultHandler) {
     DataAsyncService service =
-        getDataAsyncService(
-            request.getHeader(), resultHandler, "Query series:" + 
request.getPath());
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(
+                request.getHeader(), resultHandler, "Query series:" + 
request.getPath());
     if (service != null) {
       service.querySingleSeries(request, resultHandler);
     }
@@ -369,8 +167,9 @@ public class DataGroupServiceImpls
   public void queryMultSeries(
       MultSeriesQueryRequest request, AsyncMethodCallback<Long> resultHandler) 
throws TException {
     DataAsyncService service =
-        getDataAsyncService(
-            request.getHeader(), resultHandler, "Query series:" + 
request.getPath());
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(
+                request.getHeader(), resultHandler, "Query series:" + 
request.getPath());
     if (service != null) {
       service.queryMultSeries(request, resultHandler);
     }
@@ -380,7 +179,8 @@ public class DataGroupServiceImpls
   public void fetchSingleSeries(
       RaftNode header, long readerId, AsyncMethodCallback<ByteBuffer> 
resultHandler) {
     DataAsyncService service =
-        getDataAsyncService(header, resultHandler, "Fetch reader:" + readerId);
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(header, resultHandler, "Fetch reader:" + 
readerId);
     if (service != null) {
       service.fetchSingleSeries(header, readerId, resultHandler);
     }
@@ -394,7 +194,8 @@ public class DataGroupServiceImpls
       AsyncMethodCallback<Map<String, ByteBuffer>> resultHandler)
       throws TException {
     DataAsyncService service =
-        getDataAsyncService(header, resultHandler, "Fetch reader:" + readerId);
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(header, resultHandler, "Fetch reader:" + 
readerId);
     if (service != null) {
       service.fetchMultSeries(header, readerId, paths, resultHandler);
     }
@@ -406,7 +207,9 @@ public class DataGroupServiceImpls
       List<String> paths,
       boolean withAlias,
       AsyncMethodCallback<GetAllPathsResult> resultHandler) {
-    DataAsyncService service = getDataAsyncService(header, resultHandler, 
"Find path:" + paths);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(header, resultHandler, "Find path:" + paths);
     if (service != null) {
       service.getAllPaths(header, paths, withAlias, resultHandler);
     }
@@ -415,7 +218,8 @@ public class DataGroupServiceImpls
   @Override
   public void endQuery(
       RaftNode header, Node thisNode, long queryId, AsyncMethodCallback<Void> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(header, resultHandler, "End 
query");
+    DataAsyncService service =
+        DataGroupEngine.getInstance().getDataAsyncService(header, 
resultHandler, "End query");
     if (service != null) {
       service.endQuery(header, thisNode, queryId, resultHandler);
     }
@@ -425,15 +229,16 @@ public class DataGroupServiceImpls
   public void querySingleSeriesByTimestamp(
       SingleSeriesQueryRequest request, AsyncMethodCallback<Long> 
resultHandler) {
     DataAsyncService service =
-        getDataAsyncService(
-            request.getHeader(),
-            resultHandler,
-            "Query by timestamp:"
-                + request.getQueryId()
-                + "#"
-                + request.getPath()
-                + " of "
-                + request.getRequester());
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(
+                request.getHeader(),
+                resultHandler,
+                "Query by timestamp:"
+                    + request.getQueryId()
+                    + "#"
+                    + request.getPath()
+                    + " of "
+                    + request.getRequester());
     if (service != null) {
       service.querySingleSeriesByTimestamp(request, resultHandler);
     }
@@ -446,7 +251,8 @@ public class DataGroupServiceImpls
       List<Long> timestamps,
       AsyncMethodCallback<ByteBuffer> resultHandler) {
     DataAsyncService service =
-        getDataAsyncService(header, resultHandler, "Fetch by timestamp:" + 
readerId);
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(header, resultHandler, "Fetch by timestamp:" 
+ readerId);
     if (service != null) {
       service.fetchSingleSeriesByTimestamps(header, readerId, timestamps, 
resultHandler);
     }
@@ -455,7 +261,9 @@ public class DataGroupServiceImpls
   @Override
   public void pullTimeSeriesSchema(
       PullSchemaRequest request, AsyncMethodCallback<PullSchemaResp> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, request);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, request);
     if (service != null) {
       service.pullTimeSeriesSchema(request, resultHandler);
     }
@@ -465,7 +273,8 @@ public class DataGroupServiceImpls
   public void pullMeasurementSchema(
       PullSchemaRequest request, AsyncMethodCallback<PullSchemaResp> 
resultHandler) {
     DataAsyncService service =
-        getDataAsyncService(request.getHeader(), resultHandler, "Pull 
measurement schema");
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, "Pull 
measurement schema");
     if (service != null) {
       service.pullMeasurementSchema(request, resultHandler);
     }
@@ -474,7 +283,8 @@ public class DataGroupServiceImpls
   @Override
   public void getAllDevices(
       RaftNode header, List<String> paths, AsyncMethodCallback<Set<String>> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(header, resultHandler, "Get 
all devices");
+    DataAsyncService service =
+        DataGroupEngine.getInstance().getDataAsyncService(header, 
resultHandler, "Get all devices");
     if (service != null) {
       service.getAllDevices(header, paths, resultHandler);
     }
@@ -484,7 +294,8 @@ public class DataGroupServiceImpls
   public void getDevices(
       RaftNode header, ByteBuffer planBinary, AsyncMethodCallback<ByteBuffer> 
resultHandler)
       throws TException {
-    DataAsyncService service = getDataAsyncService(header, resultHandler, "get 
devices");
+    DataAsyncService service =
+        DataGroupEngine.getInstance().getDataAsyncService(header, 
resultHandler, "get devices");
     if (service != null) {
       service.getDevices(header, planBinary, resultHandler);
     }
@@ -496,7 +307,8 @@ public class DataGroupServiceImpls
       String path,
       int nodeLevel,
       AsyncMethodCallback<List<String>> resultHandler) {
-    DataAsyncService service = getDataAsyncService(header, resultHandler, "Get 
node list");
+    DataAsyncService service =
+        DataGroupEngine.getInstance().getDataAsyncService(header, 
resultHandler, "Get node list");
     if (service != null) {
       service.getNodeList(header, path, nodeLevel, resultHandler);
     }
@@ -504,10 +316,10 @@ public class DataGroupServiceImpls
 
   @Override
   public void getChildNodeInNextLevel(
-      RaftNode header, String path, AsyncMethodCallback<Set<String>> 
resultHandler)
-      throws TException {
+      RaftNode header, String path, AsyncMethodCallback<Set<String>> 
resultHandler) {
     DataAsyncService service =
-        getDataAsyncService(header, resultHandler, "Get child node in next 
level");
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(header, resultHandler, "Get child node in 
next level");
     if (service != null) {
       service.getChildNodeInNextLevel(header, path, resultHandler);
     }
@@ -517,7 +329,8 @@ public class DataGroupServiceImpls
   public void getChildNodePathInNextLevel(
       RaftNode header, String path, AsyncMethodCallback<Set<String>> 
resultHandler) {
     DataAsyncService service =
-        getDataAsyncService(header, resultHandler, "Get child node path in 
next level");
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(header, resultHandler, "Get child node path 
in next level");
     if (service != null) {
       service.getChildNodePathInNextLevel(header, path, resultHandler);
     }
@@ -527,7 +340,8 @@ public class DataGroupServiceImpls
   public void getAllMeasurementSchema(
       RaftNode header, ByteBuffer planBytes, AsyncMethodCallback<ByteBuffer> 
resultHandler) {
     DataAsyncService service =
-        getDataAsyncService(header, resultHandler, "Get all measurement 
schema");
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(header, resultHandler, "Get all measurement 
schema");
     if (service != null) {
       service.getAllMeasurementSchema(header, planBytes, resultHandler);
     }
@@ -536,7 +350,9 @@ public class DataGroupServiceImpls
   @Override
   public void getAggrResult(
       GetAggrResultRequest request, AsyncMethodCallback<List<ByteBuffer>> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, request);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, request);
     if (service != null) {
       service.getAggrResult(request, resultHandler);
     }
@@ -548,7 +364,8 @@ public class DataGroupServiceImpls
       List<String> timeseriesList,
       AsyncMethodCallback<List<String>> resultHandler) {
     DataAsyncService service =
-        getDataAsyncService(header, resultHandler, "Check if measurements are 
registered");
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(header, resultHandler, "Check if measurements 
are registered");
     if (service != null) {
       service.getUnregisteredTimeseries(header, timeseriesList, resultHandler);
     }
@@ -556,7 +373,9 @@ public class DataGroupServiceImpls
 
   @Override
   public void getGroupByExecutor(GroupByRequest request, 
AsyncMethodCallback<Long> resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, request);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, request);
     if (service != null) {
       service.getGroupByExecutor(request, resultHandler);
     }
@@ -569,260 +388,29 @@ public class DataGroupServiceImpls
       long startTime,
       long endTime,
       AsyncMethodCallback<List<ByteBuffer>> resultHandler) {
-    DataAsyncService service = getDataAsyncService(header, resultHandler, 
"Fetch group by");
+    DataAsyncService service =
+        DataGroupEngine.getInstance().getDataAsyncService(header, 
resultHandler, "Fetch group by");
     if (service != null) {
       service.getGroupByResult(header, executorId, startTime, endTime, 
resultHandler);
     }
   }
 
-  public void preAddNodeForDataGroup(AddNodeLog log, DataGroupMember 
targetDataGroupMember) {
-
-    // Make sure the previous add/remove node log has applied
-    metaGroupMember.syncLocalApply(log.getMetaLogIndex() - 1, false);
-
-    // Check the validity of the partition table
-    if 
(!metaGroupMember.getPartitionTable().deserialize(log.getPartitionTable())) {
-      return;
-    }
-
-    targetDataGroupMember.preAddNode(log.getNewNode());
-  }
-
-  /**
-   * Try adding the node into the group of each DataGroupMember, and if the 
DataGroupMember no
-   * longer stays in that group, also remove and stop it. If the new group 
contains this node, also
-   * create and add a new DataGroupMember for it.
-   *
-   * @param node
-   * @param result
-   */
-  public void addNode(Node node, NodeAdditionResult result) {
-    // If the node executed adding itself to the cluster, it's unnecessary to 
add new groups because
-    // they already exist.
-    if (node.equals(thisNode)) {
-      return;
-    }
-    Iterator<Entry<RaftNode, DataGroupMember>> entryIterator = 
headerGroupMap.entrySet().iterator();
-    synchronized (headerGroupMap) {
-      while (entryIterator.hasNext()) {
-        Entry<RaftNode, DataGroupMember> entry = entryIterator.next();
-        DataGroupMember dataGroupMember = entry.getValue();
-        // the member may be extruded from the group, remove and stop it if so
-        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();
-        }
-      }
-
-      if (logger.isDebugEnabled()) {
-        logger.debug(
-            "Data cluster server: start to handle new groups when adding new 
node {}", node);
-      }
-      for (PartitionGroup newGroup : result.getNewGroupList()) {
-        if (newGroup.contains(thisNode)) {
-          RaftNode header = newGroup.getHeader();
-          logger.info("Adding this node into a new group {}", newGroup);
-          DataGroupMember dataGroupMember = dataMemberFactory.create(newGroup);
-          dataGroupMember = addDataGroupMember(dataGroupMember, header);
-          dataGroupMember.pullNodeAdditionSnapshots(
-              ((SlotPartitionTable) partitionTable).getNodeSlots(header), 
node);
-        }
-      }
-    }
-  }
-
-  /**
-   * When the node joins a cluster, it also creates a new data group and a 
corresponding member
-   * which has no data. This is to make that member pull data from other nodes.
-   */
-  public void pullSnapshots() {
-    for (int raftId = 0;
-        raftId < 
ClusterDescriptor.getInstance().getConfig().getMultiRaftFactor();
-        raftId++) {
-      RaftNode raftNode = new RaftNode(thisNode, raftId);
-      List<Integer> slots = ((SlotPartitionTable) 
partitionTable).getNodeSlots(raftNode);
-      DataGroupMember dataGroupMember = headerGroupMap.get(raftNode);
-      dataGroupMember.pullNodeAdditionSnapshots(slots, thisNode);
-    }
-  }
-
-  /**
-   * Make sure the group will not receive new raft logs
-   *
-   * @param header
-   * @param dataGroupMember
-   */
-  private void removeMember(
-      RaftNode header, DataGroupMember dataGroupMember, boolean removedGroup) {
-    dataGroupMember.setReadOnly();
-    if (!removedGroup) {
-      dataGroupMember.stop();
-    } else {
-      if (dataGroupMember.getCharacter() != NodeCharacter.LEADER) {
-        new Thread(
-                () -> {
-                  try {
-                    dataGroupMember.syncLeader(null);
-                    dataGroupMember.stop();
-                  } catch (CheckConsistencyException e) {
-                    logger.warn("Failed to check consistency.", e);
-                  }
-                })
-            .start();
-      }
-    }
-    stoppedMemberManager.put(header, dataGroupMember);
-    logger.info(
-        "Data group member has removed, header {}, group is {}.",
-        header,
-        dataGroupMember.getAllNodes());
-  }
-
-  /**
-   * Set the partition table as the in-use one and build a DataGroupMember for 
each local group (the
-   * group which the local node is in) and start them.
-   *
-   * @param partitionTable
-   * @throws TTransportException
-   */
-  @SuppressWarnings("java:S1135")
-  public void buildDataGroupMembers(PartitionTable partitionTable) {
-    setPartitionTable(partitionTable);
-    // TODO-Cluster: if there are unchanged members, do not stop and restart 
them
-    // clear previous members if the partition table is reloaded
-    for (DataGroupMember value : headerGroupMap.values()) {
-      value.stop();
-    }
-
-    for (DataGroupMember value : headerGroupMap.values()) {
-      value.setUnchanged(false);
-    }
-
-    List<PartitionGroup> partitionGroups = partitionTable.getLocalGroups();
-    for (PartitionGroup partitionGroup : partitionGroups) {
-      RaftNode header = partitionGroup.getHeader();
-      DataGroupMember prevMember = headerGroupMap.get(header);
-      if (prevMember == null || 
!prevMember.getAllNodes().equals(partitionGroup)) {
-        logger.info("Building member of data group: {}", partitionGroup);
-        // no previous member or member changed
-        DataGroupMember dataGroupMember = 
dataMemberFactory.create(partitionGroup);
-        // the previous member will be replaced here
-        addDataGroupMember(dataGroupMember, header);
-        dataGroupMember.setUnchanged(true);
-      } else {
-        prevMember.setUnchanged(true);
-        prevMember.start();
-        // TODO do we nedd call other functions in addDataGroupMember() ?
-      }
-    }
-
-    // remove out-dated members of this node
-    headerGroupMap.entrySet().removeIf(e -> !e.getValue().isUnchanged());
-
-    logger.info("Data group members are ready");
-  }
-
-  public void preRemoveNodeForDataGroup(RemoveNodeLog log, DataGroupMember 
targetDataGroupMember) {
-
-    // Make sure the previous add/remove node log has applied
-    metaGroupMember.syncLocalApply(log.getMetaLogIndex() - 1, false);
-
-    // Check the validity of the partition table
-    if 
(!metaGroupMember.getPartitionTable().deserialize(log.getPartitionTable())) {
-      return;
-    }
-
-    logger.debug(
-        "Pre removing a node {} from {}",
-        log.getRemovedNode(),
-        targetDataGroupMember.getAllNodes());
-    targetDataGroupMember.preRemoveNode(log.getRemovedNode());
-  }
-
-  /**
-   * Try removing a node from the groups of each DataGroupMember. If the node 
is the header of some
-   * group, set the member to read only so that it can still provide data for 
other nodes that has
-   * not yet pulled its data. Otherwise, just change the node list of the 
member and pull new data.
-   * And create a new DataGroupMember if this node should join a new group 
because of this removal.
-   *
-   * @param node
-   * @param removalResult cluster changes due to the node removal
-   */
-  public void removeNode(Node node, NodeRemovalResult removalResult) {
-    Iterator<Entry<RaftNode, DataGroupMember>> entryIterator = 
headerGroupMap.entrySet().iterator();
-    synchronized (headerGroupMap) {
-      while (entryIterator.hasNext()) {
-        Entry<RaftNode, DataGroupMember> entry = entryIterator.next();
-        DataGroupMember dataGroupMember = entry.getValue();
-        if (dataGroupMember.getHeader().getNode().equals(node) || 
node.equals(thisNode)) {
-          entryIterator.remove();
-          removeMember(
-              entry.getKey(), dataGroupMember, 
dataGroupMember.getHeader().getNode().equals(node));
-        } else {
-          // the group should be updated
-          dataGroupMember.removeNode(node);
-        }
-      }
-
-      if (logger.isDebugEnabled()) {
-        logger.debug(
-            "Data cluster server: start to handle new groups and pulling data 
when removing node {}",
-            node);
-      }
-      // if the removed group contains the local node, the local node should 
join a new group to
-      // preserve the replication number
-      for (PartitionGroup group : partitionTable.getLocalGroups()) {
-        RaftNode header = group.getHeader();
-        if (!headerGroupMap.containsKey(header)) {
-          logger.info("{} should join a new group {}", thisNode, group);
-          DataGroupMember dataGroupMember = dataMemberFactory.create(group);
-          addDataGroupMember(dataGroupMember, header);
-        }
-        // pull new slots from the removed node
-        headerGroupMap.get(header).pullSlots(removalResult);
-      }
-    }
-  }
-
-  public void setPartitionTable(PartitionTable partitionTable) {
-    this.partitionTable = partitionTable;
-  }
-
-  /** @return The reports of every DataGroupMember in this node. */
-  public List<DataMemberReport> genMemberReports() {
-    List<DataMemberReport> dataMemberReports = new ArrayList<>();
-    for (DataGroupMember value : headerGroupMap.values()) {
-
-      dataMemberReports.add(value.genReport());
-    }
-    return dataMemberReports;
-  }
-
   @Override
   public void previousFill(
       PreviousFillRequest request, AsyncMethodCallback<ByteBuffer> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, request);
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, request);
     if (service != null) {
       service.previousFill(request, resultHandler);
     }
   }
 
-  public void closeLogManagers() {
-    for (DataGroupMember member : headerGroupMap.values()) {
-      member.closeLogManager();
-    }
-  }
-
-  public Map<RaftNode, DataGroupMember> getHeaderGroupMap() {
-    return headerGroupMap;
-  }
-
   @Override
   public void matchTerm(
       long index, long term, RaftNode header, AsyncMethodCallback<Boolean> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(header, resultHandler, 
"Match term");
+    DataAsyncService service =
+        DataGroupEngine.getInstance().getDataAsyncService(header, 
resultHandler, "Match term");
     if (service != null) {
       service.matchTerm(index, term, header, resultHandler);
     }
@@ -830,7 +418,9 @@ public class DataGroupServiceImpls
 
   @Override
   public void last(LastQueryRequest request, AsyncMethodCallback<ByteBuffer> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(request.getHeader(), 
resultHandler, "last");
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(request.getHeader(), resultHandler, "last");
     if (service != null) {
       service.last(request, resultHandler);
     }
@@ -842,7 +432,8 @@ public class DataGroupServiceImpls
       List<String> pathsToQuery,
       int level,
       AsyncMethodCallback<Integer> resultHandler) {
-    DataAsyncService service = getDataAsyncService(header, resultHandler, 
"count path");
+    DataAsyncService service =
+        DataGroupEngine.getInstance().getDataAsyncService(header, 
resultHandler, "count path");
     if (service != null) {
       service.getPathCount(header, pathsToQuery, level, resultHandler);
     }
@@ -851,7 +442,9 @@ public class DataGroupServiceImpls
   @Override
   public void onSnapshotApplied(
       RaftNode header, List<Integer> slots, AsyncMethodCallback<Boolean> 
resultHandler) {
-    DataAsyncService service = getDataAsyncService(header, resultHandler, 
"Snapshot applied");
+    DataAsyncService service =
+        DataGroupEngine.getInstance()
+            .getDataAsyncService(header, resultHandler, "Snapshot applied");
     if (service != null) {
       service.onSnapshotApplied(header, slots, resultHandler);
     }
@@ -859,168 +452,220 @@ public class DataGroupServiceImpls
 
   @Override
   public long querySingleSeries(SingleSeriesQueryRequest request) throws 
TException {
-    return getDataSyncService(request.getHeader()).querySingleSeries(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .querySingleSeries(request);
   }
 
   @Override
   public long queryMultSeries(MultSeriesQueryRequest request) throws 
TException {
-    return getDataSyncService(request.getHeader()).queryMultSeries(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .queryMultSeries(request);
   }
 
   @Override
   public ByteBuffer fetchSingleSeries(RaftNode header, long readerId) throws 
TException {
-    return getDataSyncService(header).fetchSingleSeries(header, readerId);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .fetchSingleSeries(header, readerId);
   }
 
   @Override
   public Map<String, ByteBuffer> fetchMultSeries(RaftNode header, long 
readerId, List<String> paths)
       throws TException {
-    return getDataSyncService(header).fetchMultSeries(header, readerId, paths);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .fetchMultSeries(header, readerId, paths);
   }
 
   @Override
   public long querySingleSeriesByTimestamp(SingleSeriesQueryRequest request) 
throws TException {
-    return 
getDataSyncService(request.getHeader()).querySingleSeriesByTimestamp(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .querySingleSeriesByTimestamp(request);
   }
 
   @Override
   public ByteBuffer fetchSingleSeriesByTimestamps(
       RaftNode header, long readerId, List<Long> timestamps) throws TException 
{
-    return getDataSyncService(header).fetchSingleSeriesByTimestamps(header, 
readerId, timestamps);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .fetchSingleSeriesByTimestamps(header, readerId, timestamps);
   }
 
   @Override
   public void endQuery(RaftNode header, Node thisNode, long queryId) throws 
TException {
-    getDataSyncService(header).endQuery(header, thisNode, queryId);
+    DataGroupEngine.getInstance().getDataSyncService(header).endQuery(header, 
thisNode, queryId);
   }
 
   @Override
   public GetAllPathsResult getAllPaths(RaftNode header, List<String> path, 
boolean withAlias)
       throws TException {
-    return getDataSyncService(header).getAllPaths(header, path, withAlias);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .getAllPaths(header, path, withAlias);
   }
 
   @Override
   public Set<String> getAllDevices(RaftNode header, List<String> path) throws 
TException {
-    return getDataSyncService(header).getAllDevices(header, path);
+    return 
DataGroupEngine.getInstance().getDataSyncService(header).getAllDevices(header, 
path);
   }
 
   @Override
   public List<String> getNodeList(RaftNode header, String path, int nodeLevel) 
throws TException {
-    return getDataSyncService(header).getNodeList(header, path, nodeLevel);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .getNodeList(header, path, nodeLevel);
   }
 
   @Override
   public Set<String> getChildNodeInNextLevel(RaftNode header, String path) 
throws TException {
-    return getDataSyncService(header).getChildNodeInNextLevel(header, path);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .getChildNodeInNextLevel(header, path);
   }
 
   @Override
   public Set<String> getChildNodePathInNextLevel(RaftNode header, String path) 
throws TException {
-    return getDataSyncService(header).getChildNodePathInNextLevel(header, 
path);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .getChildNodePathInNextLevel(header, path);
   }
 
   @Override
   public ByteBuffer getAllMeasurementSchema(RaftNode header, ByteBuffer 
planBinary)
       throws TException {
-    return getDataSyncService(header).getAllMeasurementSchema(header, 
planBinary);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .getAllMeasurementSchema(header, planBinary);
   }
 
   @Override
   public ByteBuffer getDevices(RaftNode header, ByteBuffer planBinary) throws 
TException {
-    return getDataSyncService(header).getDevices(header, planBinary);
+    return 
DataGroupEngine.getInstance().getDataSyncService(header).getDevices(header, 
planBinary);
   }
 
   @Override
   public List<ByteBuffer> getAggrResult(GetAggrResultRequest request) throws 
TException {
-    return getDataSyncService(request.getHeader()).getAggrResult(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .getAggrResult(request);
   }
 
   @Override
   public List<String> getUnregisteredTimeseries(RaftNode header, List<String> 
timeseriesList)
       throws TException {
-    return getDataSyncService(header).getUnregisteredTimeseries(header, 
timeseriesList);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .getUnregisteredTimeseries(header, timeseriesList);
   }
 
   @Override
   public PullSnapshotResp pullSnapshot(PullSnapshotRequest request) throws 
TException {
-    return getDataSyncService(request.getHeader()).pullSnapshot(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .pullSnapshot(request);
   }
 
   @Override
   public long getGroupByExecutor(GroupByRequest request) throws TException {
-    return getDataSyncService(request.getHeader()).getGroupByExecutor(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .getGroupByExecutor(request);
   }
 
   @Override
   public List<ByteBuffer> getGroupByResult(
       RaftNode header, long executorId, long startTime, long endTime) throws 
TException {
-    return getDataSyncService(header).getGroupByResult(header, executorId, 
startTime, endTime);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .getGroupByResult(header, executorId, startTime, endTime);
   }
 
   @Override
   public PullSchemaResp pullTimeSeriesSchema(PullSchemaRequest request) throws 
TException {
-    return 
getDataSyncService(request.getHeader()).pullTimeSeriesSchema(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .pullTimeSeriesSchema(request);
   }
 
   @Override
   public PullSchemaResp pullMeasurementSchema(PullSchemaRequest request) 
throws TException {
-    return 
getDataSyncService(request.getHeader()).pullMeasurementSchema(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .pullMeasurementSchema(request);
   }
 
   @Override
   public ByteBuffer previousFill(PreviousFillRequest request) throws 
TException {
-    return getDataSyncService(request.getHeader()).previousFill(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .previousFill(request);
   }
 
   @Override
   public ByteBuffer last(LastQueryRequest request) throws TException {
-    return getDataSyncService(request.getHeader()).last(request);
+    return 
DataGroupEngine.getInstance().getDataSyncService(request.getHeader()).last(request);
   }
 
   @Override
   public int getPathCount(RaftNode header, List<String> pathsToQuery, int 
level) throws TException {
-    return getDataSyncService(header).getPathCount(header, pathsToQuery, 
level);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .getPathCount(header, pathsToQuery, level);
   }
 
   @Override
   public boolean onSnapshotApplied(RaftNode header, List<Integer> slots) {
-    return getDataSyncService(header).onSnapshotApplied(header, slots);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .onSnapshotApplied(header, slots);
   }
 
   @Override
   public HeartBeatResponse sendHeartbeat(HeartBeatRequest request) {
-    return getDataSyncService(request.getHeader()).sendHeartbeat(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .sendHeartbeat(request);
   }
 
   @Override
   public long startElection(ElectionRequest request) {
-    return getDataSyncService(request.getHeader()).startElection(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .startElection(request);
   }
 
   @Override
   public long appendEntries(AppendEntriesRequest request) throws TException {
-    return getDataSyncService(request.getHeader()).appendEntries(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .appendEntries(request);
   }
 
   @Override
   public long appendEntry(AppendEntryRequest request) throws TException {
-    return getDataSyncService(request.getHeader()).appendEntry(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .appendEntry(request);
   }
 
   @Override
   public void sendSnapshot(SendSnapshotRequest request) throws TException {
-    getDataSyncService(request.getHeader()).sendSnapshot(request);
+    
DataGroupEngine.getInstance().getDataSyncService(request.getHeader()).sendSnapshot(request);
   }
 
   @Override
   public TSStatus executeNonQueryPlan(ExecutNonQueryReq request) throws 
TException {
-    return 
getDataSyncService(request.getHeader()).executeNonQueryPlan(request);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(request.getHeader())
+        .executeNonQueryPlan(request);
   }
 
   @Override
   public RequestCommitIndexResponse requestCommitIndex(RaftNode header) throws 
TException {
-    return getDataSyncService(header).requestCommitIndex(header);
+    return 
DataGroupEngine.getInstance().getDataSyncService(header).requestCommitIndex(header);
   }
 
   @Override
@@ -1034,13 +679,15 @@ public class DataGroupServiceImpls
 
   @Override
   public boolean matchTerm(long index, long term, RaftNode header) {
-    return getDataSyncService(header).matchTerm(index, term, header);
+    return 
DataGroupEngine.getInstance().getDataSyncService(header).matchTerm(index, term, 
header);
   }
 
   @Override
   public ByteBuffer peekNextNotNullValue(
       RaftNode header, long executorId, long startTime, long endTime) throws 
TException {
-    return getDataSyncService(header).peekNextNotNullValue(header, executorId, 
startTime, endTime);
+    return DataGroupEngine.getInstance()
+        .getDataSyncService(header)
+        .peekNextNotNullValue(header, executorId, startTime, endTime);
   }
 
   @Override
@@ -1052,7 +699,9 @@ public class DataGroupServiceImpls
       AsyncMethodCallback<ByteBuffer> resultHandler)
       throws TException {
     resultHandler.onComplete(
-        getDataSyncService(header).peekNextNotNullValue(header, executorId, 
startTime, endTime));
+        DataGroupEngine.getInstance()
+            .getDataSyncService(header)
+            .peekNextNotNullValue(header, executorId, startTime, endTime));
   }
 
   @Override
@@ -1073,24 +722,4 @@ public class DataGroupServiceImpls
       resultHandler.onError(e);
     }
   }
-
-  @Override
-  public String getHeaderGroupMapAsString() {
-    return headerGroupMap.toString();
-  }
-
-  @Override
-  public int getAsyncServiceMapSize() {
-    return asyncServiceMap.size();
-  }
-
-  @Override
-  public int getSyncServiceMapSize() {
-    return syncServiceMap.size();
-  }
-
-  @Override
-  public String getPartitionTable() {
-    return partitionTable.toString();
-  }
 }
diff --git 
a/cluster/src/test/java/org/apache/iotdb/cluster/server/member/MetaGroupMemberTest.java
 
b/cluster/src/test/java/org/apache/iotdb/cluster/server/member/MetaGroupMemberTest.java
index 54ee250..dddb0fa 100644
--- 
a/cluster/src/test/java/org/apache/iotdb/cluster/server/member/MetaGroupMemberTest.java
+++ 
b/cluster/src/test/java/org/apache/iotdb/cluster/server/member/MetaGroupMemberTest.java
@@ -68,7 +68,7 @@ import org.apache.iotdb.cluster.server.NodeCharacter;
 import org.apache.iotdb.cluster.server.Response;
 import org.apache.iotdb.cluster.server.handlers.caller.GenericHandler;
 import org.apache.iotdb.cluster.server.monitor.NodeStatusManager;
-import org.apache.iotdb.cluster.server.service.DataGroupServiceImpls;
+import org.apache.iotdb.cluster.server.service.DataGroupEngine;
 import org.apache.iotdb.cluster.server.service.MetaAsyncService;
 import org.apache.iotdb.cluster.utils.ClusterUtils;
 import org.apache.iotdb.cluster.utils.Constants;
@@ -132,13 +132,20 @@ import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicReference;
 
-import static org.apache.iotdb.cluster.server.NodeCharacter.*;
+import static org.apache.iotdb.cluster.server.NodeCharacter.ELECTOR;
+import static org.apache.iotdb.cluster.server.NodeCharacter.FOLLOWER;
+import static org.apache.iotdb.cluster.server.NodeCharacter.LEADER;
 import static org.awaitility.Awaitility.await;
-import static org.junit.Assert.*;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertTrue;
+import static org.junit.Assert.fail;
 
 public class MetaGroupMemberTest extends BaseMember {
 
-  private DataGroupServiceImpls dataGroupServiceImpls;
+  private DataGroupEngine dataGroupEngine;
   protected boolean mockDataClusterServer;
   private Node exiledNode;
 
@@ -148,7 +155,7 @@ public class MetaGroupMemberTest extends BaseMember {
   @Override
   @After
   public void tearDown() throws Exception {
-    dataGroupServiceImpls.stop();
+    dataGroupEngine.stop();
     super.tearDown();
     
ClusterDescriptor.getInstance().getConfig().setReplicationNum(prevReplicaNum);
     ClusterDescriptor.getInstance().getConfig().setSeedNodeUrls(prevSeedNodes);
@@ -171,8 +178,8 @@ public class MetaGroupMemberTest extends BaseMember {
     dummyResponse.set(Response.RESPONSE_AGREE);
     testMetaMember.setAllNodes(allNodes);
 
-    dataGroupServiceImpls =
-        new DataGroupServiceImpls(
+    dataGroupEngine =
+        new DataGroupEngine(
             new DataGroupMember.Factory(null, testMetaMember) {
               @Override
               public DataGroupMember create(PartitionGroup partitionGroup) {
@@ -181,7 +188,7 @@ public class MetaGroupMemberTest extends BaseMember {
             },
             testMetaMember);
 
-    buildDataGroups(dataGroupServiceImpls);
+    buildDataGroups(dataGroupEngine);
     testMetaMember.getThisNode().setNodeIdentifier(0);
     testMetaMember.setRouter(new 
ClusterPlanRouter(testMetaMember.getPartitionTable()));
     mockDataClusterServer = false;
@@ -327,9 +334,9 @@ public class MetaGroupMemberTest extends BaseMember {
           }
 
           @Override
-          public DataGroupServiceImpls getDataGroupEngine() {
+          public DataGroupEngine getDataGroupEngine() {
             return mockDataClusterServer
-                ? MetaGroupMemberTest.this.dataGroupServiceImpls
+                ? MetaGroupMemberTest.this.dataGroupEngine
                 : ClusterIoTDB.getInstance().getDataGroupEngine();
           }
 
@@ -544,7 +551,7 @@ public class MetaGroupMemberTest extends BaseMember {
     return metaGroupMember;
   }
 
-  private void buildDataGroups(DataGroupServiceImpls dataGroupServiceImpls) {
+  private void buildDataGroups(DataGroupEngine dataGroupServiceImpls) {
     List<PartitionGroup> partitionGroups = partitionTable.getLocalGroups();
 
     dataGroupServiceImpls.setPartitionTable(partitionTable);

Reply via email to