This is an automated email from the ASF dual-hosted git repository.
jiangtian pushed a commit to branch expr_catch_up
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/expr_catch_up by this push:
new 84d66d9dc9 add PlanBasedStateMachine
84d66d9dc9 is described below
commit 84d66d9dc901353b39536940f187a13111fb018a
Author: jt <[email protected]>
AuthorDate: Thu Jun 16 11:14:08 2022 +0800
add PlanBasedStateMachine
---
.../org/apache/iotdb/cluster/ClusterIoTDB.java | 5 +-
.../cluster/impl/NativeSingleRaftConsensus.java | 2 +-
.../iotdb/cluster/impl/PlanBasedStateMachine.java | 99 ++++++++++++++++++++++
.../cluster/server/member/DataGroupMember.java | 12 +--
.../iotdb/cluster/server/member/RaftMember.java | 5 ++
.../cluster/server/service/DataGroupEngine.java | 4 +-
6 files changed, 115 insertions(+), 12 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 441ec1f298..5a00c181e7 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
@@ -32,6 +32,7 @@ import org.apache.iotdb.cluster.config.ClusterDescriptor;
import org.apache.iotdb.cluster.coordinator.Coordinator;
import org.apache.iotdb.cluster.exception.ConfigInconsistentException;
import org.apache.iotdb.cluster.exception.StartUpCheckFailureException;
+import org.apache.iotdb.cluster.impl.PlanBasedStateMachine;
import org.apache.iotdb.cluster.metadata.CSchemaProcessor;
import org.apache.iotdb.cluster.metadata.MetaPuller;
import org.apache.iotdb.cluster.partition.slot.SlotPartitionTable;
@@ -158,7 +159,9 @@ public class ClusterIoTDB implements ClusterIoTDBMBean {
TProtocolFactory protocolFactory =
ThriftServiceThread.getProtocolFactory(
IoTDBDescriptor.getInstance().getConfig().isRpcThriftCompressionEnable());
- metaGroupMember = new MetaGroupMember(thisNode, coordinator);
+ PlanBasedStateMachine stateMachine = new PlanBasedStateMachine();
+ metaGroupMember = new MetaGroupMember(thisNode, coordinator, stateMachine);
+ stateMachine.setMetaGroupMember(metaGroupMember);
IoTDB.setClusterMode();
IoTDB.setSchemaProcessor(CSchemaProcessor.getInstance());
((CSchemaProcessor)
IoTDB.schemaProcessor).setMetaGroupMember(metaGroupMember);
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/impl/NativeSingleRaftConsensus.java
b/cluster/src/main/java/org/apache/iotdb/cluster/impl/NativeSingleRaftConsensus.java
index 622287d7a9..9d165c6afd 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/impl/NativeSingleRaftConsensus.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/impl/NativeSingleRaftConsensus.java
@@ -51,7 +51,7 @@ public class NativeSingleRaftConsensus implements
ISingleConsensus {
@Override
public ConsensusReadResponse read(IConsensusRequest request) {
- return null;
+ return new ConsensusReadResponse(null, raftMember.read(request));
}
@Override
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/impl/PlanBasedStateMachine.java
b/cluster/src/main/java/org/apache/iotdb/cluster/impl/PlanBasedStateMachine.java
new file mode 100644
index 0000000000..b98652e5de
--- /dev/null
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/impl/PlanBasedStateMachine.java
@@ -0,0 +1,99 @@
+/*
+ * 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.impl;
+
+import java.io.File;
+import org.apache.iotdb.cluster.query.ClusterPlanExecutor;
+import org.apache.iotdb.cluster.server.member.MetaGroupMember;
+import org.apache.iotdb.cluster.utils.StatusUtils;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.consensus.IStateMachine;
+import org.apache.iotdb.consensus.common.DataSet;
+import org.apache.iotdb.consensus.common.request.IConsensusRequest;
+import org.apache.iotdb.db.exception.StorageEngineException;
+import org.apache.iotdb.db.exception.metadata.StorageGroupNotSetException;
+import org.apache.iotdb.db.exception.query.QueryProcessException;
+import org.apache.iotdb.db.qp.physical.PhysicalPlan;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class PlanBasedStateMachine implements IStateMachine {
+
+ private static final Logger logger =
LoggerFactory.getLogger(PlanBasedStateMachine.class);
+
+ private ClusterPlanExecutor planExecutor;
+ private MetaGroupMember metaGroupMember;
+
+ public PlanBasedStateMachine() {
+ }
+
+ public PlanBasedStateMachine(MetaGroupMember metaGroupMember) {
+ this.metaGroupMember = metaGroupMember;
+ }
+
+ @Override
+ public void start() {
+ try {
+ planExecutor = new ClusterPlanExecutor(metaGroupMember);
+ } catch (QueryProcessException e) {
+ logger.error("Cannot initialize plan executor", e);
+ }
+ }
+
+ @Override
+ public void stop() {
+
+ }
+
+ @Override
+ public TSStatus write(IConsensusRequest request) {
+ if (!(request instanceof PhysicalPlan)) {
+ return StatusUtils.getStatus(StatusUtils.EXECUTE_STATEMENT_ERROR,
+ "Not supported request: " + request);
+ }
+ try {
+ planExecutor.processNonQuery(((PhysicalPlan) request));
+ return StatusUtils.OK;
+ } catch (QueryProcessException | StorageGroupNotSetException |
StorageEngineException e) {
+ logger.warn("Plan execution error", e);
+ return StatusUtils.getStatus(StatusUtils.EXECUTE_STATEMENT_ERROR,
+ e.getMessage());
+ }
+ }
+
+ @Override
+ public DataSet read(IConsensusRequest request) {
+ return null;
+ }
+
+ @Override
+ public boolean takeSnapshot(File snapshotDir) {
+ return false;
+ }
+
+ @Override
+ public void loadSnapshot(File latestSnapshotRootDir) {
+
+ }
+
+ public void setMetaGroupMember(MetaGroupMember metaGroupMember) {
+ this.metaGroupMember = metaGroupMember;
+ }
+}
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
index 7c0647b24d..669979bd24 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/DataGroupMember.java
@@ -26,15 +26,13 @@ import org.apache.iotdb.cluster.config.ClusterDescriptor;
import org.apache.iotdb.cluster.exception.CheckConsistencyException;
import org.apache.iotdb.cluster.exception.SnapshotInstallationException;
import org.apache.iotdb.cluster.exception.UnknownLogTypeException;
+import org.apache.iotdb.cluster.impl.PlanBasedStateMachine;
import org.apache.iotdb.cluster.log.IndirectLogDispatcher;
import org.apache.iotdb.cluster.log.Log;
-import org.apache.iotdb.cluster.log.LogApplier;
import org.apache.iotdb.cluster.log.LogParser;
import org.apache.iotdb.cluster.log.Snapshot;
import org.apache.iotdb.cluster.log.appender.BlockingLogAppender;
import org.apache.iotdb.cluster.log.appender.SlidingWindowLogAppender;
-import org.apache.iotdb.cluster.log.applier.AsyncDataLogApplier;
-import org.apache.iotdb.cluster.log.applier.DataLogApplier;
import org.apache.iotdb.cluster.log.logtypes.AddNodeLog;
import org.apache.iotdb.cluster.log.logtypes.RemoveNodeLog;
import org.apache.iotdb.cluster.log.manage.FilePartitionedSnapshotLogManager;
@@ -72,12 +70,10 @@ import org.apache.iotdb.cluster.server.monitor.Timer;
import org.apache.iotdb.cluster.server.monitor.Timer.Statistic;
import org.apache.iotdb.cluster.utils.IOUtils;
import org.apache.iotdb.cluster.utils.StatusUtils;
-import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
import org.apache.iotdb.common.rpc.thrift.TEndPoint;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
import org.apache.iotdb.commons.conf.IoTDBConstant;
-import org.apache.iotdb.commons.consensus.ConsensusGroupId;
import org.apache.iotdb.commons.consensus.DataRegionId;
import org.apache.iotdb.commons.exception.IllegalPathException;
import org.apache.iotdb.commons.exception.MetadataException;
@@ -109,7 +105,6 @@ import org.apache.iotdb.db.service.IoTDB;
import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.iotdb.tsfile.utils.Pair;
-import org.apache.thrift.protocol.TProtocolFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -348,8 +343,9 @@ public class DataGroupMember extends RaftMember implements
DataGroupMemberMBean
this.metaGroupMember = metaGroupMember;
}
- public DataGroupMember create(Node thisNode, PartitionGroup
partitionGroup, IStateMachine stateMachine) {
- return new DataGroupMember(thisNode, partitionGroup, metaGroupMember,
stateMachine);
+ public DataGroupMember create(Node thisNode, PartitionGroup
partitionGroup) {
+ return new DataGroupMember(thisNode, partitionGroup, metaGroupMember,
+ new PlanBasedStateMachine(metaGroupMember));
}
}
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
index 51f4bb78c1..b2c4d92219 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java
@@ -90,6 +90,7 @@ import org.apache.iotdb.commons.consensus.ConsensusGroupId;
import org.apache.iotdb.commons.exception.IoTDBException;
import org.apache.iotdb.commons.utils.TestOnly;
import org.apache.iotdb.consensus.IStateMachine;
+import org.apache.iotdb.consensus.common.DataSet;
import org.apache.iotdb.consensus.common.Peer;
import org.apache.iotdb.consensus.common.request.IConsensusRequest;
import org.apache.iotdb.consensus.common.response.ConsensusWriteResponse;
@@ -863,6 +864,10 @@ public abstract class RaftMember implements
RaftMemberMBean {
*/
public abstract ConsensusWriteResponse executeRequest(IConsensusRequest
plan);
+ public DataSet read(IConsensusRequest request) {
+ return stateMachine.read(request);
+ }
+
abstract ClientCategory getClientCategory();
/**
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
index 574528df24..c8ceb12334 100644
---
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
@@ -75,7 +75,7 @@ public class DataGroupEngine implements IService,
DataGroupEngineMBean {
private static TProtocolFactory protocolFactory;
private DataGroupEngine() {
- dataMemberFactory = new DataGroupMember.Factory(protocolFactory,
metaGroupMember);
+ dataMemberFactory = new DataGroupMember.Factory(metaGroupMember);
stoppedMemberManager = new StoppedMemberManager(dataMemberFactory,
thisNode);
}
@@ -88,7 +88,7 @@ public class DataGroupEngine implements IService,
DataGroupEngineMBean {
@TestOnly
public void resetFactory() {
- dataMemberFactory = new DataGroupMember.Factory(protocolFactory,
metaGroupMember);
+ dataMemberFactory = new DataGroupMember.Factory(metaGroupMember);
}
@TestOnly