This is an automated email from the ASF dual-hosted git repository.
qiaojialin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 54d0720188 uncomment the data_block_manager_port and change default
data region to StandAlone Consensus (#6317)
54d0720188 is described below
commit 54d0720188d5a7f6ccc062a6e24770a799bb8fe6
Author: Jialin Qiao <[email protected]>
AuthorDate: Fri Jun 17 09:58:23 2022 +0800
uncomment the data_block_manager_port and change default data region to
StandAlone Consensus (#6317)
---
.../resources/conf/iotdb-confignode.properties | 2 +-
.../client/ConfigNodeClientPoolFactory.java | 4 ++--
.../{ConfigNodeConf.java => ConfigNodeConfig.java} | 8 ++++----
.../iotdb/confignode/conf/ConfigNodeDescriptor.java | 4 ++--
.../confignode/conf/ConfigNodeStartupCheck.java | 2 +-
.../iotdb/confignode/manager/ConfigManager.java | 4 ++--
.../iotdb/confignode/manager/ConsensusManager.java | 4 ++--
.../iotdb/confignode/manager/PartitionManager.java | 4 ++--
.../iotdb/confignode/manager/ProcedureManager.java | 11 ++++++-----
.../apache/iotdb/confignode/persistence/UDFInfo.java | 4 ++--
.../apache/iotdb/confignode/service/ConfigNode.java | 10 +++++-----
.../service/thrift/ConfigNodeRPCService.java | 4 ++--
.../thrift/ConfigNodeRPCServiceProcessorTest.java | 10 +++++-----
.../assembly/resources/conf/iotdb-engine.properties | 2 +-
.../db/mpp/plan/scheduler/ClusterScheduler.java | 20 ++++++--------------
.../java/org/apache/iotdb/db/service/DataNode.java | 16 ++++++----------
16 files changed, 49 insertions(+), 60 deletions(-)
diff --git a/confignode/src/assembly/resources/conf/iotdb-confignode.properties
b/confignode/src/assembly/resources/conf/iotdb-confignode.properties
index f067335003..20b96dad33 100644
--- a/confignode/src/assembly/resources/conf/iotdb-confignode.properties
+++ b/confignode/src/assembly/resources/conf/iotdb-confignode.properties
@@ -60,7 +60,7 @@ target_confignode=0.0.0.0:22277
# 2. org.apache.iotdb.consensus.ratis.RatisConsensus(Raft protocol)
# 3. org.apache.iotdb.consensus.multileader.MultiLeaderConsensus(weak
consistency, high performance)
# Datatype: String
-#
data_region_consensus_protocol_class=org.apache.iotdb.consensus.multileader.MultiLeaderConsensus
+#
data_region_consensus_protocol_class=org.apache.iotdb.consensus.standalone.StandAloneConsensus
# SchemaRegion consensus protocol type
# These consensus protocols are currently supported:
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/client/ConfigNodeClientPoolFactory.java
b/confignode/src/main/java/org/apache/iotdb/confignode/client/ConfigNodeClientPoolFactory.java
index 57f9ca6ef6..f554a851c5 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/client/ConfigNodeClientPoolFactory.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/client/ConfigNodeClientPoolFactory.java
@@ -26,7 +26,7 @@ import org.apache.iotdb.commons.client.ClientPoolProperty;
import org.apache.iotdb.commons.client.IClientPoolFactory;
import
org.apache.iotdb.commons.client.async.AsyncDataNodeInternalServiceClient;
import org.apache.iotdb.commons.client.sync.SyncDataNodeInternalServiceClient;
-import org.apache.iotdb.confignode.conf.ConfigNodeConf;
+import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
import org.apache.commons.pool2.KeyedObjectPool;
@@ -34,7 +34,7 @@ import org.apache.commons.pool2.impl.GenericKeyedObjectPool;
public class ConfigNodeClientPoolFactory {
- private static final ConfigNodeConf conf =
ConfigNodeDescriptor.getInstance().getConf();
+ private static final ConfigNodeConfig conf =
ConfigNodeDescriptor.getInstance().getConf();
private ConfigNodeClientPoolFactory() {}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConf.java
b/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConfig.java
similarity index 98%
rename from
confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConf.java
rename to
confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConfig.java
index 1d601a7b03..47fb03c2d2 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConf.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConfig.java
@@ -29,7 +29,7 @@ import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.TimeUnit;
-public class ConfigNodeConf {
+public class ConfigNodeConfig {
/** could set ip or hostname */
private String rpcAddress = "0.0.0.0";
@@ -59,7 +59,7 @@ public class ConfigNodeConf {
private final String configNodeConsensusProtocolClass =
ConsensusFactory.RatisConsensus;
/** DataNode data region consensus protocol */
- private String dataRegionConsensusProtocolClass =
ConsensusFactory.MultiLeaderConsensus;
+ private String dataRegionConsensusProtocolClass =
ConsensusFactory.StandAloneConsensus;
/** DataNode schema region consensus protocol */
private String schemaRegionConsensusProtocolClass =
ConsensusFactory.StandAloneConsensus;
@@ -142,7 +142,7 @@ public class ConfigNodeConf {
/** This parameter only exists for a few days */
private boolean enableHeartbeat = true;
- ConfigNodeConf() {
+ ConfigNodeConfig() {
// empty constructor
}
@@ -298,7 +298,7 @@ public class ConfigNodeConf {
return connectionTimeoutInMS;
}
- public ConfigNodeConf setConnectionTimeoutInMS(int connectionTimeoutInMS) {
+ public ConfigNodeConfig setConnectionTimeoutInMS(int connectionTimeoutInMS) {
this.connectionTimeoutInMS = connectionTimeoutInMS;
return this;
}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeDescriptor.java
b/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeDescriptor.java
index 515669ff19..28305194c4 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeDescriptor.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeDescriptor.java
@@ -38,13 +38,13 @@ public class ConfigNodeDescriptor {
private final CommonDescriptor commonDescriptor =
CommonDescriptor.getInstance();
- private final ConfigNodeConf conf = new ConfigNodeConf();
+ private final ConfigNodeConfig conf = new ConfigNodeConfig();
private ConfigNodeDescriptor() {
loadProps();
}
- public ConfigNodeConf getConf() {
+ public ConfigNodeConfig getConf() {
return conf;
}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeStartupCheck.java
b/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeStartupCheck.java
index 8a2defa70d..40d7fb2113 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeStartupCheck.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeStartupCheck.java
@@ -50,7 +50,7 @@ public class ConfigNodeStartupCheck {
private static final Logger LOGGER =
LoggerFactory.getLogger(ConfigNodeStartupCheck.class);
- private static final ConfigNodeConf conf =
ConfigNodeDescriptor.getInstance().getConf();
+ private static final ConfigNodeConfig conf =
ConfigNodeDescriptor.getInstance().getConf();
private final File systemPropertiesFile;
private final Properties systemProperties;
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
index e234d8e30e..7621a1e4ae 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
@@ -28,7 +28,7 @@ import org.apache.iotdb.commons.partition.SchemaPartition;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.utils.AuthUtils;
import org.apache.iotdb.commons.utils.PathUtils;
-import org.apache.iotdb.confignode.conf.ConfigNodeConf;
+import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
import org.apache.iotdb.confignode.consensus.request.auth.AuthorReq;
import org.apache.iotdb.confignode.consensus.request.read.CountStorageGroupReq;
@@ -552,7 +552,7 @@ public class ConfigManager implements Manager {
@Override
public TConfigNodeRegisterResp registerConfigNode(TConfigNodeRegisterReq
req) {
// Check global configuration
- ConfigNodeConf conf = ConfigNodeDescriptor.getInstance().getConf();
+ ConfigNodeConfig conf = ConfigNodeDescriptor.getInstance().getConf();
TConfigNodeRegisterResp errorResp = new TConfigNodeRegisterResp();
errorResp.setStatus(new
TSStatus(TSStatusCode.ERROR_GLOBAL_CONFIG.getStatusCode()));
if (!req.getDataRegionConsensusProtocolClass()
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConsensusManager.java
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConsensusManager.java
index 0a666e428b..95c82d59d4 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConsensusManager.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConsensusManager.java
@@ -25,7 +25,7 @@ import org.apache.iotdb.commons.consensus.ConsensusGroupId;
import org.apache.iotdb.commons.consensus.PartitionRegionId;
import org.apache.iotdb.commons.utils.TestOnly;
import org.apache.iotdb.confignode.client.SyncConfigNodeClientPool;
-import org.apache.iotdb.confignode.conf.ConfigNodeConf;
+import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
import org.apache.iotdb.confignode.consensus.request.ConfigRequest;
import org.apache.iotdb.confignode.consensus.request.write.ApplyConfigNodeReq;
@@ -50,7 +50,7 @@ import java.util.concurrent.TimeUnit;
public class ConsensusManager {
private static final Logger LOGGER =
LoggerFactory.getLogger(ConsensusManager.class);
- private static final ConfigNodeConf conf =
ConfigNodeDescriptor.getInstance().getConf();
+ private static final ConfigNodeConfig conf =
ConfigNodeDescriptor.getInstance().getConf();
private ConsensusGroupId consensusGroupId;
private IConsensus consensusImpl;
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/PartitionManager.java
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/PartitionManager.java
index 931f16841e..718960d541 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/PartitionManager.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/PartitionManager.java
@@ -30,7 +30,7 @@ import org.apache.iotdb.commons.partition.DataPartitionTable;
import org.apache.iotdb.commons.partition.SchemaPartitionTable;
import org.apache.iotdb.commons.partition.executor.SeriesPartitionExecutor;
import org.apache.iotdb.confignode.client.SyncDataNodeClientPool;
-import org.apache.iotdb.confignode.conf.ConfigNodeConf;
+import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
import org.apache.iotdb.confignode.consensus.request.read.GetDataPartitionReq;
import
org.apache.iotdb.confignode.consensus.request.read.GetNodePathsPartitionReq;
@@ -94,7 +94,7 @@ public class PartitionManager {
/** Construct SeriesPartitionExecutor by iotdb-confignode.propertis */
private void setSeriesPartitionExecutor() {
- ConfigNodeConf conf = ConfigNodeDescriptor.getInstance().getConf();
+ ConfigNodeConfig conf = ConfigNodeDescriptor.getInstance().getConf();
this.executor =
SeriesPartitionExecutor.getSeriesPartitionExecutor(
conf.getSeriesPartitionExecutorClass(),
conf.getSeriesPartitionSlotNum());
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
index baedc62cd2..fe2f21d87b 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
@@ -21,7 +21,7 @@ package org.apache.iotdb.confignode.manager;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.utils.StatusUtils;
-import org.apache.iotdb.confignode.conf.ConfigNodeConf;
+import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
import org.apache.iotdb.confignode.persistence.ProcedureInfo;
import org.apache.iotdb.confignode.procedure.Procedure;
@@ -46,7 +46,8 @@ import java.util.concurrent.TimeUnit;
public class ProcedureManager {
private static final Logger LOGGER =
LoggerFactory.getLogger(ProcedureManager.class);
- private static final ConfigNodeConf configNodeConf =
ConfigNodeDescriptor.getInstance().getConf();
+ private static final ConfigNodeConfig CONFIG_NODE_CONFIG =
+ ConfigNodeDescriptor.getInstance().getConf();
private static final int procedureWaitTimeOut = 30;
private static final int procedureWaitRetryTimeout = 250;
@@ -68,11 +69,11 @@ public class ProcedureManager {
public void shiftExecutor(boolean running) {
if (running) {
if (!executor.isRunning()) {
- executor.init(configNodeConf.getProcedureCoreWorkerThreadsSize());
+ executor.init(CONFIG_NODE_CONFIG.getProcedureCoreWorkerThreadsSize());
executor.startWorkers();
executor.startCompletedCleaner(
- configNodeConf.getProcedureCompletedCleanInterval(),
- configNodeConf.getProcedureCompletedEvictTTL());
+ CONFIG_NODE_CONFIG.getProcedureCompletedCleanInterval(),
+ CONFIG_NODE_CONFIG.getProcedureCompletedEvictTTL());
store.start();
}
} else {
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/UDFInfo.java
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/UDFInfo.java
index cc4188bf0c..ea3a17420d 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/persistence/UDFInfo.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/persistence/UDFInfo.java
@@ -25,7 +25,7 @@ import org.apache.iotdb.commons.udf.service.UDFClassLoader;
import org.apache.iotdb.commons.udf.service.UDFExecutableManager;
import org.apache.iotdb.commons.udf.service.UDFExecutableResource;
import org.apache.iotdb.commons.udf.service.UDFRegistrationService;
-import org.apache.iotdb.confignode.conf.ConfigNodeConf;
+import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
import org.apache.iotdb.confignode.consensus.request.write.CreateFunctionReq;
import org.apache.iotdb.confignode.consensus.request.write.DropFunctionReq;
@@ -42,7 +42,7 @@ public class UDFInfo implements SnapshotProcessor {
private static final Logger LOGGER = LoggerFactory.getLogger(UDFInfo.class);
- private static final ConfigNodeConf CONFIG_NODE_CONF =
+ private static final ConfigNodeConfig CONFIG_NODE_CONF =
ConfigNodeDescriptor.getInstance().getConf();
private final UDFExecutableManager udfExecutableManager;
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNode.java
b/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNode.java
index 2895454c4f..bcf5f20cfd 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNode.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNode.java
@@ -24,7 +24,7 @@ import org.apache.iotdb.commons.service.RegisterManager;
import org.apache.iotdb.commons.udf.service.UDFClassLoaderManager;
import org.apache.iotdb.commons.udf.service.UDFExecutableManager;
import org.apache.iotdb.commons.udf.service.UDFRegistrationService;
-import org.apache.iotdb.confignode.conf.ConfigNodeConf;
+import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
import org.apache.iotdb.confignode.conf.ConfigNodeConstant;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
import org.apache.iotdb.confignode.manager.ConfigManager;
@@ -94,14 +94,14 @@ public class ConfigNode implements ConfigNodeMBean {
}
private void registerUdfServices() throws StartupException {
- final ConfigNodeConf configNodeConf =
ConfigNodeDescriptor.getInstance().getConf();
+ final ConfigNodeConfig configNodeConfig =
ConfigNodeDescriptor.getInstance().getConf();
registerManager.register(
UDFExecutableManager.setupAndGetInstance(
- configNodeConf.getTemporaryLibDir(),
configNodeConf.getUdfLibDir()));
+ configNodeConfig.getTemporaryLibDir(),
configNodeConfig.getUdfLibDir()));
registerManager.register(
-
UDFClassLoaderManager.setupAndGetInstance(configNodeConf.getUdfLibDir()));
+
UDFClassLoaderManager.setupAndGetInstance(configNodeConfig.getUdfLibDir()));
registerManager.register(
-
UDFRegistrationService.setupAndGetInstance(configNodeConf.getSystemUdfDir()));
+
UDFRegistrationService.setupAndGetInstance(configNodeConfig.getSystemUdfDir()));
}
public void active() {
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCService.java
b/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCService.java
index 7f73117c61..8ba2b961ec 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCService.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCService.java
@@ -24,14 +24,14 @@ import
org.apache.iotdb.commons.exception.runtime.RPCServiceException;
import org.apache.iotdb.commons.service.ServiceType;
import org.apache.iotdb.commons.service.ThriftService;
import org.apache.iotdb.commons.service.ThriftServiceThread;
-import org.apache.iotdb.confignode.conf.ConfigNodeConf;
+import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
import org.apache.iotdb.confignode.rpc.thrift.ConfigIService;
/** ConfigNodeRPCServer exposes the interface that interacts with the DataNode
*/
public class ConfigNodeRPCService extends ThriftService implements
ConfigNodeRPCServiceMBean {
- private static final ConfigNodeConf conf =
ConfigNodeDescriptor.getInstance().getConf();
+ private static final ConfigNodeConfig conf =
ConfigNodeDescriptor.getInstance().getConf();
private ConfigNodeRPCServiceProcessor configNodeRPCServiceProcessor;
diff --git
a/confignode/src/test/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessorTest.java
b/confignode/src/test/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessorTest.java
index 0109d1e8f2..9d972616b8 100644
---
a/confignode/src/test/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessorTest.java
+++
b/confignode/src/test/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessorTest.java
@@ -37,7 +37,7 @@ import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.udf.service.UDFClassLoaderManager;
import org.apache.iotdb.commons.udf.service.UDFExecutableManager;
import org.apache.iotdb.commons.udf.service.UDFRegistrationService;
-import org.apache.iotdb.confignode.conf.ConfigNodeConf;
+import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
import org.apache.iotdb.confignode.conf.ConfigNodeConstant;
import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor;
import org.apache.iotdb.confignode.conf.ConfigNodeStartupCheck;
@@ -98,11 +98,11 @@ public class ConfigNodeRPCServiceProcessorTest {
@BeforeClass
public static void beforeClass() throws StartupException,
ConfigurationException, IOException {
- final ConfigNodeConf configNodeConf =
ConfigNodeDescriptor.getInstance().getConf();
+ final ConfigNodeConfig configNodeConfig =
ConfigNodeDescriptor.getInstance().getConf();
UDFExecutableManager.setupAndGetInstance(
- configNodeConf.getTemporaryLibDir(), configNodeConf.getUdfLibDir());
- UDFClassLoaderManager.setupAndGetInstance(configNodeConf.getUdfLibDir());
-
UDFRegistrationService.setupAndGetInstance(configNodeConf.getSystemUdfDir());
+ configNodeConfig.getTemporaryLibDir(),
configNodeConfig.getUdfLibDir());
+ UDFClassLoaderManager.setupAndGetInstance(configNodeConfig.getUdfLibDir());
+
UDFRegistrationService.setupAndGetInstance(configNodeConfig.getSystemUdfDir());
ConfigNodeStartupCheck.getInstance().startUpCheck();
}
diff --git a/server/src/assembly/resources/conf/iotdb-engine.properties
b/server/src/assembly/resources/conf/iotdb-engine.properties
index 01f12621fa..6345b92d76 100644
--- a/server/src/assembly/resources/conf/iotdb-engine.properties
+++ b/server/src/assembly/resources/conf/iotdb-engine.properties
@@ -31,7 +31,7 @@ rpc_port=6667
### Shuffle Configuration
####################
# Datatype: int
-# data_block_manager_port=8777
+data_block_manager_port=8777
# Datatype: int
# data_block_manager_core_pool_size=1
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/ClusterScheduler.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/ClusterScheduler.java
index 398f0f8272..9437c22536 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/ClusterScheduler.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/scheduler/ClusterScheduler.java
@@ -50,20 +50,15 @@ import java.util.concurrent.ScheduledExecutorService;
public class ClusterScheduler implements IScheduler {
private static final Logger logger =
LoggerFactory.getLogger(ClusterScheduler.class);
- private MPPQueryContext queryContext;
// The stateMachine of the QueryExecution owned by this QueryScheduler
- private QueryStateMachine stateMachine;
- private QueryType queryType;
+ private final QueryStateMachine stateMachine;
+ private final QueryType queryType;
// The fragment instances which should be sent to corresponding Nodes.
- private List<FragmentInstance> instances;
+ private final List<FragmentInstance> instances;
- private ExecutorService executor;
- private ExecutorService writeOperationExecutor;
- private ScheduledExecutorService scheduledExecutor;
-
- private IFragInstanceDispatcher dispatcher;
- private IFragInstanceStateTracker stateTracker;
- private IQueryTerminator queryTerminator;
+ private final IFragInstanceDispatcher dispatcher;
+ private final IFragInstanceStateTracker stateTracker;
+ private final IQueryTerminator queryTerminator;
public ClusterScheduler(
MPPQueryContext queryContext,
@@ -74,12 +69,9 @@ public class ClusterScheduler implements IScheduler {
ExecutorService writeOperationExecutor,
ScheduledExecutorService scheduledExecutor,
IClientManager<TEndPoint, SyncDataNodeInternalServiceClient>
internalServiceClientManager) {
- this.queryContext = queryContext;
this.stateMachine = stateMachine;
this.instances = instances;
this.queryType = queryType;
- this.executor = executor;
- this.scheduledExecutor = scheduledExecutor;
this.dispatcher =
new FragmentInstanceDispatcherImpl(
queryType, executor, writeOperationExecutor,
internalServiceClientManager);
diff --git a/server/src/main/java/org/apache/iotdb/db/service/DataNode.java
b/server/src/main/java/org/apache/iotdb/db/service/DataNode.java
index 95c9d4015c..d42c04d370 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/DataNode.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/DataNode.java
@@ -86,7 +86,7 @@ public class DataNode implements DataNodeMBean {
*/
private static final int DEFAULT_JOIN_RETRY = 10;
- private TEndPoint thisNode = new TEndPoint();
+ private final TEndPoint thisNode = new TEndPoint();
private DataNode() {
// we do not init anything here, so that we can re-initialize the instance
in IT.
@@ -95,8 +95,6 @@ public class DataNode implements DataNodeMBean {
private static final RegisterManager registerManager = new RegisterManager();
public static ServiceProvider serviceProvider;
- // private IClientManager clientManager;
-
public static DataNode getInstance() {
return DataNodeHolder.INSTANCE;
}
@@ -123,12 +121,12 @@ public class DataNode implements DataNodeMBean {
try {
// setup InternalService
setUpInternalService();
- // contact with config node to join into the cluster
- prepareJoinCluster();
+ // register current DataNode to ConfigNode
+ registerInConfigNode();
// setup DataNode
active();
// send message to config node stating that data node is ready
- joinCluster();
+ activateCurrentDataNode();
// setup rpc service
setUpRPCService();
logger.info("Congratulation, IoTDB DataNode is set up successfully. Now,
enjoy yourself!");
@@ -166,7 +164,7 @@ public class DataNode implements DataNodeMBean {
}
/** register DataNode with ConfigNode */
- private void prepareJoinCluster() throws StartupException {
+ private void registerInConfigNode() throws StartupException {
int retry = DEFAULT_JOIN_RETRY;
ConfigNodeInfo.getInstance()
@@ -320,7 +318,7 @@ public class DataNode implements DataNodeMBean {
}
/** send a message to ConfigNode after DataNode is available */
- private void joinCluster() throws StartupException {
+ private void activateCurrentDataNode() throws StartupException {
int retry = DEFAULT_JOIN_RETRY;
ConfigNodeInfo.getInstance()
@@ -444,8 +442,6 @@ public class DataNode implements DataNodeMBean {
Thread.setDefaultUncaughtExceptionHandler(new
IoTDBDefaultThreadExceptionHandler());
}
- private void dataNodeIdChecker() {}
-
private static class DataNodeHolder {
private static final DataNode INSTANCE = new DataNode();