This is an automated email from the ASF dual-hosted git repository.
tanxinyu 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 4e7687b [IOTDB-2730] start config service (#5287)
4e7687b is described below
commit 4e7687b54035e3ffdc057d6d335b169029d7e50e
Author: wangchao316 <[email protected]>
AuthorDate: Mon Mar 21 09:14:00 2022 +0800
[IOTDB-2730] start config service (#5287)
* [IOTDB-2730] alter iotdb-commons to commons
* [IOTDB-2730] config node server
* move startupcheck to commons
* alter test
* alter test
* alter test
* alter test
* rename commons to node-commons
* alter comment
* [IOTDB-2730] start config node service
* alter conf
* alter exception
* add Licensed
* fix issue
* alter thrift
---
.../org/apache/iotdb/cluster/ClusterIoTDB.java | 2 +-
.../cluster/ClusterIoTDBServerCommandLine.java | 2 +-
.../cluster/utils/nodetool/ClusterMonitor.java | 2 +-
.../resources/conf/iotdb-confignode.properties | 30 +++++-
.../iotdb/confignode/conf/ConfigNodeConf.java | 38 ++++++-
.../iotdb/confignode/conf/ConfigNodeConfCheck.java | 31 +++---
.../confignode/conf/ConfigNodeDescriptor.java | 117 +++++++++++++++------
.../iotdb/confignode/service/ConfigNode.java | 2 +-
.../confignode/service/ConfigNodeCommandLine.java | 5 +-
.../ConfigNodeMBean.java} | 15 +--
.../service/thrift/server/ConfigNodeRPCServer.java | 16 +--
.../thrift/server/ConfigNodeRPCServerMBean.java} | 15 +--
.../server/ConfigNodeRPCServerProcessor.java | 22 +++-
.../iotdb/db/integration/IoTDBCheckConfigIT.java | 2 +-
.../iotdb/commons/concurrent/ThreadName.java | 5 +-
.../commons}/exception/ConfigurationException.java | 3 +-
.../iotdb/commons/service/ThriftService.java | 2 +-
.../org/apache/iotdb/db/conf/IoTDBConfigCheck.java | 2 +-
.../java/org/apache/iotdb/db/service/IoTDB.java | 2 +-
.../org/apache/iotdb/db/service/IoTDBMBean.java | 4 +-
.../src/main/thrift/confignode.thrift | 29 ++++-
21 files changed, 239 insertions(+), 107 deletions(-)
diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
index 2206555..087eb8f 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDB.java
@@ -58,6 +58,7 @@ import org.apache.iotdb.cluster.utils.ClusterUtils;
import org.apache.iotdb.cluster.utils.nodetool.ClusterMonitor;
import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
import org.apache.iotdb.commons.conf.IoTDBConstant;
+import org.apache.iotdb.commons.exception.ConfigurationException;
import org.apache.iotdb.commons.exception.StartupException;
import org.apache.iotdb.commons.service.JMXService;
import org.apache.iotdb.commons.service.RegisterManager;
@@ -66,7 +67,6 @@ import org.apache.iotdb.commons.utils.TestOnly;
import org.apache.iotdb.db.conf.IoTDBConfig;
import org.apache.iotdb.db.conf.IoTDBConfigCheck;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
-import org.apache.iotdb.db.exception.ConfigurationException;
import org.apache.iotdb.db.exception.query.QueryProcessException;
import org.apache.iotdb.db.service.IoTDB;
import org.apache.iotdb.db.service.basic.ServiceProvider;
diff --git
a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDBServerCommandLine.java
b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDBServerCommandLine.java
index 2cc43b2..220fbb6 100644
---
a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDBServerCommandLine.java
+++
b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterIoTDBServerCommandLine.java
@@ -19,7 +19,7 @@
package org.apache.iotdb.cluster;
import org.apache.iotdb.commons.ServerCommandLine;
-import org.apache.iotdb.db.exception.ConfigurationException;
+import org.apache.iotdb.commons.exception.ConfigurationException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
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 086433d..758fcde 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
@@ -92,7 +92,7 @@ public class ClusterMonitor implements ClusterMonitorMBean,
IService {
private void startCollectClusterStatus() {
// monitor all nodes' live status
LOGGER.info("start metric node status and leader distribution");
-
IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(ThreadName.Cluster_Monitor.getName())
+
IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(ThreadName.CLUSTER_MONITOR.getName())
.scheduleAtFixedRate(
() -> {
MetaGroupMember metaGroupMember =
ClusterIoTDB.getInstance().getMetaGroupMember();
diff --git a/confignode/src/assembly/resources/conf/iotdb-confignode.properties
b/confignode/src/assembly/resources/conf/iotdb-confignode.properties
index f51a7d4..22ebc3f 100644
--- a/confignode/src/assembly/resources/conf/iotdb-confignode.properties
+++ b/confignode/src/assembly/resources/conf/iotdb-confignode.properties
@@ -73,4 +73,32 @@ config_node_address_lists=host0:22278,host1:22278,host2:22278
# thrift init buffer size
# Datatype: int
-# thrift_init_buffer_size=1024
\ No newline at end of file
+# thrift_init_buffer_size=1024
+
+####################
+### Directory Configuration
+####################
+
+# system dir
+# If this property is unset, system will save the data in the default relative
path directory under the IoTDB folder(i.e., %IOTDB_HOME%/data/system).
+# If it is absolute, system will save the data in exact location it points to.
+# If it is relative, system will save the data in the relative path directory
it indicates under the IoTDB folder.
+# For windows platform
+# If its prefix is a drive specifier followed by "\\", or if its prefix is
"\\\\", then the path is absolute. Otherwise, it is relative.
+# system_dir=data\\system
+# For Linux platform
+# If its prefix is "/", then the path is absolute. Otherwise, it is relative.
+# system_dir=data/system
+
+
+# data dirs
+# If this property is unset, system will save the data in the default relative
path directory under the IoTDB folder(i.e., %IOTDB_HOME%/data/data).
+# If it is absolute, system will save the data in exact location it points to.
+# If it is relative, system will save the data in the relative path directory
it indicates under the IoTDB folder.
+# Note: If data_dir is assigned an empty string(i.e.,zero-size), it will be
handled as a relative path.
+# For windows platform
+# If its prefix is a drive specifier followed by "\\", or if its prefix is
"\\\\", then the path is absolute. Otherwise, it is relative.
+# data_dirs=data\\data
+# For Linux platform
+# If its prefix is "/", then the path is absolute. Otherwise, it is relative.
+# data_dirs=data/data
\ No newline at end of file
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConf.java
b/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConf.java
index 5d17434..467fddf 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConf.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConf.java
@@ -18,18 +18,21 @@
*/
package org.apache.iotdb.confignode.conf;
+import org.apache.iotdb.commons.conf.IoTDBConstant;
import org.apache.iotdb.rpc.RpcUtils;
+import java.io.File;
+
public class ConfigNodeConf {
/** could set ip or hostname */
- private String rpcAddress;
+ private String rpcAddress = "0.0.0.0";
/** used for communication between data node and config node */
- private int rpcPort;
+ private int rpcPort = 22277;
/** used for communication between data node and data node */
- private Long internalPort;
+ private int internalPort = 22278;
/** every node should have the same config_node_address_lists */
private String addressLists;
@@ -58,6 +61,15 @@ public class ConfigNodeConf {
/** just for test wait for 60 second by default. */
private int thriftServerAwaitTimeForStopService = 60;
+ /** System directory, including version file for each storage group and
metadata */
+ private String systemDir =
+ ConfigNodeConstant.DATA_DIR + File.separator +
IoTDBConstant.SYSTEM_FOLDER_NAME;
+
+ /** Data directory of data. It can be settled as dataDirs = {"data1",
"data2", "data3"}; */
+ private String[] dataDirs = {
+ ConfigNodeConstant.DATA_DIR + File.separator + ConfigNodeConstant.DATA_DIR
+ };
+
public ConfigNodeConf() {
// empty constructor
}
@@ -134,11 +146,11 @@ public class ConfigNodeConf {
this.rpcPort = rpcPort;
}
- public Long getInternalPort() {
+ public int getInternalPort() {
return internalPort;
}
- public void setInternalPort(Long internalPort) {
+ public void setInternalPort(int internalPort) {
this.internalPort = internalPort;
}
@@ -157,4 +169,20 @@ public class ConfigNodeConf {
public void setThriftServerAwaitTimeForStopService(int
thriftServerAwaitTimeForStopService) {
this.thriftServerAwaitTimeForStopService =
thriftServerAwaitTimeForStopService;
}
+
+ public String getSystemDir() {
+ return systemDir;
+ }
+
+ public void setSystemDir(String systemDir) {
+ this.systemDir = systemDir;
+ }
+
+ public String[] getDataDirs() {
+ return dataDirs;
+ }
+
+ public void setDataDirs(String[] dataDirs) {
+ this.dataDirs = dataDirs;
+ }
}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConfCheck.java
b/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConfCheck.java
index ca43603..4c6ec59 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConfCheck.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/conf/ConfigNodeConfCheck.java
@@ -18,7 +18,8 @@
*/
package org.apache.iotdb.confignode.conf;
-import org.apache.iotdb.confignode.exception.conf.RepeatConfigurationException;
+import org.apache.iotdb.commons.exception.ConfigurationException;
+import org.apache.iotdb.commons.exception.StartupException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -44,17 +45,19 @@ public class ConfigNodeConfCheck {
private Properties specialProperties;
- public void checkConfig() throws RepeatConfigurationException, IOException {
+ private ConfigNodeConfCheck() {
+ specialProperties = new Properties();
+ }
- String propsDir = ConfigNodeDescriptor.getInstance().getPropsDir();
- if (propsDir == null) {
- // Skip configuration check when developer mode or test mode
- return;
+ public void checkConfig() throws ConfigurationException, IOException,
StartupException {
+ // if systemDir does not exist, create systemDir
+ File systemDir = new File(conf.getSystemDir());
+ if (systemDir.isDirectory() && !systemDir.exists()) {
+ systemDir.mkdirs();
}
- specialProperties = new Properties();
File specialPropertiesFile =
- new File(propsDir + File.separator +
ConfigNodeConstant.SPECIAL_CONF_NAME);
+ new File(conf.getSystemDir() + File.separator +
ConfigNodeConstant.SPECIAL_CONF_NAME);
if (!specialPropertiesFile.exists()) {
if (specialPropertiesFile.createNewFile()) {
LOGGER.info(
@@ -66,7 +69,7 @@ public class ConfigNodeConfCheck {
LOGGER.error(
"Can't create special configuration file {} for ConfigNode.
IoTDB-ConfigNode is shutdown.",
specialPropertiesFile.getAbsolutePath());
- System.exit(-1);
+ throw new StartupException("Can't create special configuration file");
}
}
@@ -93,13 +96,13 @@ public class ConfigNodeConfCheck {
}
/** Ensure that special parameters are consistent with each startup except
the first one */
- private void checkSpecialProperties() throws RepeatConfigurationException {
+ private void checkSpecialProperties() throws ConfigurationException {
int specialDeviceGroupCount =
Integer.parseInt(
specialProperties.getProperty(
"device_group_count",
String.valueOf(conf.getDeviceGroupCount())));
if (specialDeviceGroupCount != conf.getDeviceGroupCount()) {
- throw new RepeatConfigurationException(
+ throw new ConfigurationException(
"device_group_count",
String.valueOf(conf.getDeviceGroupCount()),
String.valueOf(specialDeviceGroupCount));
@@ -110,7 +113,7 @@ public class ConfigNodeConfCheck {
"device_group_hash_executor_class",
conf.getDeviceGroupHashExecutorClass());
if (!Objects.equals(
specialDeviceGroupHashExecutorClass,
conf.getDeviceGroupHashExecutorClass())) {
- throw new RepeatConfigurationException(
+ throw new ConfigurationException(
"device_group_hash_executor_class",
conf.getDeviceGroupHashExecutorClass(),
specialDeviceGroupHashExecutorClass);
@@ -129,8 +132,4 @@ public class ConfigNodeConfCheck {
public static ConfigNodeConfCheck getInstance() {
return ConfigNodeConfCheckHolder.INSTANCE;
}
-
- private ConfigNodeConfCheck() {
- LOGGER.info("Starting IoTDB Cluster ConfigNode " +
ConfigNodeConstant.VERSION);
- }
}
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 a5e4a8a..373d83a 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
@@ -29,7 +29,6 @@ import java.net.URL;
import java.util.Properties;
public class ConfigNodeDescriptor {
-
private static final Logger LOGGER =
LoggerFactory.getLogger(ConfigNodeDescriptor.class);
private final ConfigNodeConf conf = new ConfigNodeConf();
@@ -42,42 +41,52 @@ public class ConfigNodeDescriptor {
return conf;
}
- public String getPropsDir() {
- // Check if CONFIG_NODE_CONF is set
- String propsDir = System.getProperty(ConfigNodeConstant.CONFIGNODE_CONF,
null);
- if (propsDir == null) {
- // Check if CONFIG_NODE_HOME is set
- propsDir = System.getProperty(ConfigNodeConstant.CONFIGNODE_HOME, null);
- if (propsDir == null) {
- // When start ConfigNode with script, CONFIG_NODE_CONF and
CONFIG_NODE_HOME must be set.
- // Therefore, this case is TestOnly
- // TODO: Specify a test dir
- }
- propsDir = propsDir + File.separator + ConfigNodeConstant.CONF_DIR;
- }
-
- return propsDir;
- }
-
+ /**
+ * get props url location
+ *
+ * @return url object if location exit, otherwise null.
+ */
public URL getPropsUrl() {
- String url = getPropsDir();
-
- if (url == null) {
- return null;
+ // Check if a config-directory was specified first.
+ String urlString = System.getProperty(ConfigNodeConstant.CONFIGNODE_CONF,
null);
+ // If it wasn't, check if a home directory was provided (This usually
contains a config)
+ if (urlString == null) {
+ urlString = System.getProperty(ConfigNodeConstant.CONFIGNODE_HOME, null);
+ if (urlString != null) {
+ urlString =
+ urlString
+ + File.separatorChar
+ + "conf"
+ + File.separatorChar
+ + ConfigNodeConstant.CONF_NAME;
+ } else {
+ // If this too wasn't provided, try to find a default config in the
root of the classpath.
+ URL uri = ConfigNodeConf.class.getResource("/" +
ConfigNodeConstant.CONF_NAME);
+ if (uri != null) {
+ return uri;
+ }
+ LOGGER.warn(
+ "Cannot find IOTDB_HOME or IOTDB_CONF environment variable when
loading "
+ + "config file {}, use default configuration",
+ ConfigNodeConstant.CONF_NAME);
+ // update all data seriesPath
+ // conf.updatePath();
+ return null;
+ }
}
-
- // Add props prefix
- if (!url.startsWith("file:") && !url.startsWith("classpath:")) {
- url = "file:" + url;
+ // If a config location was provided, but it doesn't end with a properties
file,
+ // append the default location.
+ else if (!urlString.endsWith(".properties")) {
+ urlString += (File.separatorChar + ConfigNodeConstant.CONF_NAME);
}
- // Add props suffix
- if (!url.endsWith(".properties")) {
- url += File.separator + ConfigNodeConstant.CONF_NAME;
+ // If the url doesn't start with "file:" or "classpath:", it's provided as
a no path.
+ // So we need to add it to make it a real URL.
+ if (!urlString.startsWith("file:") && !urlString.startsWith("classpath:"))
{
+ urlString = "file:" + urlString;
}
-
try {
- return new URL(url);
+ return new URL(urlString);
} catch (MalformedURLException e) {
return null;
}
@@ -107,6 +116,52 @@ public class ConfigNodeDescriptor {
properties.getProperty(
"device_group_hash_executor_class",
conf.getDeviceGroupHashExecutorClass()));
+ conf.setRpcAddress(properties.getProperty("config_node_rpc_address",
conf.getRpcAddress()));
+
+ conf.setRpcPort(
+ Integer.parseInt(
+ properties.getProperty("config_node_rpc_port",
String.valueOf(conf.getRpcPort()))));
+
+ conf.setInternalPort(
+ Integer.parseInt(
+ properties.getProperty(
+ "config_node_internal_port",
String.valueOf(conf.getInternalPort()))));
+
+ conf.setAddressLists(
+ properties.getProperty("config_node_address_lists",
conf.getAddressLists()));
+
+ conf.setRpcAdvancedCompressionEnable(
+ Boolean.parseBoolean(
+ properties.getProperty(
+ "rpc_advanced_compression_enable",
+ String.valueOf(conf.isRpcAdvancedCompressionEnable()))));
+
+ conf.setRpcThriftCompressionEnabled(
+ Boolean.parseBoolean(
+ properties.getProperty(
+ "rpc_thrift_compression_enable",
+ String.valueOf(conf.isRpcThriftCompressionEnabled()))));
+
+ conf.setRpcMaxConcurrentClientNum(
+ Integer.parseInt(
+ properties.getProperty(
+ "rpc_max_concurrent_client_num",
+ String.valueOf(conf.getRpcMaxConcurrentClientNum()))));
+
+ conf.setThriftDefaultBufferSize(
+ Integer.parseInt(
+ properties.getProperty(
+ "thrift_init_buffer_size",
String.valueOf(conf.getThriftDefaultBufferSize()))));
+
+ conf.setThriftMaxFrameSize(
+ Integer.parseInt(
+ properties.getProperty(
+ "thrift_max_frame_size",
String.valueOf(conf.getThriftMaxFrameSize()))));
+
+ conf.setSystemDir(properties.getProperty("system_dir",
conf.getSystemDir()));
+
+ conf.setDataDirs(properties.getProperty("data_dirs",
conf.getDataDirs()[0]).split(","));
+
} catch (IOException e) {
LOGGER.warn("Couldn't load ConfigNode conf file, use default config", e);
}
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 77dffdc..6069f70 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
@@ -31,7 +31,7 @@ import
org.apache.iotdb.confignode.service.thrift.server.ConfigNodeRPCServerProc
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-public class ConfigNode {
+public class ConfigNode implements ConfigNodeMBean {
private static final Logger LOGGER =
LoggerFactory.getLogger(ConfigNode.class);
private final String mbeanName =
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNodeCommandLine.java
b/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNodeCommandLine.java
index 01f8a84..63bb8f6 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNodeCommandLine.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNodeCommandLine.java
@@ -19,8 +19,9 @@
package org.apache.iotdb.confignode.service;
import org.apache.iotdb.commons.ServerCommandLine;
+import org.apache.iotdb.commons.exception.ConfigurationException;
+import org.apache.iotdb.commons.exception.StartupException;
import org.apache.iotdb.confignode.conf.ConfigNodeConfCheck;
-import org.apache.iotdb.confignode.exception.ConfigNodeException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -63,7 +64,7 @@ public class ConfigNodeCommandLine extends ServerCommandLine {
try {
// Check parameters
ConfigNodeConfCheck.getInstance().checkConfig();
- } catch (ConfigNodeException | IOException e) {
+ } catch (IOException | ConfigurationException | StartupException e) {
LOGGER.error("Meet error when doing start checking", e);
return -1;
}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/exception/conf/RepeatConfigurationException.java
b/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNodeMBean.java
similarity index 60%
rename from
confignode/src/main/java/org/apache/iotdb/confignode/exception/conf/RepeatConfigurationException.java
rename to
confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNodeMBean.java
index 515f098..cf88956 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/exception/conf/RepeatConfigurationException.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNodeMBean.java
@@ -16,17 +16,6 @@
* specific language governing permissions and limitations
* under the License.
*/
-package org.apache.iotdb.confignode.exception.conf;
+package org.apache.iotdb.confignode.service;
-import org.apache.iotdb.confignode.exception.ConfigNodeException;
-
-/** Throws when there exists some special parameters are repeatedly defined */
-public class RepeatConfigurationException extends ConfigNodeException {
-
- public RepeatConfigurationException(String parameter, String badValue,
String correctValue) {
- super(
- String.format(
- "Parameter %s can not be %s, because you're already set to: %s.",
- parameter, badValue, correctValue));
- }
-}
+public interface ConfigNodeMBean {}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/server/ConfigNodeRPCServer.java
b/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/server/ConfigNodeRPCServer.java
index 19180bc..e1001d7 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/server/ConfigNodeRPCServer.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/server/ConfigNodeRPCServer.java
@@ -19,6 +19,7 @@
package org.apache.iotdb.confignode.service.thrift.server;
import org.apache.iotdb.commons.concurrent.ThreadName;
+import org.apache.iotdb.commons.conf.IoTDBConstant;
import org.apache.iotdb.commons.exception.runtime.RPCServiceException;
import org.apache.iotdb.commons.service.ServiceType;
import org.apache.iotdb.commons.service.ThriftService;
@@ -28,8 +29,7 @@ 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 ConfigNodeRPCServer extends ThriftService {
-
+public class ConfigNodeRPCServer extends ThriftService implements
ConfigNodeRPCServerMBean {
private ConfigNodeConf config = ConfigNodeDescriptor.getInstance().getConf();
private ConfigNodeRPCServerProcessor configNodeRPCServerProcessor;
@@ -49,15 +49,15 @@ public class ConfigNodeRPCServer extends ThriftService {
@Override
public void initSyncedServiceImpl(Object configNodeRPCServerProcessor) {
this.configNodeRPCServerProcessor = (ConfigNodeRPCServerProcessor)
configNodeRPCServerProcessor;
+
+ super.mbeanName =
+ String.format(
+ "%s:%s=%s", this.getClass().getPackage(), IoTDBConstant.JMX_TYPE,
getID().getJmxName());
super.initSyncedServiceImpl(this.configNodeRPCServerProcessor);
}
@Override
public void initTProcessor() throws InstantiationException {
- if (processor == null) {
- throw new InstantiationException("ConfigNodeRPCServerProcessor is null");
- }
-
processor = new ConfigIService.Processor<>(configNodeRPCServerProcessor);
}
@@ -69,7 +69,7 @@ public class ConfigNodeRPCServer extends ThriftService {
new ThriftServiceThread(
processor,
getID().getName(),
- ThreadName.CLUSTER_RPC_CLIENT.getName(),
+ ThreadName.CONFIG_NODE_RPC_CLIENT.getName(),
getBindIP(),
getBindPort(),
config.getRpcMaxConcurrentClientNum(),
@@ -79,7 +79,7 @@ public class ConfigNodeRPCServer extends ThriftService {
} catch (RPCServiceException e) {
throw new IllegalAccessException(e.getMessage());
}
- thriftServiceThread.setName(ThreadName.CLUSTER_RPC_SERVICE.getName());
+ thriftServiceThread.setName(ThreadName.CONFIG_NODE_RPC_SERVER.getName());
}
@Override
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/exception/startup/StartupException.java
b/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/server/ConfigNodeRPCServerMBean.java
similarity index 64%
rename from
confignode/src/main/java/org/apache/iotdb/confignode/exception/startup/StartupException.java
rename to
confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/server/ConfigNodeRPCServerMBean.java
index 242c8c5..5b97780 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/exception/startup/StartupException.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/server/ConfigNodeRPCServerMBean.java
@@ -16,18 +16,7 @@
* specific language governing permissions and limitations
* under the License.
*/
-package org.apache.iotdb.confignode.exception.startup;
-import org.apache.iotdb.confignode.exception.ConfigNodeException;
+package org.apache.iotdb.confignode.service.thrift.server;
-/** Throws when there exists errors when startup checks */
-public class StartupException extends ConfigNodeException {
-
- public StartupException(String name, String message) {
- super(String.format("Failed to start [%s], because [%s]", name, message));
- }
-
- public StartupException(String message) {
- super(message);
- }
-}
+public interface ConfigNodeRPCServerMBean {}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/server/ConfigNodeRPCServerProcessor.java
b/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/server/ConfigNodeRPCServerProcessor.java
index 0589dc9..4aabb73 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/server/ConfigNodeRPCServerProcessor.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/server/ConfigNodeRPCServerProcessor.java
@@ -20,12 +20,15 @@ package org.apache.iotdb.confignode.service.thrift.server;
import org.apache.iotdb.confignode.manager.ConfigManager;
import org.apache.iotdb.confignode.rpc.thrift.ConfigIService;
+import org.apache.iotdb.confignode.rpc.thrift.DataNodeRegisterReq;
+import org.apache.iotdb.confignode.rpc.thrift.DataNodesInfo;
import org.apache.iotdb.confignode.rpc.thrift.DataPartitionInfo;
import org.apache.iotdb.confignode.rpc.thrift.DeleteStorageGroupReq;
+import org.apache.iotdb.confignode.rpc.thrift.DeviceGroupHashInfo;
import org.apache.iotdb.confignode.rpc.thrift.GetDataPartitionReq;
import org.apache.iotdb.confignode.rpc.thrift.GetSchemaPartitionReq;
import org.apache.iotdb.confignode.rpc.thrift.SchemaPartitionInfo;
-import org.apache.iotdb.confignode.rpc.thrift.SetStorageGroupoReq;
+import org.apache.iotdb.confignode.rpc.thrift.SetStorageGroupReq;
import org.apache.iotdb.service.rpc.thrift.TSStatus;
import org.apache.thrift.TException;
@@ -40,7 +43,7 @@ public class ConfigNodeRPCServerProcessor implements
ConfigIService.Iface {
}
@Override
- public TSStatus setStorageGroup(SetStorageGroupoReq req) throws TException {
+ public TSStatus setStorageGroup(SetStorageGroupReq req) throws TException {
return null;
}
@@ -50,6 +53,11 @@ public class ConfigNodeRPCServerProcessor implements
ConfigIService.Iface {
}
@Override
+ public DeviceGroupHashInfo getDeviceGroupHashInfo() throws TException {
+ return null;
+ }
+
+ @Override
public SchemaPartitionInfo getSchemaPartition(GetSchemaPartitionReq req)
throws TException {
return null;
}
@@ -59,6 +67,16 @@ public class ConfigNodeRPCServerProcessor implements
ConfigIService.Iface {
return null;
}
+ @Override
+ public TSStatus registerDataNode(DataNodeRegisterReq req) throws TException {
+ return null;
+ }
+
+ @Override
+ public DataNodesInfo getDataNodesInfo() throws TException {
+ return null;
+ }
+
public void handleClientExit() {}
// TODO: Interfaces for data operations
diff --git
a/integration/src/test/java/org/apache/iotdb/db/integration/IoTDBCheckConfigIT.java
b/integration/src/test/java/org/apache/iotdb/db/integration/IoTDBCheckConfigIT.java
index da985ce..d909bf3 100644
---
a/integration/src/test/java/org/apache/iotdb/db/integration/IoTDBCheckConfigIT.java
+++
b/integration/src/test/java/org/apache/iotdb/db/integration/IoTDBCheckConfigIT.java
@@ -18,10 +18,10 @@
*/
package org.apache.iotdb.db.integration;
+import org.apache.iotdb.commons.exception.ConfigurationException;
import org.apache.iotdb.db.conf.IoTDBConfigCheck;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory;
-import org.apache.iotdb.db.exception.ConfigurationException;
import org.apache.iotdb.db.service.IoTDB;
import org.apache.iotdb.db.utils.EnvironmentUtils;
import org.apache.iotdb.itbase.category.LocalStandaloneTest;
diff --git
a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
index c12faea..ba8d1ae 100644
---
a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
+++
b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
@@ -69,8 +69,9 @@ public enum ThreadName {
CLUSTER_DATA_RPC_CLIENT("ClusterDataRPC-Client"),
CLUSTER_DATA_HEARTBEAT_RPC_SERVICE("ClusterDataHeartbeatRPC"),
CLUSTER_DATA_HEARTBEAT_RPC_CLIENT("ClusterDataHeartbeatRPC-Client"),
- Cluster_Monitor("ClusterMonitor"),
- ;
+ CLUSTER_MONITOR("ClusterMonitor"),
+ CONFIG_NODE_RPC_SERVER("ConfigNodeRpcServer"),
+ CONFIG_NODE_RPC_CLIENT("ConfigNodeRPC-Client");
private final String name;
diff --git
a/server/src/main/java/org/apache/iotdb/db/exception/ConfigurationException.java
b/node-commons/src/main/java/org/apache/iotdb/commons/exception/ConfigurationException.java
similarity index 93%
rename from
server/src/main/java/org/apache/iotdb/db/exception/ConfigurationException.java
rename to
node-commons/src/main/java/org/apache/iotdb/commons/exception/ConfigurationException.java
index f615bda..1c7d36e 100644
---
a/server/src/main/java/org/apache/iotdb/db/exception/ConfigurationException.java
+++
b/node-commons/src/main/java/org/apache/iotdb/commons/exception/ConfigurationException.java
@@ -17,9 +17,8 @@
* under the License.
*/
-package org.apache.iotdb.db.exception;
+package org.apache.iotdb.commons.exception;
-import org.apache.iotdb.commons.exception.IoTDBException;
import org.apache.iotdb.rpc.TSStatusCode;
public class ConfigurationException extends IoTDBException {
diff --git
a/node-commons/src/main/java/org/apache/iotdb/commons/service/ThriftService.java
b/node-commons/src/main/java/org/apache/iotdb/commons/service/ThriftService.java
index bc93f6d..e792d88 100644
---
a/node-commons/src/main/java/org/apache/iotdb/commons/service/ThriftService.java
+++
b/node-commons/src/main/java/org/apache/iotdb/commons/service/ThriftService.java
@@ -35,7 +35,7 @@ public abstract class ThriftService implements IService {
public static final String STATUS_UP = "UP";
public static final String STATUS_DOWN = "DOWN";
- protected final String mbeanName =
+ protected String mbeanName =
String.format(
"%s:%s=%s", IoTDBConstant.IOTDB_PACKAGE, IoTDBConstant.JMX_TYPE,
getID().getJmxName());
protected AbstractThriftServiceThread thriftServiceThread;
diff --git
a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfigCheck.java
b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfigCheck.java
index 5a48019..c285d82 100644
--- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfigCheck.java
+++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfigCheck.java
@@ -19,8 +19,8 @@
package org.apache.iotdb.db.conf;
import org.apache.iotdb.commons.conf.IoTDBConstant;
+import org.apache.iotdb.commons.exception.ConfigurationException;
import org.apache.iotdb.db.engine.fileSystem.SystemFileFactory;
-import org.apache.iotdb.db.exception.ConfigurationException;
import org.apache.iotdb.db.metadata.upgrade.MetadataUpgrader;
import org.apache.iotdb.tsfile.common.conf.TSFileConfig;
import org.apache.iotdb.tsfile.common.conf.TSFileDescriptor;
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 53f6361..9b19894 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
@@ -20,6 +20,7 @@ package org.apache.iotdb.db.service;
import org.apache.iotdb.commons.concurrent.IoTDBDefaultThreadExceptionHandler;
import org.apache.iotdb.commons.conf.IoTDBConstant;
+import org.apache.iotdb.commons.exception.ConfigurationException;
import org.apache.iotdb.commons.exception.StartupException;
import org.apache.iotdb.commons.service.JMXService;
import org.apache.iotdb.commons.service.RegisterManager;
@@ -34,7 +35,6 @@ import
org.apache.iotdb.db.engine.compaction.CompactionTaskManager;
import org.apache.iotdb.db.engine.cq.ContinuousQueryService;
import org.apache.iotdb.db.engine.flush.FlushManager;
import org.apache.iotdb.db.engine.trigger.service.TriggerRegistrationService;
-import org.apache.iotdb.db.exception.ConfigurationException;
import org.apache.iotdb.db.exception.query.QueryProcessException;
import org.apache.iotdb.db.metadata.SchemaEngine;
import org.apache.iotdb.db.protocol.influxdb.meta.InfluxDBMetaManager;
diff --git a/server/src/main/java/org/apache/iotdb/db/service/IoTDBMBean.java
b/server/src/main/java/org/apache/iotdb/db/service/IoTDBMBean.java
index 968bc08..aa178d0 100644
--- a/server/src/main/java/org/apache/iotdb/db/service/IoTDBMBean.java
+++ b/server/src/main/java/org/apache/iotdb/db/service/IoTDBMBean.java
@@ -18,10 +18,10 @@
*/
package org.apache.iotdb.db.service;
-import org.apache.iotdb.db.exception.StorageEngineException;
+import org.apache.iotdb.commons.exception.ShutdownException;
@FunctionalInterface
public interface IoTDBMBean {
- void stop() throws StorageEngineException;
+ void stop() throws ShutdownException;
}
diff --git a/thrift-confignode/src/main/thrift/confignode.thrift
b/thrift-confignode/src/main/thrift/confignode.thrift
index 89a8f47..6c28535 100644
--- a/thrift-confignode/src/main/thrift/confignode.thrift
+++ b/thrift-confignode/src/main/thrift/confignode.thrift
@@ -21,7 +21,7 @@ include "rpc.thrift"
namespace java org.apache.iotdb.confignode.rpc.thrift
namespace py iotdb.thrift.confignode
-struct SetStorageGroupoReq {
+struct SetStorageGroupReq {
1: required string storageGroup
}
@@ -53,12 +53,37 @@ struct DataPartitionInfo {
2: required map<i32, list<i32>> dataRegionIDsMap
}
+struct DeviceGroupHashInfo {
+ 1: required i32 deviceGroupCount
+ 2: required string hashClass
+}
+
+struct DataNodeRegisterReq {
+ 1: required rpc.EndPoint endPoint
+}
+
+struct DataNodeRegisterResp {
+ 1: required rpc.TSStatus registerResult
+ 2: optional i32 dataNodeID
+}
+
+struct DataNodesInfo {
+ 1: required map<i32, list<rpc.EndPoint>> dataNodesMap
+}
+
+
service ConfigIService {
- rpc.TSStatus setStorageGroup(SetStorageGroupoReq req)
+ rpc.TSStatus setStorageGroup(SetStorageGroupReq req)
rpc.TSStatus deleteStorageGroup(DeleteStorageGroupReq req)
+ DeviceGroupHashInfo getDeviceGroupHashInfo()
+
SchemaPartitionInfo getSchemaPartition(GetSchemaPartitionReq req)
DataPartitionInfo getDataPartition(GetDataPartitionReq req)
+
+ rpc.TSStatus registerDataNode(DataNodeRegisterReq req)
+
+ DataNodesInfo getDataNodesInfo()
}
\ No newline at end of file