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

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


The following commit(s) were added to refs/heads/cluster by this push:
     new 4cb1df3  add change consistency level
4cb1df3 is described below

commit 4cb1df35225e46d905d6fdfe89ea3301008dc4d0
Author: lta <[email protected]>
AuthorDate: Thu Apr 4 17:28:01 2019 +0800

    add change consistency level
---
 .../apache/iotdb/cluster/callback/BatchQPTask.java |  12 ++-
 .../iotdb/cluster/callback/SingleQPTask.java       |   2 +-
 .../apache/iotdb/cluster/config/ClusterConfig.java |  17 +++-
 .../iotdb/cluster/config/ClusterConstant.java      |  28 +++++
 .../iotdb/cluster/config/ClusterDescriptor.java    |   5 +-
 .../cluster/entity/raft/DataStateMachine.java      |   4 +-
 .../cluster/entity/raft/MetadataStateManchine.java |   4 +-
 .../apache/iotdb/cluster/qp/ClusterQPExecutor.java |  35 +++++++
 .../cluster/qp/executor/NonQueryExecutor.java      |   4 +-
 .../cluster/qp/executor/QueryMetadataExecutor.java |   4 +-
 .../cluster/rpc/impl/RaftNodeAsClientManager.java  |   2 +-
 .../processor/DataGroupNonQueryAsyncProcessor.java |   2 +-
 .../processor/MetaGroupNonQueryAsyncProcessor.java |   2 +-
 .../processor/QueryTimeSeriesAsyncProcessor.java   |   2 +-
 .../cluster/rpc/service/TSServiceClusterImpl.java  |  30 ++++++
 .../iotdb/cluster/qp/ClusterQPExecutorTest.java    | 113 +++++++++++++++++++++
 .../org/apache/iotdb/db/service/TSServiceImpl.java |  17 ++++
 17 files changed, 263 insertions(+), 20 deletions(-)

diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/callback/BatchQPTask.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/callback/BatchQPTask.java
index 3fffba7..464f630 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/callback/BatchQPTask.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/callback/BatchQPTask.java
@@ -53,6 +53,11 @@ public class BatchQPTask extends QPTask {
   private Map<String, List<Integer>> planIndexMap;
 
   /**
+   * Task thread map
+   */
+  private Map<String, Thread> taskThreadMap;
+
+  /**
    * Batch result
    */
   private int[] batchResult;
@@ -146,7 +151,12 @@ public class BatchQPTask extends QPTask {
 
   @Override
   public void shutdown() {
-    //TODO
+    for(Thread taskThread:taskThreadMap.values()){
+      if(taskThread.isAlive()){
+        taskThread.interrupt();
+      }
+    }
+    this.taskCountDownLatch.countDown();
   }
 
   public boolean isAllSuccessful() {
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/callback/SingleQPTask.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/callback/SingleQPTask.java
index 8c32c70..34c1dfd 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/callback/SingleQPTask.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/callback/SingleQPTask.java
@@ -49,6 +49,6 @@ public class SingleQPTask extends QPTask {
 
   @Override
   public void shutdown() {
-
+    this.taskCountDownLatch.countDown();
   }
 }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java
index 950a596..2adef86 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java
@@ -122,7 +122,7 @@ public class ClusterConfig {
     // empty constructor
   }
 
-  public void updatePath() {
+  public void setDefaultPath() {
     IoTDBConfig conf = IoTDBDescriptor.getInstance().getConfig();
     String iotdbDataDir = conf.getDataDir();
     iotdbDataDir = FilePathUtils.regularizePath(iotdbDataDir);
@@ -130,12 +130,19 @@ public class ClusterConfig {
     this.raftSnapshotPath = raftDir + File.separatorChar + 
DEFAULT_RAFT_SNAPSHOT_DIR;
     this.raftLogPath = raftDir + File.separatorChar + DEFAULT_RAFT_LOG_DIR;
     this.raftMetadataPath = raftDir + File.separatorChar + 
DEFAULT_RAFT_METADATA_DIR;
+  }
+
+  public void createAllPath(){
+    createPath(this.raftSnapshotPath);
+    createPath(this.raftLogPath);
+    createPath(this.raftMetadataPath);
+  }
+
+  private void createPath(String path){
     try {
-      FileUtils.forceMkdir(new File(this.raftSnapshotPath));
-      FileUtils.forceMkdir(new File(this.raftLogPath));
-      FileUtils.forceMkdir(new File(this.raftMetadataPath));
+      FileUtils.forceMkdir(new File(path));
     } catch (IOException e) {
-      LOGGER.warn("Raft dir already exists.");
+      LOGGER.warn("Path {} already exists.", path);
     }
   }
 
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConstant.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConstant.java
new file mode 100644
index 0000000..71547ee
--- /dev/null
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConstant.java
@@ -0,0 +1,28 @@
+/**
+ * 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.config;
+
+public class ClusterConstant {
+
+  public static final String SET_READ_METADATA_CONSISTENCY_LEVEL_PATTERN = 
"set\\s+read\\s+metadata\\s+level\\s+\\d+";
+
+  public static final String SET_READ_DATA_CONSISTENCY_LEVEL_PATTERN = 
"set\\s+read\\s+data\\s+level\\s+\\d+";
+
+  public static final int MAX_CONSISTENCY_LEVEL = 2;
+}
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterDescriptor.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterDescriptor.java
index 276d26a..4a4c6ef 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterDescriptor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterDescriptor.java
@@ -49,7 +49,7 @@ public class ClusterDescriptor {
    * load an property file and set ClusterConfig variables.
    */
   private void loadProps() {
-    conf.updatePath();
+    conf.setDefaultPath();
     InputStream inputStream;
     String url = System.getProperty(IoTDBConstant.IOTDB_CONF, null);
     if (url == null) {
@@ -61,6 +61,7 @@ public class ClusterDescriptor {
             "Cannot find IOTDB_HOME or CLUSTER_CONF environment variable when 
loading "
                 + "config file {}, use default configuration",
             ClusterConfig.CONFIG_NAME);
+        conf.createAllPath();
         return;
       }
     } else {
@@ -71,6 +72,7 @@ public class ClusterDescriptor {
       inputStream = new FileInputStream(new File(url));
     } catch (FileNotFoundException e) {
       LOGGER.warn("Fail to find config file {}", url, e);
+      conf.createAllPath();
       return;
     }
 
@@ -141,6 +143,7 @@ public class ClusterDescriptor {
     } catch (Exception e) {
       LOGGER.warn("Incorrect format in config file, use default 
configuration", e);
     } finally {
+      conf.createAllPath();
       try {
         inputStream.close();
       } catch (IOException e) {
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/entity/raft/DataStateMachine.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/entity/raft/DataStateMachine.java
index da7a431..2bcc152 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/entity/raft/DataStateMachine.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/entity/raft/DataStateMachine.java
@@ -127,7 +127,7 @@ public class DataStateMachine extends StateMachineAdapter {
               nullReadTask.await();
             } catch (InterruptedException e) {
               status.setCode(-1);
-              status.setErrorMsg(e.toString());
+              status.setErrorMsg(e.getMessage());
             }
           }
           qpExecutor.processNonQuery(plan);
@@ -136,7 +136,7 @@ public class DataStateMachine extends StateMachineAdapter {
           }
         } catch (ProcessorException | IOException e) {
           LOGGER.error("Execute physical plan error", e);
-          status = new Status(-1, e.toString());
+          status = new Status(-1, e.getMessage());
           if (closure != null) {
             response.addResult(false);
           }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/entity/raft/MetadataStateManchine.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/entity/raft/MetadataStateManchine.java
index 0c9b43b..57e7436 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/entity/raft/MetadataStateManchine.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/entity/raft/MetadataStateManchine.java
@@ -117,13 +117,13 @@ public class MetadataStateManchine extends 
StateMachineAdapter {
           }
         } catch (IOException | PathErrorException e) {
           LOGGER.error("Execute metadata plan error", e);
-          status = new Status(-1, e.toString());
+          status = new Status(-1, e.getMessage());
           if (closure != null) {
             response.addResult(false);
           }
         } catch (ProcessorException e) {
           LOGGER.error("Execute author plan error", e);
-          status = new Status(-1, e.toString());
+          status = new Status(-1, e.getMessage());
           if (closure != null) {
             response.addResult(false);
           }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/qp/ClusterQPExecutor.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/qp/ClusterQPExecutor.java
index aac5f1f..800c723 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/qp/ClusterQPExecutor.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/qp/ClusterQPExecutor.java
@@ -28,6 +28,7 @@ import java.util.concurrent.atomic.AtomicInteger;
 import org.apache.iotdb.cluster.callback.QPTask;
 import org.apache.iotdb.cluster.callback.QPTask.TaskState;
 import org.apache.iotdb.cluster.config.ClusterConfig;
+import org.apache.iotdb.cluster.config.ClusterConstant;
 import org.apache.iotdb.cluster.config.ClusterDescriptor;
 import org.apache.iotdb.cluster.entity.Server;
 import org.apache.iotdb.cluster.exception.RaftConnectionException;
@@ -45,14 +46,24 @@ import org.slf4j.LoggerFactory;
 public abstract class ClusterQPExecutor {
 
   private static final Logger LOGGER = 
LoggerFactory.getLogger(ClusterQPExecutor.class);
+
   private static final ClusterConfig CLUSTER_CONFIG = 
ClusterDescriptor.getInstance().getConfig();
+
   protected static final String METADATA_GROUP_ID = 
CLUSTER_CONFIG.METADATA_GROUP_ID;
+
+  /**
+   * Raft as client manager.
+   */
   protected static final RaftNodeAsClientManager CLIENT_MANAGER = 
RaftNodeAsClientManager
       .getInstance();
+
   protected Router router = Router.getInstance();
+
   private PhysicalNode localNode = new PhysicalNode(CLUSTER_CONFIG.getIp(),
       CLUSTER_CONFIG.getPort());
+
   protected MManager mManager = MManager.getInstance();
+
   protected final Server server = Server.getInstance();
 
   /**
@@ -253,4 +264,28 @@ public abstract class ClusterQPExecutor {
       currentTask.shutdown();
     }
   }
+
+  public void setReadMetadataConsistencyLevel(int level) throws Exception {
+    if (level <= ClusterConstant.MAX_CONSISTENCY_LEVEL) {
+      this.readMetadataConsistencyLevel = level;
+    } else {
+      throw new Exception(String.format("Consistency level %d not support", 
level));
+    }
+  }
+
+  public void setReadDataConsistencyLevel(int level) throws Exception {
+    if (level <= ClusterConstant.MAX_CONSISTENCY_LEVEL) {
+      this.readDataConsistencyLevel = level;
+    } else {
+      throw new Exception(String.format("Consistency level %d not support", 
level));
+    }
+  }
+
+  public int getReadMetadataConsistencyLevel() {
+    return readMetadataConsistencyLevel;
+  }
+
+  public int getReadDataConsistencyLevel() {
+    return readDataConsistencyLevel;
+  }
 }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/qp/executor/NonQueryExecutor.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/qp/executor/NonQueryExecutor.java
index f14a59d..3b6f9d8 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/qp/executor/NonQueryExecutor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/qp/executor/NonQueryExecutor.java
@@ -212,7 +212,7 @@ public class NonQueryExecutor extends ClusterQPExecutor {
         } catch (Exception e) {
           result[i] = Statement.EXECUTE_FAILED;
           batchResult.setAllSuccessful(false);
-          batchResult.setBatchErrorMessage(e.toString());
+          batchResult.setBatchErrorMessage(e.getMessage());
         }
       }
     }
@@ -229,7 +229,7 @@ public class NonQueryExecutor extends ClusterQPExecutor {
         subTaskMap.put(groupId, singleQPTask);
       } catch (IOException e) {
         batchResult.setAllSuccessful(false);
-        batchResult.setBatchErrorMessage(e.toString());
+        batchResult.setBatchErrorMessage(e.getMessage());
         for (int index : planIndexMap.get(groupId)) {
           result[index] = Statement.EXECUTE_FAILED;
         }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/qp/executor/QueryMetadataExecutor.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/qp/executor/QueryMetadataExecutor.java
index 5c02e61..d9e517a 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/qp/executor/QueryMetadataExecutor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/qp/executor/QueryMetadataExecutor.java
@@ -181,7 +181,7 @@ public class QueryMetadataExecutor extends 
ClusterQPExecutor {
                   response.addTimeSeries(mManager.getShowTimeseriesPath(path));
                 }
               } catch (final PathErrorException e) {
-                response = 
QueryTimeSeriesResponse.createErrorInstance(groupId, e.toString());
+                response = 
QueryTimeSeriesResponse.createErrorInstance(groupId, e.getMessage());
               }
             } else {
               response = QueryTimeSeriesResponse.createErrorInstance(groupId, 
status.getErrorMsg());
@@ -227,7 +227,7 @@ public class QueryMetadataExecutor extends 
ClusterQPExecutor {
                 response = QueryStorageGroupResponse
                     
.createSuccessInstance(metadataHolder.getFsm().getAllStorageGroups());
               } catch (final PathErrorException e) {
-                response = 
QueryStorageGroupResponse.createErrorInstance(e.toString());
+                response = 
QueryStorageGroupResponse.createErrorInstance(e.getMessage());
               }
             } else {
               response = 
QueryStorageGroupResponse.createErrorInstance(status.getErrorMsg());
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/rpc/impl/RaftNodeAsClientManager.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/rpc/impl/RaftNodeAsClientManager.java
index 9330d63..2719340 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/rpc/impl/RaftNodeAsClientManager.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/rpc/impl/RaftNodeAsClientManager.java
@@ -201,7 +201,7 @@ public class RaftNodeAsClientManager {
                   }
                 }, TASK_TIMEOUT_MS);
       } catch (RemotingException | InterruptedException e) {
-        LOGGER.error(e.toString());
+        LOGGER.error(e.getMessage());
         throw new RaftConnectionException(e);
       }
       releaseClient();
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/DataGroupNonQueryAsyncProcessor.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/DataGroupNonQueryAsyncProcessor.java
index 709caf2..d2ef52e 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/DataGroupNonQueryAsyncProcessor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/DataGroupNonQueryAsyncProcessor.java
@@ -84,7 +84,7 @@ public class DataGroupNonQueryAsyncProcessor extends
             .wrap(SerializerManager.getSerializer(SerializerManager.Hessian2)
                 .serialize(dataGroupNonQueryRequest)));
       } catch (final CodecException e) {
-        response.setErrorMsg(e.toString());
+        response.setErrorMsg(e.getMessage());
         response.addResult(false);
         asyncContext.sendResponse(response);
       }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/MetaGroupNonQueryAsyncProcessor.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/MetaGroupNonQueryAsyncProcessor.java
index 8c2c12a..2049651 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/MetaGroupNonQueryAsyncProcessor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/MetaGroupNonQueryAsyncProcessor.java
@@ -84,7 +84,7 @@ public class MetaGroupNonQueryAsyncProcessor extends
                 .serialize(metaGroupNonQueryRequest)));
       } catch (final CodecException e) {
         response.addResult(false);
-        response.setErrorMsg(e.toString());
+        response.setErrorMsg(e.getMessage());
         asyncContext.sendResponse(response);
       }
 
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/QueryTimeSeriesAsyncProcessor.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/QueryTimeSeriesAsyncProcessor.java
index 9c254ac..e4ef357 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/QueryTimeSeriesAsyncProcessor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/rpc/processor/QueryTimeSeriesAsyncProcessor.java
@@ -68,7 +68,7 @@ public class QueryTimeSeriesAsyncProcessor extends 
BasicAsyncUserProcessor<Query
                   response.addTimeSeries(mManager.getShowTimeseriesPath(path));
                 }
               } catch (final PathErrorException e) {
-                response = 
QueryTimeSeriesResponse.createErrorInstance(groupId, e.toString());
+                response = 
QueryTimeSeriesResponse.createErrorInstance(groupId, e.getMessage());
               }
             } else {
               response = QueryTimeSeriesResponse.createErrorInstance(groupId, 
status.getErrorMsg());
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/rpc/service/TSServiceClusterImpl.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/rpc/service/TSServiceClusterImpl.java
index 6b42362..0a2c456 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/rpc/service/TSServiceClusterImpl.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/rpc/service/TSServiceClusterImpl.java
@@ -18,12 +18,15 @@
  */
 package org.apache.iotdb.cluster.rpc.service;
 
+import com.alipay.sofa.jraft.util.OnlyForTest;
 import java.io.IOException;
 import java.sql.Statement;
 import java.util.Arrays;
 import java.util.List;
 import java.util.Set;
+import java.util.regex.Pattern;
 import java.util.stream.Collectors;
+import org.apache.iotdb.cluster.config.ClusterConstant;
 import org.apache.iotdb.cluster.qp.executor.NonQueryExecutor;
 import org.apache.iotdb.cluster.qp.executor.QueryMetadataExecutor;
 import org.apache.iotdb.db.auth.AuthException;
@@ -184,6 +187,28 @@ public class TSServiceClusterImpl extends TSServiceImpl {
   }
 
   @Override
+  public boolean execSetConsistencyLevel(String statement) throws Exception {
+    if (statement == null) {
+      return false;
+    }
+    statement = statement.toLowerCase().trim();
+    if 
(Pattern.matches(ClusterConstant.SET_READ_METADATA_CONSISTENCY_LEVEL_PATTERN, 
statement)) {
+      String[] splits = statement.split("\\s+");
+      int level = Integer.valueOf(splits[splits.length-1]);
+        nonQueryExecutor.get().setReadMetadataConsistencyLevel(level);
+      return true;
+    } else if (Pattern
+        .matches(ClusterConstant.SET_READ_DATA_CONSISTENCY_LEVEL_PATTERN, 
statement)) {
+      String[] splits = statement.split("\\s+");
+      int level = Integer.valueOf(splits[splits.length-1]);
+      nonQueryExecutor.get().setReadDataConsistencyLevel(level);
+      return true;
+    } else{
+      return false;
+    }
+  }
+
+  @Override
   protected boolean executeNonQuery(PhysicalPlan plan) throws 
ProcessorException {
     return nonQueryExecutor.get().processNonQuery(plan);
   }
@@ -213,4 +238,9 @@ public class TSServiceClusterImpl extends TSServiceImpl {
       throws InterruptedException, ProcessorException {
     return queryMetadataExecutor.get().processMetadataInStringQuery();
   }
+
+  @OnlyForTest
+  public NonQueryExecutor getNonQueryExecutor() {
+    return nonQueryExecutor.get();
+  }
 }
diff --git 
a/cluster/src/test/java/org/apache/iotdb/cluster/qp/ClusterQPExecutorTest.java 
b/cluster/src/test/java/org/apache/iotdb/cluster/qp/ClusterQPExecutorTest.java
new file mode 100644
index 0000000..efe6431
--- /dev/null
+++ 
b/cluster/src/test/java/org/apache/iotdb/cluster/qp/ClusterQPExecutorTest.java
@@ -0,0 +1,113 @@
+/**
+ * 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.qp;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+import java.io.IOException;
+import java.util.concurrent.CountDownLatch;
+import org.apache.iotdb.cluster.config.ClusterConfig;
+import org.apache.iotdb.cluster.config.ClusterDescriptor;
+import org.apache.iotdb.cluster.qp.executor.NonQueryExecutor;
+import org.apache.iotdb.cluster.rpc.service.TSServiceClusterImpl;
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+
+public class ClusterQPExecutorTest {
+
+  private static final ClusterConfig CLUSTER_CONFIG = 
ClusterDescriptor.getInstance().getConfig();
+
+  private TSServiceClusterImpl impl;
+
+  private NonQueryExecutor executor;
+
+  @Before
+  public void setUp() throws Exception {
+    impl = new TSServiceClusterImpl();
+    impl.initClusterService();
+
+    executor = impl.getNonQueryExecutor();
+  }
+
+  @After
+  public void tearDown() {
+    impl.closeClusterService();
+  }
+
+  @Test
+  public void setReadMetadataConsistencyLevel() throws Exception {
+    assertEquals(CLUSTER_CONFIG.getReadMetadataConsistencyLevel(),
+        executor.getReadMetadataConsistencyLevel());
+    boolean exec;
+    exec= impl.execSetConsistencyLevel("set read metadata level 1");
+    assertTrue(exec);
+    assertEquals(1, executor.getReadMetadataConsistencyLevel());
+
+    exec= impl.execSetConsistencyLevel("show timeseries");
+    assertEquals(1, executor.getReadMetadataConsistencyLevel());
+    assertFalse(exec);
+
+    exec= impl.execSetConsistencyLevel("set read metadata level 2");
+    assertTrue(exec);
+    assertEquals(2, executor.getReadMetadataConsistencyLevel());
+
+    exec = impl.execSetConsistencyLevel("set read metadata level -2");
+    assertEquals(2, executor.getReadMetadataConsistencyLevel());
+    assertFalse(exec);
+
+    try {
+      impl.execSetConsistencyLevel("set read metadata level 90");
+    } catch (Exception e) {
+      assertEquals("Consistency level 90 not support", e.getMessage());
+    }
+    assertEquals(2, executor.getReadMetadataConsistencyLevel());
+  }
+
+  @Test
+  public void setReadDataConsistencyLevel() throws Exception {
+    assertEquals(CLUSTER_CONFIG.getReadDataConsistencyLevel(),
+        executor.getReadDataConsistencyLevel());
+    boolean exec;
+    exec= impl.execSetConsistencyLevel("set read data level 1");
+    assertTrue(exec);
+    assertEquals(1, executor.getReadDataConsistencyLevel());
+
+    exec= impl.execSetConsistencyLevel("show timeseries");
+    assertEquals(1, executor.getReadDataConsistencyLevel());
+    assertFalse(exec);
+
+    exec= impl.execSetConsistencyLevel("set read metadata level 2");
+    assertTrue(exec);
+    assertEquals(2, executor.getReadMetadataConsistencyLevel());
+
+    exec = impl.execSetConsistencyLevel("set read metadata level -2");
+    assertEquals(2, executor.getReadMetadataConsistencyLevel());
+    assertFalse(exec);
+
+    try {
+      impl.execSetConsistencyLevel("set read metadata level 90");
+    } catch (Exception e) {
+      assertEquals("Consistency level 90 not support", e.getMessage());
+    }
+    assertEquals(2, executor.getReadMetadataConsistencyLevel());
+  }
+}
\ No newline at end of file
diff --git a/iotdb/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java 
b/iotdb/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java
index a144e73..e8206a1 100644
--- a/iotdb/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java
+++ b/iotdb/src/main/java/org/apache/iotdb/db/service/TSServiceImpl.java
@@ -472,6 +472,16 @@ public class TSServiceImpl implements TSIService.Iface, 
ServerContext {
         return getTSExecuteStatementResp(TS_StatusCode.ERROR_STATUS, 
e.getMessage());
       }
 
+      try{
+        if (execSetConsistencyLevel(statement)) {
+          return getTSExecuteStatementResp(TS_StatusCode.SUCCESS_STATUS,
+              "Execute set consistency level successfully");
+        }
+      }catch (Exception e){
+        LOGGER.error("meet error while executing set consistency level 
command!", e);
+        return getTSExecuteStatementResp(TS_StatusCode.ERROR_STATUS, 
e.getMessage());
+      }
+
       PhysicalPlan physicalPlan;
       try {
         physicalPlan = processor.parseSQLToPhysicalPlan(statement, 
zoneIds.get());
@@ -495,6 +505,13 @@ public class TSServiceImpl implements TSIService.Iface, 
ServerContext {
     }
   }
 
+  /**
+   * Set consistency level
+   */
+  public boolean execSetConsistencyLevel(String statement) throws Exception {
+    return false;
+  }
+
   @Override
   public TSExecuteStatementResp executeQueryStatement(TSExecuteStatementReq 
req) throws TException {
 

Reply via email to