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

Reply via email to