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

commit 4f92cc5d0a12055ea61fa08ca487e509374e7078
Author: xiangdong huang <[email protected]>
AuthorDate: Tue Jul 27 19:00:53 2021 +0800

    a temporary submit..
---
 .../resources/conf/iotdb-cluster.properties        |   5 -
 .../{ClusterMain.java => ClusterIoTDB.java}        | 161 ++++++++++++++-------
 .../iotdb/cluster/config/ClusterDescriptor.java    |  12 +-
 .../iotdb/cluster/server/ClusterRPCService.java    |  65 +++++----
 .../cluster/server/ClusterRPCServiceMBean.java     |  34 +++++
 ...ClientServer.java => ClusterTSServiceImpl.java} | 133 +++--------------
 .../iotdb/cluster/server/MetaClusterServer.java    |  27 ++--
 .../cluster/server/member/MetaGroupMember.java     |  30 ++--
 .../cluster/utils/nodetool/ClusterMonitor.java     |   4 +-
 .../cluster/integration/BaseSingleNodeTest.java    |   1 +
 .../clusterinfo/ClusterInfoServiceImplTest.java    |   8 +-
 .../org/apache/iotdb/db/concurrent/ThreadName.java |   4 +-
 .../apache/iotdb/db/qp/executor/PlanExecutor.java  |   4 +-
 .../iotdb/db/query/control/TracingManager.java     |  33 ++---
 .../java/org/apache/iotdb/db/service/IoTDB.java    |  20 +--
 .../org/apache/iotdb/db/service/RPCService.java    |  13 +-
 .../iotdb/db/service/RPCServiceThriftHandler.java  |   2 +-
 .../org/apache/iotdb/db/service/ServiceType.java   |   3 +
 .../iotdb/db/service/thrift/ThriftService.java     |   5 -
 .../db/service/thrift/ThriftServiceThread.java     |   4 +-
 .../iotdb/db/query/control/TracingManagerTest.java |   3 -
 21 files changed, 273 insertions(+), 298 deletions(-)

diff --git a/cluster/src/assembly/resources/conf/iotdb-cluster.properties 
b/cluster/src/assembly/resources/conf/iotdb-cluster.properties
index 2bdac88..4228249 100644
--- a/cluster/src/assembly/resources/conf/iotdb-cluster.properties
+++ b/cluster/src/assembly/resources/conf/iotdb-cluster.properties
@@ -64,11 +64,6 @@ seed_nodes=127.0.0.1:9003
 # WARNING: this must be consistent across all nodes in the cluster
 # rpc_thrift_compression_enable=false
 
-# max client connections created by thrift
-# this configuration applies separately to data/meta/client connections and 
thus does not control
-# the number of global connections
-# max_concurrent_client_num=10000
-
 # number of replications for one partition
 default_replica_num=1
 
diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterMain.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
similarity index 74%
rename from cluster/src/main/java/org/apache/iotdb/cluster/ClusterMain.java
rename to cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
index adc4661..cec70b2 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterMain.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
@@ -27,15 +27,20 @@ import 
org.apache.iotdb.cluster.exception.StartUpCheckFailureException;
 import org.apache.iotdb.cluster.partition.slot.SlotPartitionTable;
 import org.apache.iotdb.cluster.partition.slot.SlotStrategy;
 import org.apache.iotdb.cluster.rpc.thrift.Node;
+import org.apache.iotdb.cluster.server.ClusterRPCService;
 import org.apache.iotdb.cluster.server.MetaClusterServer;
 import org.apache.iotdb.cluster.server.Response;
 import org.apache.iotdb.cluster.server.clusterinfo.ClusterInfoServer;
 import org.apache.iotdb.cluster.utils.ClusterUtils;
+import org.apache.iotdb.cluster.utils.nodetool.ClusterMonitor;
 import org.apache.iotdb.db.conf.IoTDBConfigCheck;
 import org.apache.iotdb.db.conf.IoTDBConstant;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.exception.StartupException;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
+import org.apache.iotdb.db.service.IoTDB;
+import org.apache.iotdb.db.service.JMXService;
+import org.apache.iotdb.db.service.RegisterManager;
 import org.apache.iotdb.db.utils.TestOnly;
 
 import org.apache.thrift.TException;
@@ -53,9 +58,12 @@ import java.util.Set;
 
 import static org.apache.iotdb.cluster.utils.ClusterUtils.UNKNOWN_CLIENT_IP;
 
