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-",
