This is an automated email from the ASF dual-hosted git repository.
JackieTien97 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 a185d9c7e94 Extend the monitoring framework to support multiple node
types (#18470)
a185d9c7e94 is described below
commit a185d9c7e943fb1b2699e980eb7ab0ad35feece2
Author: suchenglong <[email protected]>
AuthorDate: Tue Sep 1 09:18:40 2026 +0800
Extend the monitoring framework to support multiple node types (#18470)
---
.../iotdb/confignode/service/ConfigNode.java | 2 +-
.../apache/iotdb/db/i18n/DataNodeMiscMessages.java | 2 -
.../apache/iotdb/db/i18n/DataNodeMiscMessages.java | 2 -
.../db/service/metrics/DataNodeMetricsHelper.java | 1 +
.../metrics/config/MetricConfigDescriptor.java | 73 ++++++++++++----------
.../org/apache/iotdb/metrics/utils/NodeType.java | 3 +-
iotdb-core/node-commons/pom.xml | 8 +++
.../apache/iotdb/commons/i18n/ServiceMessages.java | 4 ++
.../apache/iotdb/commons/i18n/ServiceMessages.java | 4 ++
.../commons/service/metric}/ProcessMetrics.java | 8 +--
10 files changed, 65 insertions(+), 42 deletions(-)
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNode.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNode.java
index ffd5ec2d33b..cf49a3916bf 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNode.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/service/ConfigNode.java
@@ -39,6 +39,7 @@ import org.apache.iotdb.commons.service.RegisterManager;
import org.apache.iotdb.commons.service.ServiceType;
import org.apache.iotdb.commons.service.metric.JvmGcMonitorMetrics;
import org.apache.iotdb.commons.service.metric.MetricService;
+import org.apache.iotdb.commons.service.metric.ProcessMetrics;
import org.apache.iotdb.commons.service.metric.cpu.CpuUsageMetrics;
import org.apache.iotdb.commons.utils.StatusUtils;
import org.apache.iotdb.commons.utils.TestOnly;
@@ -59,7 +60,6 @@ import
org.apache.iotdb.confignode.rpc.thrift.TConfigNodeRegisterResp;
import org.apache.iotdb.confignode.rpc.thrift.TNodeVersionInfo;
import org.apache.iotdb.confignode.service.thrift.ConfigNodeRPCService;
import
org.apache.iotdb.confignode.service.thrift.ConfigNodeRPCServiceProcessor;
-import org.apache.iotdb.db.service.metrics.ProcessMetrics;
import org.apache.iotdb.metrics.config.MetricConfigDescriptor;
import org.apache.iotdb.metrics.metricsets.UpTimeMetrics;
import org.apache.iotdb.metrics.metricsets.disk.DiskMetrics;
diff --git
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
index e76211672ff..4aef9cd74f0 100644
---
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
+++
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
@@ -446,8 +446,6 @@ public final class DataNodeMiscMessages {
//
---------------------------------------------------------------------------
// service – metrics
//
---------------------------------------------------------------------------
- public static final String FAILED_GET_PROCESS_RESIDENT_MEMORY =
- "Failed to get process resident memory for pid {}";
public static final String DATANODE_PORT_CHECK_SUCCESSFUL = "DataNode port
check successful.";
//
---------------------------------------------------------------------------
diff --git
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
index c9a3f49710f..bf909f27745 100644
---
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
+++
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java
@@ -445,8 +445,6 @@ public final class DataNodeMiscMessages {
//
---------------------------------------------------------------------------
// service – metrics
//
---------------------------------------------------------------------------
- public static final String FAILED_GET_PROCESS_RESIDENT_MEMORY =
- "获取进程 {} 的常驻内存失败";
public static final String DATANODE_PORT_CHECK_SUCCESSFUL = "DataNode
端口检查通过。";
//
---------------------------------------------------------------------------
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/DataNodeMetricsHelper.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/DataNodeMetricsHelper.java
index 9f6f4e57201..d857f546e2d 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/DataNodeMetricsHelper.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/DataNodeMetricsHelper.java
@@ -29,6 +29,7 @@ import org.apache.iotdb.commons.conf.IoTDBConstant;
import org.apache.iotdb.commons.service.metric.JvmGcMonitorMetrics;
import org.apache.iotdb.commons.service.metric.MetricService;
import org.apache.iotdb.commons.service.metric.PerformanceOverviewMetrics;
+import org.apache.iotdb.commons.service.metric.ProcessMetrics;
import org.apache.iotdb.commons.service.metric.cpu.CpuUsageMetrics;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.pipe.metric.PipeDataNodeMetrics;
diff --git
a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfigDescriptor.java
b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfigDescriptor.java
index def3d1e50c0..0d83d387409 100644
---
a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfigDescriptor.java
+++
b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfigDescriptor.java
@@ -33,13 +33,24 @@ public class MetricConfigDescriptor {
/** The metric config of metric service. */
private static final MetricConfig metricConfig = new MetricConfig();
+ private static final String CONFIG_NODE_PREFIX = "cn_";
+ private static final String DATA_NODE_PREFIX = "dn_";
+
private MetricConfigDescriptor() {
// empty constructor
}
/** Load properties into metric config. */
public void loadProps(Properties properties, boolean isConfigNode) {
- MetricConfig loadConfig = generateFromProperties(properties, isConfigNode);
+ loadProps(properties, isConfigNode ? CONFIG_NODE_PREFIX :
DATA_NODE_PREFIX);
+ }
+
+ /**
+ * Load properties into metric config with a node-specific prefix (e.g.
{@code "cn_"}, {@code
+ * "dn_"}, {@code "sn_"}).
+ */
+ public void loadProps(Properties properties, String prefix) {
+ MetricConfig loadConfig = generateFromProperties(properties, prefix);
metricConfig.copy(loadConfig);
}
@@ -49,7 +60,16 @@ public class MetricConfigDescriptor {
* @return reload level of metric service
*/
public ReloadLevel loadHotProps(Properties properties, boolean isConfigNode)
{
- MetricConfig newMetricConfig = generateFromProperties(properties,
isConfigNode);
+ return loadHotProps(properties, isConfigNode ? CONFIG_NODE_PREFIX :
DATA_NODE_PREFIX);
+ }
+
+ /**
+ * Load properties into metric config when reload service with a
node-specific prefix.
+ *
+ * @return reload level of metric service
+ */
+ public ReloadLevel loadHotProps(Properties properties, String prefix) {
+ MetricConfig newMetricConfig = generateFromProperties(properties, prefix);
ReloadLevel reloadLevel = ReloadLevel.NOTHING;
if (!metricConfig.equals(newMetricConfig)) {
if
(!metricConfig.getMetricLevel().equals(newMetricConfig.getMetricLevel())
@@ -73,7 +93,7 @@ public class MetricConfigDescriptor {
}
/** Load properties into metric config. */
- private MetricConfig generateFromProperties(Properties properties, boolean
isConfigNode) {
+ private MetricConfig generateFromProperties(Properties properties, String
prefix) {
MetricConfig loadConfig = new MetricConfig();
String reporterList =
@@ -85,16 +105,13 @@ public class MetricConfigDescriptor {
.map(ReporterType::toString)
.collect(Collectors.toSet())),
properties,
- isConfigNode);
+ prefix);
loadConfig.setMetricReporterList(reporterList);
loadConfig.setMetricLevel(
MetricLevel.valueOf(
getProperty(
- "metric_level",
- String.valueOf(loadConfig.getMetricLevel()),
- properties,
- isConfigNode)));
+ "metric_level", String.valueOf(loadConfig.getMetricLevel()),
properties, prefix)));
loadConfig.setAsyncCollectPeriodInSecond(
Integer.parseInt(
@@ -102,7 +119,7 @@ public class MetricConfigDescriptor {
"metric_async_collect_period",
String.valueOf(loadConfig.getAsyncCollectPeriodInSecond()),
properties,
- isConfigNode)));
+ prefix)));
loadConfig.setPrometheusReporterPort(
Integer.parseInt(
@@ -110,7 +127,7 @@ public class MetricConfigDescriptor {
"metric_prometheus_reporter_port",
String.valueOf(loadConfig.getPrometheusReporterPort()),
properties,
- isConfigNode)));
+ prefix)));
loadConfig.setPrometheusReporterUsername(
getPropertyWithoutPrefix(
@@ -139,8 +156,7 @@ public class MetricConfigDescriptor {
IoTDBReporterConfig reporterConfig = loadConfig.getIoTDBReporterConfig();
reporterConfig.setHost(
- getProperty(
- "metric_iotdb_reporter_host", reporterConfig.getHost(),
properties, isConfigNode));
+ getProperty("metric_iotdb_reporter_host", reporterConfig.getHost(),
properties, prefix));
reporterConfig.setPort(
Integer.valueOf(
@@ -148,21 +164,15 @@ public class MetricConfigDescriptor {
"metric_iotdb_reporter_port",
String.valueOf(reporterConfig.getPort()),
properties,
- isConfigNode)));
+ prefix)));
reporterConfig.setUsername(
getProperty(
- "metric_iotdb_reporter_username",
- reporterConfig.getUsername(),
- properties,
- isConfigNode));
+ "metric_iotdb_reporter_username", reporterConfig.getUsername(),
properties, prefix));
reporterConfig.setPassword(
getProperty(
- "metric_iotdb_reporter_password",
- reporterConfig.getPassword(),
- properties,
- isConfigNode));
+ "metric_iotdb_reporter_password", reporterConfig.getPassword(),
properties, prefix));
reporterConfig.setMaxConnectionNumber(
Integer.valueOf(
@@ -170,14 +180,11 @@ public class MetricConfigDescriptor {
"metric_iotdb_reporter_max_connection_number",
String.valueOf(reporterConfig.getMaxConnectionNumber()),
properties,
- isConfigNode)));
+ prefix)));
reporterConfig.setLocation(
getProperty(
- "metric_iotdb_reporter_location",
- reporterConfig.getLocation(),
- properties,
- isConfigNode));
+ "metric_iotdb_reporter_location", reporterConfig.getLocation(),
properties, prefix));
reporterConfig.setPushPeriodInSecond(
Integer.valueOf(
@@ -185,8 +192,9 @@ public class MetricConfigDescriptor {
"metric_iotdb_reporter_push_period",
String.valueOf(reporterConfig.getPushPeriodInSecond()),
properties,
- isConfigNode)));
- if (!isConfigNode) {
+ prefix)));
+
+ if (DATA_NODE_PREFIX.equals(prefix)) {
loadConfig.setInternalReportType(
InternalReporterType.valueOf(
properties.getProperty(
@@ -197,11 +205,12 @@ public class MetricConfigDescriptor {
return loadConfig;
}
- /** Get property from confignode or datanode. */
+ /**
+ * Get property with a node-specific prefix (e.g. {@code "cn_"}, {@code
"dn_"}, {@code "sn_"}).
+ */
private String getProperty(
- String target, String defaultValue, Properties properties, boolean
isConfigNode) {
- return Optional.ofNullable(
- properties.getProperty((isConfigNode ? "cn_" : "dn_") + target,
defaultValue))
+ String target, String defaultValue, Properties properties, String
prefix) {
+ return Optional.ofNullable(properties.getProperty(prefix + target,
defaultValue))
.map(String::trim)
.orElse(defaultValue);
}
diff --git
a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/utils/NodeType.java
b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/utils/NodeType.java
index 1ffa95030c3..e2c826780f7 100644
---
a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/utils/NodeType.java
+++
b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/utils/NodeType.java
@@ -21,7 +21,8 @@ package org.apache.iotdb.metrics.utils;
public enum NodeType {
CONFIGNODE,
- DATANODE;
+ DATANODE,
+ STREAMNODE;
@Override
public String toString() {
diff --git a/iotdb-core/node-commons/pom.xml b/iotdb-core/node-commons/pom.xml
index 2baf82bf21c..3bb68a29562 100644
--- a/iotdb-core/node-commons/pom.xml
+++ b/iotdb-core/node-commons/pom.xml
@@ -164,6 +164,14 @@
<groupId>com.github.luben</groupId>
<artifactId>zstd-jni</artifactId>
</dependency>
+ <dependency>
+ <groupId>net.java.dev.jna</groupId>
+ <artifactId>jna</artifactId>
+ </dependency>
+ <dependency>
+ <groupId>net.java.dev.jna</groupId>
+ <artifactId>jna-platform</artifactId>
+ </dependency>
<dependency>
<groupId>org.reflections</groupId>
<artifactId>reflections</artifactId>
diff --git
a/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/ServiceMessages.java
b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/ServiceMessages.java
index b76c030498f..8c8bfbc0d29 100644
---
a/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/ServiceMessages.java
+++
b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/ServiceMessages.java
@@ -136,6 +136,10 @@ public final class ServiceMessages {
// ---- CpuUsageMetrics ----
public static final String CPU_USAGE_UPDATE_TIME = "Time for update cpu
usage is {} ns";
+ // ---- ProcessMetrics ----
+ public static final String FAILED_GET_PROCESS_RESIDENT_MEMORY =
+ "Failed to get process resident memory for pid {}";
+
private ServiceMessages() {}
public static final String UNKNOWN_SERVICE_TYPE = "Unknown ServiceType: ";
diff --git
a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/ServiceMessages.java
b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/ServiceMessages.java
index d350a74ce08..a96ae6f449d 100644
---
a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/ServiceMessages.java
+++
b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/ServiceMessages.java
@@ -136,6 +136,10 @@ public final class ServiceMessages {
// ---- CpuUsageMetrics ----
public static final String CPU_USAGE_UPDATE_TIME = "CPU 使用率更新耗时 {} 纳秒";
+ // ---- ProcessMetrics ----
+ public static final String FAILED_GET_PROCESS_RESIDENT_MEMORY =
+ "获取进程 {} 的常驻内存失败";
+
private ServiceMessages() {}
public static final String UNKNOWN_SERVICE_TYPE = "未知服务类型:";
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/ProcessMetrics.java
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/ProcessMetrics.java
similarity index 97%
rename from
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/ProcessMetrics.java
rename to
iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/ProcessMetrics.java
index 75335614894..9d2ddfdb332 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/metrics/ProcessMetrics.java
+++
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/ProcessMetrics.java
@@ -7,7 +7,7 @@
* "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
+ * 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
@@ -17,10 +17,10 @@
* under the License.
*/
-package org.apache.iotdb.db.service.metrics;
+package org.apache.iotdb.commons.service.metric;
+import org.apache.iotdb.commons.i18n.ServiceMessages;
import org.apache.iotdb.commons.service.metric.enums.Tag;
-import org.apache.iotdb.db.i18n.DataNodeMiscMessages;
import org.apache.iotdb.metrics.AbstractMetricService;
import org.apache.iotdb.metrics.MetricConstant;
import org.apache.iotdb.metrics.config.MetricConfig;
@@ -253,7 +253,7 @@ public class ProcessMetrics implements IMetricSet {
return 0L;
}
} catch (Exception e) {
- LOGGER.debug(DataNodeMiscMessages.FAILED_GET_PROCESS_RESIDENT_MEMORY,
CONFIG.getPid(), e);
+ LOGGER.debug(ServiceMessages.FAILED_GET_PROCESS_RESIDENT_MEMORY,
CONFIG.getPid(), e);
return 0L;
}
}