-public class ClusterMain {
+//we do not inherent IoTDB instance, as it may break the singleton mode of 
IoTDB.
+public class ClusterIoTDB {
 
-  private static final Logger logger = 
LoggerFactory.getLogger(ClusterMain.class);
+  private static final Logger logger = 
LoggerFactory.getLogger(ClusterIoTDB.class);
+  private final String mbeanName =
+      String.format("%s:%s=%s", IoTDBConstant.IOTDB_PACKAGE, 
IoTDBConstant.JMX_TYPE, "ClusterIoTDB");
 
   // establish the cluster as a seed
   private static final String MODE_START = "-s";
@@ -65,7 +73,12 @@ public class ClusterMain {
   // metaport-of-removed-node
   private static final String MODE_REMOVE = "-r";
 
-  private static MetaClusterServer metaServer;
+  private MetaClusterServer metaServer;
+
+  private IoTDB iotdb = IoTDB.getInstance();
+
+  // Cluster IoTDB uses a individual registerManager with its parent.
+  private RegisterManager registerManager = new RegisterManager();
 
   public static void main(String[] args) {
     if (args.length < 1) {
@@ -76,7 +89,6 @@ public class ClusterMain {
               + "-a: start the node as a new node\n"
               + "-r: remove the node out of the cluster\n",
           IoTDBConstant.IOTDB_CONF);
-
       return;
     }
 
@@ -103,50 +115,18 @@ public class ClusterMain {
     String mode = args[0];
     logger.info("Running mode {}", mode);
 
+    ClusterIoTDB cluster = ClusterIoTDBHolder.INSTANCE;
+    //we start IoTDB kernel first.
+    cluster.iotdb.active();
+
+    //then we start the cluster module.
     if (MODE_START.equals(mode)) {
-      try {
-        metaServer = new MetaClusterServer();
-        startServerCheck();
-        preStartCustomize();
-        metaServer.start();
-        metaServer.buildCluster();
-        // Currently, we do not register ClusterInfoService as a JMX Bean,
-        // so we use startService() rather than start()
-        ClusterInfoServer.getInstance().startService();
-      } catch (TTransportException
-          | StartupException
-          | QueryProcessException
-          | StartUpCheckFailureException
-          | ConfigInconsistentException e) {
-        metaServer.stop();
-        logger.error("Fail to start meta server", e);
-      }
+      cluster.activeStartNodeMode();
     } else if (MODE_ADD.equals(mode)) {
-      try {
-        long startTime = System.currentTimeMillis();
-        metaServer = new MetaClusterServer();
-        preStartCustomize();
-        metaServer.start();
-        metaServer.joinCluster();
-        // Currently, we do not register ClusterInfoService as a JMX Bean,
-        // so we use startService() rather than start()
-        ClusterInfoServer.getInstance().startService();
-
-        logger.info(
-            "Adding this node {} to cluster costs {} ms",
-            metaServer.getMember().getThisNode(),
-            (System.currentTimeMillis() - startTime));
-      } catch (TTransportException
-          | StartupException
-          | QueryProcessException
-          | StartUpCheckFailureException
-          | ConfigInconsistentException e) {
-        metaServer.stop();
-        logger.error("Fail to join cluster", e);
-      }
+      cluster.activeAddNodeMode();
     } else if (MODE_REMOVE.equals(mode)) {
       try {
-        doRemoveNode(args);
+        cluster.doRemoveNode(args);
       } catch (IOException e) {
         logger.error("Fail to remove node in cluster", e);
       }
@@ -155,7 +135,61 @@ public class ClusterMain {
     }
   }
 
-  private static void startServerCheck() throws StartupException {
+  public void activeStartNodeMode() {
+    try {
+      metaServer = new MetaClusterServer();
+      startServerCheck();
+      preStartCustomize();
+      metaServer.start();
+      metaServer.buildCluster();
+      // Currently, we do not register ClusterInfoService as a JMX Bean,
+      // so we use startService() rather than start()
+      ClusterInfoServer.getInstance().startService();
+      // JMX based DBA API
+      registerManager.register(ClusterMonitor.INSTANCE);
+      // we must wait until the metaGroup established.
+      // So that the ClusterRPCService can work.
+      registerManager.register(ClusterRPCService.getInstance());
+    } catch (TTransportException
+        | StartupException
+        | QueryProcessException
+        | StartUpCheckFailureException
+        | ConfigInconsistentException e) {
+      stop();
+      logger.error("Fail to start meta server", e);
+    }
+  }
+
+  public void activeAddNodeMode() {
+    try {
+      long startTime = System.currentTimeMillis();
+      metaServer = new MetaClusterServer();
+      preStartCustomize();
+      metaServer.start();
+      metaServer.joinCluster();
+      // Currently, we do not register ClusterInfoService as a JMX Bean,
+      // so we use startService() rather than start()
+      ClusterInfoServer.getInstance().startService();
+      // JMX based DBA API
+      registerManager.register(ClusterMonitor.INSTANCE);
+      //finally, we start the RPC service
+      registerManager.register(ClusterRPCService.getInstance());
+      logger.info(
+          "Adding this node {} to cluster costs {} ms",
+          metaServer.getMember().getThisNode(),
+          (System.currentTimeMillis() - startTime));
+    } catch (TTransportException
+        | StartupException
+        | QueryProcessException
+        | StartUpCheckFailureException
+        | ConfigInconsistentException e) {
+      stop();
+      logger.error("Fail to join cluster", e);
+    }
+  }
+
+
+  private void startServerCheck() throws StartupException {
     ClusterConfig config = ClusterDescriptor.getInstance().getConfig();
     // check the initial replicateNum and refuse to start when the 
replicateNum <= 0
     if (config.getReplicationNum() <= 0) {
@@ -218,7 +252,7 @@ public class ClusterMain {
     }
   }
 
-  private static void doRemoveNode(String[] args) throws IOException {
+  private void doRemoveNode(String[] args) throws IOException {
     if (args.length != 3) {
       logger.error("Usage: -r <ip> <metaPort>");
       return;
@@ -255,7 +289,7 @@ public class ClusterMain {
     }
   }
 
-  private static void handleNodeRemovalResp(Long response, Node nodeToRemove, 
long startTime) {
+  private void handleNodeRemovalResp(Long response, Node nodeToRemove, long 
startTime) {
     if (response == Response.RESPONSE_AGREE) {
       logger.info(
           "Node {} is successfully removed, cost {}ms",
@@ -273,13 +307,13 @@ public class ClusterMain {
     }
   }
 
-  public static MetaClusterServer getMetaServer() {
+  public MetaClusterServer getMetaServer() {
     return metaServer;
   }
 
   /** Developers may perform pre-start customizations here for debugging or 
experiments. */
   @SuppressWarnings("java:S125") // leaving examples
-  private static void preStartCustomize() {
+  private void preStartCustomize() {
     // customize data distribution
     // The given example tries to divide storage groups like "root.sg_1", 
"root.sg_2"... into k
     // nodes evenly, and use default strategy for other groups
@@ -324,8 +358,35 @@ public class ClusterMain {
         });
   }
 
+
+
   @TestOnly
-  public static void setMetaClusterServer(MetaClusterServer metaClusterServer) 
{
+  public void setMetaClusterServer(MetaClusterServer metaClusterServer) {
     metaServer = metaClusterServer;
   }
+
+  public void stop() {
+    deactivate();
+  }
+
+  private void deactivate() {
+    logger.info("Deactivating Cluster IoTDB...");
+    metaServer.stop();
+    registerManager.deregisterAll();
+    JMXService.deregisterMBean(mbeanName);
+    logger.info("ClusterIoTDB is deactivated.");
+    //stop the iotdb kernel
+    iotdb.stop();
+  }
+
+
+  public static ClusterIoTDB getInstance() {
+    return ClusterIoTDBHolder.INSTANCE;
+  }
+  private static class ClusterIoTDBHolder {
+
+    private static final ClusterIoTDB INSTANCE = new ClusterIoTDB();
+
+    private ClusterIoTDBHolder() {}
+  }
 }
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 9d983ac..38c1a35 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
@@ -151,11 +151,6 @@ public class ClusterDescriptor {
             properties.getProperty(
                 "cluster_info_public_port", 
Integer.toString(config.getClusterInfoRpcPort()))));
 
-    config.setMaxConcurrentClientNum(
-        Integer.parseInt(
-            properties.getProperty(
-                "max_concurrent_client_num", 
String.valueOf(config.getMaxConcurrentClientNum()))));
-
     config.setMultiRaftFactor(
         Integer.parseInt(
             properties.getProperty(
@@ -373,18 +368,13 @@ public class ClusterDescriptor {
 
   /**
    * This method is for setting hot modified properties of the cluster. 
Currently, we support
-   * max_concurrent_client_num, connection_timeout_ms, max_resolved_log_size
+   * connection_timeout_ms, max_resolved_log_size
    *
    * @param properties
    * @throws QueryProcessException
    */
   public void loadHotModifiedProps(Properties properties) {
 
-    config.setMaxConcurrentClientNum(
-        Integer.parseInt(
-            properties.getProperty(
-                "max_concurrent_client_num", 
String.valueOf(config.getMaxConcurrentClientNum()))));
-
     config.setConnectionTimeoutInMS(
         Integer.parseInt(
             properties.getProperty(
diff --git a/server/src/main/java/org/apache/iotdb/db/service/RPCService.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/ClusterRPCService.java
similarity index 61%
copy from server/src/main/java/org/apache/iotdb/db/service/RPCService.java
copy to 
cluster/src/main/java/org/apache/iotdb/cluster/server/ClusterRPCService.java
index 5bbddec..8c79ac5 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/RPCService.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/ClusterRPCService.java
@@ -16,45 +16,43 @@
  * specific language governing permissions and limitations
  * under the License.
  */
-package org.apache.iotdb.db.service;
 
+package org.apache.iotdb.cluster.server;
+
+import org.apache.iotdb.cluster.config.ClusterDescriptor;
+import org.apache.iotdb.cluster.coordinator.Coordinator;
+import org.apache.iotdb.cluster.server.member.MetaGroupMember;
 import org.apache.iotdb.db.concurrent.ThreadName;
 import org.apache.iotdb.db.conf.IoTDBConfig;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.exception.query.QueryProcessException;
 import org.apache.iotdb.db.exception.runtime.RPCServiceException;
+import org.apache.iotdb.db.service.RPCServiceThriftHandler;
+import org.apache.iotdb.db.service.ServiceType;
 import org.apache.iotdb.db.service.thrift.ThriftService;
 import org.apache.iotdb.db.service.thrift.ThriftServiceThread;
 import org.apache.iotdb.service.rpc.thrift.TSIService.Processor;
 
-/** A service to handle jdbc request from client. */
-public class RPCService extends ThriftService implements RPCServiceMBean {
-
-  private TSServiceImpl impl;
+public class ClusterRPCService  extends ThriftService implements 
ClusterRPCServiceMBean {
 
-  private RPCService() {}
+  private ClusterTSServiceImpl impl;
 
-  public static RPCService getInstance() {
-    return RPCServiceHolder.INSTANCE;
-  }
+  private ClusterRPCService() {}
 
   @Override
-  public int getRPCPort() {
-    IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
-    return config.getRpcPort();
+  public ThriftService getImplementation() {
+    return ClusterRPCServiceHolder.INSTANCE;
   }
 
   @Override
-  public ThriftService getImplementation() {
-    return getInstance();
+  public ServiceType getID() {
+    return ServiceType.CLUSTER_RPC_SERVICE;
   }
 
   @Override
   public void initTProcessor()
-      throws ClassNotFoundException, IllegalAccessException, 
InstantiationException {
-    impl =
-        (TSServiceImpl)
-            
Class.forName(IoTDBDescriptor.getInstance().getConfig().getRpcImplClassName())
-                .newInstance();
+      throws IllegalAccessException, InstantiationException {
+    impl = new ClusterTSServiceImpl();
     processor = new Processor<>(impl);
   }
 
@@ -66,7 +64,7 @@ public class RPCService extends ThriftService implements 
RPCServiceMBean {
           new ThriftServiceThread(
               processor,
               getID().getName(),
-              ThreadName.RPC_CLIENT.getName(),
+              ThreadName.CLUSTER_RPC_CLIENT.getName(),
               config.getRpcAddress(),
               config.getRpcPort(),
               config.getRpcMaxConcurrentClientNum(),
@@ -76,7 +74,7 @@ public class RPCService extends ThriftService implements 
RPCServiceMBean {
     } catch (RPCServiceException e) {
       throw new IllegalAccessException(e.getMessage());
     }
-    thriftServiceThread.setName(ThreadName.RPC_SERVICE.getName());
+    thriftServiceThread.setName(ThreadName.CLUSTER_RPC_SERVICE.getName());
   }
 
   @Override
@@ -86,18 +84,31 @@ public class RPCService extends ThriftService implements 
RPCServiceMBean {
 
   @Override
   public int getBindPort() {
-    return IoTDBDescriptor.getInstance().getConfig().getRpcPort();
+    return ClusterDescriptor.getInstance().getConfig().getClusterRpcPort();
   }
 
   @Override
-  public ServiceType getID() {
-    return ServiceType.RPC_SERVICE;
+  public int getRPCPort() {
+    return getBindPort();
+  }
+
+  public static ClusterRPCService getInstance() {
+    return ClusterRPCServiceHolder.INSTANCE;
+  }
+
+  public void assignExecutorToServiceImpl(MetaGroupMember member) throws 
QueryProcessException {
+    this.impl.setExecutor(member);
+  }
+
+  public void assignCoordinator(Coordinator coordinator) {
+    this.impl.setCoordinator(coordinator);
   }
 
-  private static class RPCServiceHolder {
 
-    private static final RPCService INSTANCE = new RPCService();
+  private static class ClusterRPCServiceHolder {
+
+    private static final ClusterRPCService INSTANCE = new ClusterRPCService();
 
-    private RPCServiceHolder() {}
+    private ClusterRPCServiceHolder() {}
   }
 }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/ClusterRPCServiceMBean.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/ClusterRPCServiceMBean.java
new file mode 100644
index 0000000..d7dde8a
--- /dev/null
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/ClusterRPCServiceMBean.java
@@ -0,0 +1,34 @@
+/*
+ * 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;
+
+import org.apache.iotdb.db.exception.StartupException;
+
+public interface ClusterRPCServiceMBean {
+
+  String getRPCServiceStatus();
+
+  int getRPCPort();
+
+  void startService() throws StartupException;
+
+  void restartService() throws StartupException;
+
+  void stopService();
+}
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/ClusterTSServiceImpl.java
similarity index 63%
rename from 
cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java
rename to 
cluster/src/main/java/org/apache/iotdb/cluster/server/ClusterTSServiceImpl.java
index 1200c2d..e43e503 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/server/ClientServer.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/ClusterTSServiceImpl.java
@@ -21,7 +21,6 @@ package org.apache.iotdb.cluster.server;
 
 import org.apache.iotdb.cluster.client.async.AsyncDataClient;
 import org.apache.iotdb.cluster.client.sync.SyncDataClient;
-import org.apache.iotdb.cluster.config.ClusterConfig;
 import org.apache.iotdb.cluster.config.ClusterDescriptor;
 import org.apache.iotdb.cluster.coordinator.Coordinator;
 import org.apache.iotdb.cluster.metadata.CMManager;
@@ -32,7 +31,6 @@ import org.apache.iotdb.cluster.rpc.thrift.Node;
 import org.apache.iotdb.cluster.rpc.thrift.RaftNode;
 import org.apache.iotdb.cluster.server.handlers.caller.GenericHandler;
 import org.apache.iotdb.cluster.server.member.MetaGroupMember;
-import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.exception.StorageEngineException;
 import org.apache.iotdb.db.exception.metadata.MetadataException;
 import org.apache.iotdb.db.exception.query.QueryProcessException;
@@ -41,74 +39,50 @@ import org.apache.iotdb.db.qp.physical.PhysicalPlan;
 import org.apache.iotdb.db.query.context.QueryContext;
 import org.apache.iotdb.db.service.IoTDB;
 import org.apache.iotdb.db.service.TSServiceImpl;
-import org.apache.iotdb.db.utils.CommonUtils;
-import org.apache.iotdb.rpc.RpcTransportFactory;
 import org.apache.iotdb.rpc.RpcUtils;
 import org.apache.iotdb.rpc.TSStatusCode;
-import org.apache.iotdb.service.rpc.thrift.TSIService.Processor;
 import org.apache.iotdb.service.rpc.thrift.TSStatus;
 import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
 
 import org.apache.thrift.TException;
-import org.apache.thrift.protocol.TBinaryProtocol;
-import org.apache.thrift.protocol.TCompactProtocol;
 import org.apache.thrift.protocol.TProtocol;
-import org.apache.thrift.protocol.TProtocolFactory;
 import org.apache.thrift.server.ServerContext;
-import org.apache.thrift.server.TServer;
 import org.apache.thrift.server.TServerEventHandler;
-import org.apache.thrift.server.TThreadPoolServer;
-import org.apache.thrift.transport.TServerSocket;
-import org.apache.thrift.transport.TServerTransport;
 import org.apache.thrift.transport.TTransport;
-import org.apache.thrift.transport.TTransportException;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.io.IOException;
-import java.net.InetSocketAddress;
 import java.util.List;
 import java.util.Map;
 import java.util.Map.Entry;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.ExecutorService;
-import java.util.concurrent.Executors;
-import java.util.concurrent.SynchronousQueue;
-import java.util.concurrent.ThreadFactory;
-import java.util.concurrent.ThreadPoolExecutor;
-import java.util.concurrent.atomic.AtomicLong;
 import java.util.concurrent.atomic.AtomicReference;
 
 /**
- * ClientServer is the cluster version of TSServiceImpl, which is responsible 
for the processing of
+ * ClusterTSServiceImpl is the cluster version of TSServiceImpl, which is 
responsible for the processing of
  * the user requests (sqls and session api). It inherits the basic procedures 
from TSServiceImpl,
  * but redirect the queries of data and metadata to a MetaGroupMember of the 
local node.
  */
-public class ClientServer extends TSServiceImpl {
+public class ClusterTSServiceImpl extends TSServiceImpl {
 
-  private static final Logger logger = 
LoggerFactory.getLogger(ClientServer.class);
+  private static final Logger logger = 
LoggerFactory.getLogger(ClusterTSServiceImpl.class);
   /**
-   * The Coordinator of the local node. Through this node ClientServer queries 
data and meta from
+   * The Coordinator of the local node. Through this node queries data and 
meta from
    * the cluster and performs data manipulations to the cluster.
    */
   private Coordinator coordinator;
 
-  public void setCoordinator(Coordinator coordinator) {
-    this.coordinator = coordinator;
-  }
-
-  /** The single thread pool that runs poolServer to unblock the main thread. 
*/
-  private ExecutorService serverService;
-
-  /**
-   * Using the poolServer, ClientServer will listen to a socket to accept 
thrift requests like an
-   * HttpServer.
-   */
-  private TServer poolServer;
 
-  /** The socket poolServer will listen to. Async service requires nonblocking 
socket */
-  private TServerTransport serverTransport;
+//  /**
+//   * Using the poolServer, ClusterTSServiceImpl will listen to a socket to 
accept thrift requests like an
+//   * HttpServer.
+//   */
+//  private TServer poolServer;
+//
+//  /** The socket poolServer will listen to. Async service requires 
nonblocking socket */
+//  private TServerTransport serverTransport;
 
   /**
    * queryId -> queryContext map. When a query ends either normally or 
accidentally, the resources
@@ -116,89 +90,20 @@ public class ClientServer extends TSServiceImpl {
    */
   private Map<Long, RemoteQueryContext> queryContextMap = new 
ConcurrentHashMap<>();
 
-  public ClientServer(MetaGroupMember metaGroupMember) throws 
QueryProcessException {
+  public ClusterTSServiceImpl(MetaGroupMember metaGroupMember) throws 
QueryProcessException {
     super();
     this.processor = new ClusterPlanner();
-    this.executor = new ClusterPlanExecutor(metaGroupMember);
   }
 
-  /**
-   * Create a thrift server to listen to the client port and accept requests 
from clients. This
-   * server is run in a separate thread. Calling the method twice does not 
induce side effects.
-   *
-   * @throws TTransportException
-   */
-  public void start() throws TTransportException {
-    if (serverService != null) {
-      return;
-    }
-
-    serverService = Executors.newSingleThreadExecutor(r -> new Thread(r, 
"ClusterClientServer"));
-    ClusterConfig config = ClusterDescriptor.getInstance().getConfig();
-
-    // this defines how thrift parse the requests bytes to a request
-    TProtocolFactory protocolFactory;
-    if 
(IoTDBDescriptor.getInstance().getConfig().isRpcThriftCompressionEnable()) {
-      protocolFactory = new TCompactProtocol.Factory();
-    } else {
-      protocolFactory = new TBinaryProtocol.Factory();
-    }
-    serverTransport =
-        new TServerSocket(
-            new InetSocketAddress(
-                IoTDBDescriptor.getInstance().getConfig().getRpcAddress(),
-                config.getClusterRpcPort()));
-    // async service also requires nonblocking server, and HsHaServer is 
basically more efficient a
-    // nonblocking server
-    int maxConcurrentClientNum =
-        Math.max(CommonUtils.getCpuCores(), 
config.getMaxConcurrentClientNum());
-    TThreadPoolServer.Args poolArgs =
-        new TThreadPoolServer.Args(serverTransport)
-            .maxWorkerThreads(maxConcurrentClientNum)
-            .minWorkerThreads(CommonUtils.getCpuCores());
-    poolArgs.executorService(
-        new ThreadPoolExecutor(
-            poolArgs.minWorkerThreads,
-            poolArgs.maxWorkerThreads,
-            poolArgs.stopTimeoutVal,
-            poolArgs.stopTimeoutUnit,
-            new SynchronousQueue<>(),
-            new ThreadFactory() {
-              private AtomicLong threadIndex = new AtomicLong(0);
-
-              @Override
-              public Thread newThread(Runnable r) {
-                return new Thread(r, "ClusterClient-" + 
threadIndex.incrementAndGet());
-              }
-            }));
-    // ClientServer will do the following processing when the HsHaServer has 
parsed a request
-    poolArgs.processor(new Processor<>(this));
-    poolArgs.protocolFactory(protocolFactory);
-    // nonblocking server requests FramedTransport
-    poolArgs.transportFactory(RpcTransportFactory.INSTANCE);
-
-    poolServer = new TThreadPoolServer(poolArgs);
-    // mainly for handling client exit events
-    poolServer.setServerEventHandler(new EventHandler());
-
-    serverService.submit(() -> poolServer.serve());
-    logger.info("Client service is set up");
+  public void setExecutor(MetaGroupMember metaGroupMember) throws 
QueryProcessException {
+    this.executor = new ClusterPlanExecutor(metaGroupMember);
   }
 
-  /**
-   * Stop the thrift server, close the socket and shutdown the thread pool. 
Calling the method twice
-   * does not induce side effects.
-   */
-  public void stop() {
-    if (serverService == null) {
-      return;
-    }
-
-    poolServer.stop();
-    serverService.shutdownNow();
-    serverTransport.close();
+  public void setCoordinator(Coordinator coordinator) {
+    this.coordinator = coordinator;
   }
 
+
   /**
    * Redirect the plan to the local Coordinator so that it will be processed 
cluster-wide.
    *
@@ -235,7 +140,7 @@ public class ClientServer extends TSServiceImpl {
 
     @Override
     public void deleteContext(ServerContext serverContext, TProtocol input, 
TProtocol output) {
-      ClientServer.this.handleClientExit();
+      ClusterTSServiceImpl.this.handleClientExit();
     }
 
     @Override
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/MetaClusterServer.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/MetaClusterServer.java
index fce7a87..b0f8a25 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/server/MetaClusterServer.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/server/MetaClusterServer.java
@@ -18,6 +18,7 @@
  */
 package org.apache.iotdb.cluster.server;
 
+import org.apache.iotdb.cluster.ClusterIoTDB;
 import org.apache.iotdb.cluster.config.ClusterDescriptor;
 import org.apache.iotdb.cluster.coordinator.Coordinator;
 import org.apache.iotdb.cluster.exception.ConfigInconsistentException;
@@ -64,10 +65,8 @@ public class MetaClusterServer extends RaftServer
   // each node only contains one MetaGroupMember
   private MetaGroupMember member;
   private Coordinator coordinator;
-  // the single-node IoTDB instance
-  private IoTDB ioTDB;
-  // to register the ClusterMonitor that helps monitoring the cluster
-  private RegisterManager registerManager = new RegisterManager();
+
+
   private MetaAsyncService asyncService;
   private MetaSyncService syncService;
   private MetaHeartbeatServer metaHeartbeatServer;
@@ -94,28 +93,24 @@ public class MetaClusterServer extends RaftServer
   public void start() throws TTransportException, StartupException {
     super.start();
     metaHeartbeatServer.start();
-    ioTDB = new IoTDB();
+
     IoTDB.setMetaManager(CMManager.getInstance());
     ((CMManager) IoTDB.metaManager).setMetaGroupMember(member);
     ((CMManager) IoTDB.metaManager).setCoordinator(coordinator);
-    ioTDB.active();
+    //TODO FIXME move this out of MetaClusterServer
+    IoTDB.getInstance().active();
+
     member.start();
-    // JMX based DBA API
-    registerManager.register(ClusterMonitor.INSTANCE);
+
   }
 
   /** Also stops the IoTDB instance, the MetaGroupMember and the 
ClusterMonitor. */
   @Override
   public void stop() {
-    if (ioTDB == null) {
-      return;
-    }
+
     metaHeartbeatServer.stop();
     super.stop();
-    ioTDB.stop();
-    ioTDB = null;
     member.stop();
-    registerManager.deregisterAll();
   }
 
   /** Build a initial cluster with other nodes on the seed list. */
@@ -371,8 +366,4 @@ public class MetaClusterServer extends RaftServer
     this.member = metaGroupMember;
   }
 
-  @TestOnly
-  public IoTDB getIoTDB() {
-    return ioTDB;
-  }
 }
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 b78d388..39ebf3c 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
@@ -65,7 +65,7 @@ import 
org.apache.iotdb.cluster.rpc.thrift.SendSnapshotRequest;
 import org.apache.iotdb.cluster.rpc.thrift.StartUpStatus;
 import org.apache.iotdb.cluster.rpc.thrift.TSMetaService;
 import org.apache.iotdb.cluster.rpc.thrift.TSMetaService.AsyncClient;
-import org.apache.iotdb.cluster.server.ClientServer;
+import org.apache.iotdb.cluster.server.ClusterTSServiceImpl;
 import org.apache.iotdb.cluster.server.DataClusterServer;
 import org.apache.iotdb.cluster.server.HardLinkCleaner;
 import org.apache.iotdb.cluster.server.NodeCharacter;
@@ -213,7 +213,7 @@ public class MetaGroupMember extends RaftMember {
    * an override of TSServiceImpl, which redirect JDBC and Session requests to 
the MetaGroupMember
    * so they can be processed cluster-wide
    */
-  private ClientServer clientServer;
+  private ClusterTSServiceImpl clusterTSServiceImpl;
 
   private DataClientProvider dataClientProvider;
 
@@ -276,7 +276,7 @@ public class MetaGroupMember extends RaftMember {
     Factory dataMemberFactory = new Factory(factory, this);
     dataClusterServer = new DataClusterServer(thisNode, dataMemberFactory, 
this);
     dataHeartbeatServer = new DataHeartbeatServer(thisNode, dataClusterServer);
-    clientServer = new ClientServer(this);
+    clusterTSServiceImpl = new ClusterTSServiceImpl(this);
     startUpStatus = getNewStartUpStatus();
 
     // try loading the partition table if there was a previous cluster
@@ -333,7 +333,7 @@ public class MetaGroupMember extends RaftMember {
   }
 
   /**
-   * Stop the heartbeat and catch-up thread pool, DataClusterServer, 
ClientServer and reportThread.
+   * Stop the heartbeat and catch-up thread pool, DataClusterServer, 
ClusterTSServiceImpl and reportThread.
    * Calling the method twice does not induce side effects.
    */
   @Override
@@ -345,8 +345,8 @@ public class MetaGroupMember extends RaftMember {
     if (getDataHeartbeatServer() != null) {
       getDataHeartbeatServer().stop();
     }
-    if (clientServer != null) {
-      clientServer.stop();
+    if (clusterTSServiceImpl != null) {
+      clusterTSServiceImpl.stop();
     }
     if (reportThread != null) {
       reportThread.shutdownNow();
@@ -371,14 +371,14 @@ public class MetaGroupMember extends RaftMember {
   }
 
   /**
-   * Start DataClusterServer and ClientServer so this node will be able to 
respond to other nodes
+   * Start DataClusterServer and ClusterTSServiceImpl so this node will be 
able to respond to other nodes
    * and clients.
    */
   protected void initSubServers() throws TTransportException, StartupException 
{
     getDataClusterServer().start();
     getDataHeartbeatServer().start();
-    clientServer.setCoordinator(this.coordinator);
-    clientServer.start();
+    clusterTSServiceImpl.setCoordinator(this.coordinator);
+    clusterTSServiceImpl.start();
   }
 
   /**
@@ -549,7 +549,7 @@ public class MetaGroupMember extends RaftMember {
 
   /**
    * Send a join cluster request to "node". If the joining is accepted, set 
the partition table,
-   * start DataClusterServer and ClientServer and initialize DataGroupMembers.
+   * start DataClusterServer and ClusterTSServiceImpl and initialize 
DataGroupMembers.
    *
    * @return true if the node has successfully joined the cluster, false 
otherwise.
    */
@@ -669,7 +669,7 @@ public class MetaGroupMember extends RaftMember {
 
   /**
    * Deserialize a partition table from the buffer, save it locally, add nodes 
from the partition
-   * table and start DataClusterServer and ClientServer.
+   * table and start DataClusterServer and ClusterTSServiceImpl.
    */
   public synchronized void acceptPartitionTable(
       ByteBuffer partitionTableBuffer, boolean needSerialization) {
@@ -715,7 +715,7 @@ public class MetaGroupMember extends RaftMember {
   /**
    * Process a HeartBeatResponse from a follower. If the follower has provided 
its identifier, try
    * registering for it and if all nodes have registered and there is no 
available partition table,
-   * initialize a new one and start the ClientServer and DataClusterServer. If 
the follower requires
+   * initialize a new one and start the ClusterTSServiceImpl and 
DataClusterServer. If the follower requires
    * a partition table, add it to the blind node list so that at the next 
heartbeat this node will
    * send it a partition table
    */
@@ -800,7 +800,7 @@ public class MetaGroupMember extends RaftMember {
   }
 
   /**
-   * Start the DataClusterServer and ClientServer so this node can serve other 
nodes and clients.
+   * Start the DataClusterServer and ClusterTSServiceImpl` so this node can 
serve other nodes and clients.
    * Also build DataGroupMembers using the partition table.
    */
   protected synchronized void startSubServers() {
@@ -1814,8 +1814,8 @@ public class MetaGroupMember extends RaftMember {
                     // ignore
                   }
                   super.stop();
-                  if (clientServer != null) {
-                    clientServer.stop();
+                  if (clusterTSServiceImpl != null) {
+                    clusterTSServiceImpl.stop();
                   }
                   logger.info("{} has been removed from the cluster", name);
                 })
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/utils/nodetool/ClusterMonitor.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/utils/nodetool/ClusterMonitor.java
index ce941f5..9250785 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/utils/nodetool/ClusterMonitor.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/utils/nodetool/ClusterMonitor.java
@@ -18,7 +18,7 @@
  */
 package org.apache.iotdb.cluster.utils.nodetool;
 
-import org.apache.iotdb.cluster.ClusterMain;
+import org.apache.iotdb.cluster.ClusterIoTDB;
 import org.apache.iotdb.cluster.config.ClusterConstant;
 import org.apache.iotdb.cluster.config.ClusterDescriptor;
 import org.apache.iotdb.cluster.partition.PartitionGroup;
@@ -201,7 +201,7 @@ public class ClusterMonitor implements ClusterMonitorMBean, 
IService {
   }
 
   private MetaGroupMember getMetaGroupMember() {
-    MetaClusterServer metaClusterServer = ClusterMain.getMetaServer();
+    MetaClusterServer metaClusterServer = 
ClusterIoTDB.getInstance().getMetaServer();
     if (metaClusterServer == null) {
       return null;
     }
diff --git 
a/cluster/src/test/java/org/apache/iotdb/cluster/integration/BaseSingleNodeTest.java
 
b/cluster/src/test/java/org/apache/iotdb/cluster/integration/BaseSingleNodeTest.java
index 70fbb66..7565e9d 100644
--- 
a/cluster/src/test/java/org/apache/iotdb/cluster/integration/BaseSingleNodeTest.java
+++ 
b/cluster/src/test/java/org/apache/iotdb/cluster/integration/BaseSingleNodeTest.java
@@ -51,6 +51,7 @@ public abstract class BaseSingleNodeTest {
 
   @After
   public void tearDown() throws Exception {
+    //TODO fixme
     metaServer.stop();
     recoverConfigs();
     EnvironmentUtils.cleanEnv();
diff --git 
a/cluster/src/test/java/org/apache/iotdb/cluster/server/clusterinfo/ClusterInfoServiceImplTest.java
 
b/cluster/src/test/java/org/apache/iotdb/cluster/server/clusterinfo/ClusterInfoServiceImplTest.java
index 5dde20f..499efce 100644
--- 
a/cluster/src/test/java/org/apache/iotdb/cluster/server/clusterinfo/ClusterInfoServiceImplTest.java
+++ 
b/cluster/src/test/java/org/apache/iotdb/cluster/server/clusterinfo/ClusterInfoServiceImplTest.java
@@ -19,7 +19,7 @@
 
 package org.apache.iotdb.cluster.server.clusterinfo;
 
-import org.apache.iotdb.cluster.ClusterMain;
+import org.apache.iotdb.cluster.ClusterIoTDB;
 import org.apache.iotdb.cluster.rpc.thrift.DataPartitionEntry;
 import org.apache.iotdb.cluster.rpc.thrift.Node;
 import org.apache.iotdb.cluster.server.MetaClusterServer;
@@ -52,7 +52,7 @@ public class ClusterInfoServiceImplTest {
     metaClusterServer.getMember().stop();
     metaClusterServer.setMetaGroupMember(metaGroupMember);
 
-    ClusterMain.setMetaClusterServer(metaClusterServer);
+    ClusterIoTDB.setMetaClusterServer(metaClusterServer);
 
     metaClusterServer.getIoTDB().metaManager.setStorageGroup(new 
PartialPath("root", "sg"));
     // metaClusterServer.getMember()
@@ -61,11 +61,11 @@ public class ClusterInfoServiceImplTest {
 
   @After
   public void tearDown() throws MetadataException {
-    ClusterMain.getMetaServer()
+    ClusterIoTDB.getMetaServer()
         .getIoTDB()
         .metaManager
         .deleteStorageGroups(Collections.singletonList(new PartialPath("root", 
"sg")));
-    ClusterMain.getMetaServer().stop();
+    ClusterIoTDB.getMetaServer().stop();
   }
 
   @Test
diff --git 
a/server/src/main/java/org/apache/iotdb/db/concurrent/ThreadName.java 
b/server/src/main/java/org/apache/iotdb/db/concurrent/ThreadName.java
index 0850b75..bbbd6f1 100644
--- a/server/src/main/java/org/apache/iotdb/db/concurrent/ThreadName.java
+++ b/server/src/main/java/org/apache/iotdb/db/concurrent/ThreadName.java
@@ -46,7 +46,9 @@ public enum ThreadName {
   QUERY_SERVICE("Query"),
   WINDOW_EVALUATION_SERVICE("WindowEvaluationTaskPoolManager"),
   CONTINUOUS_QUERY_SERVICE("ContinuousQueryTaskPoolManager"),
-  CLUSTER_INFO_SERVICE("ClusterInfoClient");
+  CLUSTER_INFO_SERVICE("ClusterInfoClient"),
+  CLUSTER_RPC_SERVICE("ClusterRPC"),
+  CLUSTER_RPC_CLIENT("Cluster-RPC-Client");
 
   private final String name;
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java 
b/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java
index ebb26e6..4e0649e 100644
--- a/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java
+++ b/server/src/main/java/org/apache/iotdb/db/qp/executor/PlanExecutor.java
@@ -471,9 +471,7 @@ public class PlanExecutor implements IPlanExecutor {
     if (!plan.isTracingOn()) {
       TracingManager.getInstance().close();
     } else {
-      if (!TracingManager.getInstance().getWriterStatus()) {
-        TracingManager.getInstance().openTracingWriteStream();
-      }
+      TracingManager.getInstance().initTracingManager();
     }
   }
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/query/control/TracingManager.java 
b/server/src/main/java/org/apache/iotdb/db/query/control/TracingManager.java
index 450c4f4..6b61323 100644
--- a/server/src/main/java/org/apache/iotdb/db/query/control/TracingManager.java
+++ b/server/src/main/java/org/apache/iotdb/db/query/control/TracingManager.java
@@ -44,11 +44,18 @@ public class TracingManager {
   private BufferedWriter writer;
   private Map<Long, Long> queryStartTime = new ConcurrentHashMap<>();
 
-  public TracingManager(String dirName, String logFileName) {
-    initTracingManager(dirName, logFileName);
+  private TracingManager() {
+    initTracingManager();
   }
 
-  public void initTracingManager(String dirName, String logFileName) {
+  public void initTracingManager() {
+    if (this.writer != null) {
+      //the tracing manager has been initialized.
+      return;
+    }
+    String dirName = IoTDBDescriptor.getInstance().getConfig().getTracingDir();
+    String logFileName = IoTDBConstant.TRACING_LOG;
+
     File tracingDir = SystemFileFactory.INSTANCE.getFile(dirName);
     if (!tracingDir.exists()) {
       if (tracingDir.mkdirs()) {
@@ -203,31 +210,17 @@ public class TracingManager {
 
   public void close() {
     try {
+      queryStartTime.clear();
       writer.close();
+      writer = null;
     } catch (IOException e) {
       logger.error("Meeting error while Close the tracing log stream : {}", 
e.getMessage());
     }
   }
 
-  public boolean getWriterStatus() {
-    try {
-      writer.flush();
-      return true;
-    } catch (IOException e) {
-      return false;
-    }
-  }
-
-  public void openTracingWriteStream() {
-    initTracingManager(
-        IoTDBDescriptor.getInstance().getConfig().getTracingDir(), 
IoTDBConstant.TRACING_LOG);
-  }
-
   private static class TracingManagerHelper {
 
-    private static final TracingManager INSTANCE =
-        new TracingManager(
-            IoTDBDescriptor.getInstance().getConfig().getTracingDir(), 
IoTDBConstant.TRACING_LOG);
+    private static final TracingManager INSTANCE = new TracingManager();
 
     private TracingManagerHelper() {}
   }
diff --git a/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java 
b/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
index 1154e3d..9cd89e6 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/IoTDB.java
@@ -152,6 +152,14 @@ public class IoTDB implements IoTDBMBean {
 
   private void deactivate() {
     logger.info("Deactivating IoTDB...");
+    //some user may call Tracing on but do not close tracing.
+    //so, when remove the system, we have to close the tracing
+    if 
(IoTDBDescriptor.getInstance().getConfig().isEnablePerformanceTracing()) {
+      TracingManager.getInstance().close();
+    }
+    PrimitiveArrayManager.close();
+    SystemInfo.getInstance().close();
+
     registerManager.deregisterAll();
     JMXService.deregisterMBean(mbeanName);
     logger.info("IoTDB is deactivated.");
@@ -174,15 +182,9 @@ public class IoTDB implements IoTDBMBean {
   }
 
   public void shutdown() throws Exception {
-    logger.info("Deactivating IoTDB...");
-    if 
(IoTDBDescriptor.getInstance().getConfig().isEnablePerformanceTracing()) {
-      TracingManager.getInstance().close();
-    }
-    registerManager.shutdownAll();
-    PrimitiveArrayManager.close();
-    SystemInfo.getInstance().close();
-    JMXService.deregisterMBean(mbeanName);
-    logger.info("IoTDB is deactivated.");
+    stop();
+
+    logger.info("IoTDB is shutdown.");
   }
 
   private void setUncaughtExceptionHandler() {
diff --git a/server/src/main/java/org/apache/iotdb/db/service/RPCService.java 
b/server/src/main/java/org/apache/iotdb/db/service/RPCService.java
index 5bbddec..be5f2f7 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/RPCService.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/RPCService.java
@@ -31,19 +31,11 @@ public class RPCService extends ThriftService implements 
RPCServiceMBean {
 
   private TSServiceImpl impl;
 
-  private RPCService() {}
-
   public static RPCService getInstance() {
     return RPCServiceHolder.INSTANCE;
   }
 
   @Override
-  public int getRPCPort() {
-    IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
-    return config.getRpcPort();
-  }
-
-  @Override
   public ThriftService getImplementation() {
     return getInstance();
   }
@@ -94,6 +86,11 @@ public class RPCService extends ThriftService implements 
RPCServiceMBean {
     return ServiceType.RPC_SERVICE;
   }
 
+  @Override
+  public int getRPCPort() {
+    return getBindPort();
+  }
+
   private static class RPCServiceHolder {
 
     private static final RPCService INSTANCE = new RPCService();
diff --git 
a/server/src/main/java/org/apache/iotdb/db/service/RPCServiceThriftHandler.java 
b/server/src/main/java/org/apache/iotdb/db/service/RPCServiceThriftHandler.java
index b0f51bd..dacde67 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/service/RPCServiceThriftHandler.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/service/RPCServiceThriftHandler.java
@@ -24,7 +24,7 @@ import org.apache.thrift.transport.TTransport;
 public class RPCServiceThriftHandler implements TServerEventHandler {
   private TSServiceImpl serviceImpl;
 
-  RPCServiceThriftHandler(TSServiceImpl serviceImpl) {
+  public RPCServiceThriftHandler(TSServiceImpl serviceImpl) {
     this.serviceImpl = serviceImpl;
   }
 
diff --git a/server/src/main/java/org/apache/iotdb/db/service/ServiceType.java 
b/server/src/main/java/org/apache/iotdb/db/service/ServiceType.java
index 2f61d90..febacac 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/ServiceType.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/ServiceType.java
@@ -55,6 +55,9 @@ public enum ServiceType {
   SYSTEMINFO_SERVICE("MemTable Monitor Service", "MemTable, Monitor"),
   CONTINUOUS_QUERY_SERVICE("Continuous Query Service", "Continuous Query 
Service"),
   CLUSTER_INFO_SERVICE("Cluster Monitor Service (thrift-based)", "Cluster 
Monitor-Thrift"),
+
+  CLUSTER_RPC_SERVICE("Cluster RPC ServerService", "ClusterRPCService"),
+
   ;
 
   private final String name;
diff --git 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftService.java 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftService.java
index d975743..dfb2526 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftService.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftService.java
@@ -64,11 +64,6 @@ public abstract class ThriftService implements IService {
     }
   }
 
-  public int getRPCPort() {
-    IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
-    return config.getRpcPort();
-  }
-
   public abstract ThriftService getImplementation();
 
   @Override
diff --git 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftServiceThread.java
 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftServiceThread.java
index 612d187..2564d05 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftServiceThread.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/service/thrift/ThriftServiceThread.java
@@ -61,7 +61,7 @@ public class ThriftServiceThread extends Thread {
       String bindAddress,
       int port,
       int maxWorkerThreads,
-      int timeoutMs,
+      int timeoutSecond,
       TServerEventHandler serverEventHandler,
       boolean compress) {
     if (compress) {
@@ -77,7 +77,7 @@ public class ThriftServiceThread extends Thread {
           new TThreadPoolServer.Args(serverTransport)
               .maxWorkerThreads(maxWorkerThreads)
               .minWorkerThreads(CommonUtils.getCpuCores())
-              .stopTimeoutVal(timeoutMs);
+              .stopTimeoutVal(timeoutSecond);
       poolArgs.executorService =
           IoTDBThreadPoolFactory.createThriftRpcClientThreadPool(poolArgs, 
threadsName);
       poolArgs.processor(processor);
diff --git 
a/server/src/test/java/org/apache/iotdb/db/query/control/TracingManagerTest.java
 
b/server/src/test/java/org/apache/iotdb/db/query/control/TracingManagerTest.java
index 3f1f0a9..c1783e3 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/query/control/TracingManagerTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/query/control/TracingManagerTest.java
@@ -64,9 +64,6 @@ public class TracingManagerTest {
 
   @Test
   public void tracingQueryTest() throws IOException {
-    if (!tracingManager.getWriterStatus()) {
-      tracingManager.openTracingWriteStream();
-    }
     String[] ans = {
       "Query Id: 10 - Query Statement: " + sql,
       "Query Id: 10 - Start time: 2020-12-",

Reply via email to