From df86a4c039e1f8b9cfb4c84ce989abac29c45697 Mon Sep 17 00:00:00 2001 From: suchenglong <404083629@qq.com> Date: Thu, 13 Aug 2026 17:55:50 +0800 Subject: [PATCH 1/4] support streamnode --- .../config/MetricConfigDescriptor.java | 74 +++++++++++-------- 1 file changed, 42 insertions(+), 32 deletions(-) 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 def3d1e50c057..85e9a89c14588 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 void loadProps(Properties properties, boolean isConfigNode) { * @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 ReloadLevel loadHotProps(Properties properties, boolean isConfigNode) { } /** 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 @@ private MetricConfig generateFromProperties(Properties properties, boolean isCon .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 @@ private MetricConfig generateFromProperties(Properties properties, boolean isCon "metric_async_collect_period", String.valueOf(loadConfig.getAsyncCollectPeriodInSecond()), properties, - isConfigNode))); + prefix))); loadConfig.setPrometheusReporterPort( Integer.parseInt( @@ -110,7 +127,7 @@ private MetricConfig generateFromProperties(Properties properties, boolean isCon "metric_prometheus_reporter_port", String.valueOf(loadConfig.getPrometheusReporterPort()), properties, - isConfigNode))); + prefix))); loadConfig.setPrometheusReporterUsername( getPropertyWithoutPrefix( @@ -139,8 +156,7 @@ private MetricConfig generateFromProperties(Properties properties, boolean isCon 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 @@ private MetricConfig generateFromProperties(Properties properties, boolean isCon "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 @@ private MetricConfig generateFromProperties(Properties properties, boolean isCon "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,10 @@ private MetricConfig generateFromProperties(Properties properties, boolean isCon "metric_iotdb_reporter_push_period", String.valueOf(reporterConfig.getPushPeriodInSecond()), properties, - isConfigNode))); - if (!isConfigNode) { + prefix))); + // Internal reporter writes metrics into IoTDB internal tables, which only DataNode + // (prefix "dn_") has storage to host. ConfigNode and StreamNode skip this config. + if (DATA_NODE_PREFIX.equals(prefix)) { loadConfig.setInternalReportType( InternalReporterType.valueOf( properties.getProperty( @@ -197,11 +206,12 @@ private MetricConfig generateFromProperties(Properties properties, boolean isCon 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); } From b1c63883a958d4d6e63d5f9f3ab0349e7c85cdbe Mon Sep 17 00:00:00 2001 From: suchenglong <404083629@qq.com> Date: Fri, 14 Aug 2026 08:38:25 +0800 Subject: [PATCH 2/4] remove unnecessary comments --- .../apache/iotdb/metrics/config/MetricConfigDescriptor.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) 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 85e9a89c14588..0d83d38740972 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 @@ -193,8 +193,7 @@ private MetricConfig generateFromProperties(Properties properties, String prefix String.valueOf(reporterConfig.getPushPeriodInSecond()), properties, prefix))); - // Internal reporter writes metrics into IoTDB internal tables, which only DataNode - // (prefix "dn_") has storage to host. ConfigNode and StreamNode skip this config. + if (DATA_NODE_PREFIX.equals(prefix)) { loadConfig.setInternalReportType( InternalReporterType.valueOf( From 091cb227019f1a7583bb38ebcbb0d28927eda308 Mon Sep 17 00:00:00 2001 From: suchenglong <404083629@qq.com> Date: Fri, 28 Aug 2026 09:39:42 +0800 Subject: [PATCH 3/4] metric add stream node type --- .../src/main/java/org/apache/iotdb/metrics/utils/NodeType.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) 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 1ffa95030c391..e2c826780f7db 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 @@ public enum NodeType { CONFIGNODE, - DATANODE; + DATANODE, + STREAMNODE; @Override public String toString() { From 4c49d016bf9f447d62e6856a70ed9454358c45e6 Mon Sep 17 00:00:00 2001 From: suchenglong <404083629@qq.com> Date: Mon, 31 Aug 2026 17:20:39 +0800 Subject: [PATCH 4/4] =?UTF-8?q?Move=20and=20refactor=20the=20ProcessMetric?= =?UTF-8?q?s=20class=20to=20the=20node=E2=80=91commons=20module.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../org/apache/iotdb/confignode/service/ConfigNode.java | 2 +- .../en/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java | 2 -- .../zh/org/apache/iotdb/db/i18n/DataNodeMiscMessages.java | 2 -- .../iotdb/db/service/metrics/DataNodeMetricsHelper.java | 1 + iotdb-core/node-commons/pom.xml | 8 ++++++++ .../en/org/apache/iotdb/commons/i18n/ServiceMessages.java | 4 ++++ .../zh/org/apache/iotdb/commons/i18n/ServiceMessages.java | 4 ++++ .../iotdb/commons/service/metric}/ProcessMetrics.java | 8 ++++---- 8 files changed, 22 insertions(+), 9 deletions(-) rename iotdb-core/{datanode/src/main/java/org/apache/iotdb/db/service/metrics => node-commons/src/main/java/org/apache/iotdb/commons/service/metric}/ProcessMetrics.java (97%) 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 ffd5ec2d33b7d..cf49a3916bf2a 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.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.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 e76211672ff71..4aef9cd74f050 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 @@ private 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 c9a3f49710f80..bf909f2774521 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 @@ private 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 e2204e8cf0b57..1c257537d2b77 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.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/node-commons/pom.xml b/iotdb-core/node-commons/pom.xml index 2baf82bf21cb0..3bb68a29562c3 100644 --- a/iotdb-core/node-commons/pom.xml +++ b/iotdb-core/node-commons/pom.xml @@ -164,6 +164,14 @@ com.github.luben zstd-jni + + net.java.dev.jna + jna + + + net.java.dev.jna + jna-platform + org.reflections reflections 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 b76c030498fd4..8c8bfbc0d298b 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 d350a74ce085d..a96ae6f449d4a 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 75335614894bc..9d2ddfdb33278 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 @@ private long getResidentMemory() { 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; } }