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