From 650e9ae94926134bdefde39d9a48d1a442e7b468 Mon Sep 17 00:00:00 2001 From: Tian Jiang Date: Tue, 1 Sep 2026 16:49:44 +0800 Subject: [PATCH 1/4] Optimize Prometheus reporter snapshot updates --- iotdb-core/metrics/ReadMe.md | 1 + .../iotdb/metrics/i18n/MetricsMessages.java | 1 + .../iotdb/metrics/i18n/MetricsMessages.java | 1 + .../iotdb/metrics/config/MetricConfig.java | 14 ++++ .../config/MetricConfigDescriptor.java | 25 ++++++ .../prometheus/PrometheusReporter.java | 81 ++++++++++++++++++- .../metrics/config/MetricConfigTest.java | 4 + .../conf/iotdb-system.properties.template | 6 ++ .../iotdb/commons/concurrent/ThreadName.java | 2 + .../commons/service/metric/MetricService.java | 12 ++- 10 files changed, 142 insertions(+), 5 deletions(-) diff --git a/iotdb-core/metrics/ReadMe.md b/iotdb-core/metrics/ReadMe.md index f0eb361f0a2b1..6ab575abe0984 100644 --- a/iotdb-core/metrics/ReadMe.md +++ b/iotdb-core/metrics/ReadMe.md @@ -63,6 +63,7 @@ Configure the metrics module through `iotdb-system.properties`. The main options | `dn(cn)_metric_level` | Initial metric level. | `OFF`, `CORE`, `IMPORTANT`, `NORMAL`, `ALL` | | `cn_metric_prometheus_reporter_port` | Prometheus HTTP port for ConfigNode. | `9091` | | `dn_metric_prometheus_reporter_port` | Prometheus HTTP port for DataNode. | `9092` | +| `prometheus_reporter_async_update` | Serve a cached Prometheus snapshot refreshed every 15 seconds. | `true` | More details, see the User Guide and the `iotdb-system.properties.template` file. diff --git a/iotdb-core/metrics/interface/src/main/i18n/en/org/apache/iotdb/metrics/i18n/MetricsMessages.java b/iotdb-core/metrics/interface/src/main/i18n/en/org/apache/iotdb/metrics/i18n/MetricsMessages.java index 778da2b32bfba..4cd7a7fd0a1d5 100644 --- a/iotdb-core/metrics/interface/src/main/i18n/en/org/apache/iotdb/metrics/i18n/MetricsMessages.java +++ b/iotdb-core/metrics/interface/src/main/i18n/en/org/apache/iotdb/metrics/i18n/MetricsMessages.java @@ -128,5 +128,6 @@ private MetricsMessages() {} public static final String LOG_IOTDBSESSIONREPORTER_START_WRITE_ARG_ARG_E79CDDAE = "IoTDBSessionReporter start, write to {}:{}"; public static final String LOG_PROMETHEUSREPORTER_STARTED_USE_PORT_ARG_A688FFC8 = "PrometheusReporter started, use port {}"; public static final String LOG_DETECTED_ERROR_TAKING_METRIC_TIMER_SNAPSHOT_WILL_DISCARD_METRIC_B7154169 = "Detected an error when taking metric timer snapshot, will discard this metric"; + public static final String LOG_PROMETHEUSREPORTER_FAILED_TO_UPDATE_METRICS_SNAPSHOT_ASYNCHRONOUSLY_F19FE4E3 = "PrometheusReporter failed to update metrics snapshot asynchronously"; } diff --git a/iotdb-core/metrics/interface/src/main/i18n/zh/org/apache/iotdb/metrics/i18n/MetricsMessages.java b/iotdb-core/metrics/interface/src/main/i18n/zh/org/apache/iotdb/metrics/i18n/MetricsMessages.java index fd68351a93974..968694f6eb006 100644 --- a/iotdb-core/metrics/interface/src/main/i18n/zh/org/apache/iotdb/metrics/i18n/MetricsMessages.java +++ b/iotdb-core/metrics/interface/src/main/i18n/zh/org/apache/iotdb/metrics/i18n/MetricsMessages.java @@ -124,5 +124,6 @@ private MetricsMessages() {} public static final String LOG_IOTDBSESSIONREPORTER_START_WRITE_ARG_ARG_E79CDDAE = "IoTDBSessionReporter 启动,写入 {}:{}"; public static final String LOG_PROMETHEUSREPORTER_STARTED_USE_PORT_ARG_A688FFC8 = "PrometheusReporter 已启动,使用端口 {}"; public static final String LOG_DETECTED_ERROR_TAKING_METRIC_TIMER_SNAPSHOT_WILL_DISCARD_METRIC_B7154169 = "获取 metric timer 快照时检测到错误,将丢弃该 metric"; + public static final String LOG_PROMETHEUSREPORTER_FAILED_TO_UPDATE_METRICS_SNAPSHOT_ASYNCHRONOUSLY_F19FE4E3 = "PrometheusReporter 异步更新监控项快照失败"; } diff --git a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfig.java b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfig.java index 2e0eaf5cb0b31..c5d55162f8cf5 100644 --- a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfig.java +++ b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/config/MetricConfig.java @@ -52,6 +52,9 @@ public class MetricConfig { /** The export port for prometheus to get metrics. */ private Integer prometheusReporterPort = 9091; + /** Whether Prometheus metrics are collected asynchronously into a cached snapshot. */ + private boolean prometheusReporterAsyncUpdate = true; + private String prometheusReporterUsername = ""; private String prometheusReporterPassword = ""; @@ -140,6 +143,14 @@ public void setPrometheusReporterPort(Integer prometheusReporterPort) { this.prometheusReporterPort = prometheusReporterPort; } + public boolean isPrometheusReporterAsyncUpdate() { + return prometheusReporterAsyncUpdate; + } + + public void setPrometheusReporterAsyncUpdate(boolean prometheusReporterAsyncUpdate) { + this.prometheusReporterAsyncUpdate = prometheusReporterAsyncUpdate; + } + public boolean prometheusNeedAuth() { return prometheusReporterUsername != null && !prometheusReporterUsername.isEmpty(); } @@ -264,6 +275,7 @@ public void copy(MetricConfig newMetricConfig) { metricLevel = newMetricConfig.getMetricLevel(); asyncCollectPeriodInSecond = newMetricConfig.getAsyncCollectPeriodInSecond(); prometheusReporterPort = newMetricConfig.getPrometheusReporterPort(); + prometheusReporterAsyncUpdate = newMetricConfig.isPrometheusReporterAsyncUpdate(); prometheusReporterUsername = newMetricConfig.getPrometheusReporterUsername(); prometheusReporterPassword = newMetricConfig.getPrometheusReporterPassword(); internalReporterType = newMetricConfig.getInternalReportType(); @@ -287,6 +299,7 @@ public boolean equals(Object obj) { && metricLevel.equals(anotherMetricConfig.getMetricLevel()) && asyncCollectPeriodInSecond.equals(anotherMetricConfig.getAsyncCollectPeriodInSecond()) && prometheusReporterPort.equals(anotherMetricConfig.getPrometheusReporterPort()) + && prometheusReporterAsyncUpdate == anotherMetricConfig.isPrometheusReporterAsyncUpdate() && iotdbReporterConfig.equals(anotherMetricConfig.getIoTDBReporterConfig()) && internalReporterType.equals(anotherMetricConfig.getInternalReportType()); } @@ -298,6 +311,7 @@ public int hashCode() { metricLevel, asyncCollectPeriodInSecond, prometheusReporterPort, + prometheusReporterAsyncUpdate, iotdbReporterConfig, internalReporterType); } 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..50dd0127ddfcd 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 @@ -112,6 +112,11 @@ private MetricConfig generateFromProperties(Properties properties, boolean isCon properties, isConfigNode))); + loadConfig.setPrometheusReporterAsyncUpdate( + Boolean.parseBoolean( + getPrometheusReporterAsyncUpdateProperty( + properties, isConfigNode, loadConfig.isPrometheusReporterAsyncUpdate()))); + loadConfig.setPrometheusReporterUsername( getPropertyWithoutPrefix( "metric_prometheus_reporter_username", @@ -213,6 +218,26 @@ private String getPropertyWithoutPrefix( .orElse(defaultValue); } + private String getPrometheusReporterAsyncUpdateProperty( + Properties properties, boolean isConfigNode, boolean defaultValue) { + String value = properties.getProperty("prometheus_reporter_async_update"); + if (value == null) { + // Keep accepting the metric-prefixed forms for compatibility with node-specific configs. + value = properties.getProperty("metric_prometheus_reporter_async_update"); + } + if (value == null) { + value = + properties.getProperty( + (isConfigNode ? "cn_" : "dn_") + "prometheus_reporter_async_update"); + } + if (value == null) { + value = + properties.getProperty( + (isConfigNode ? "cn_" : "dn_") + "metric_prometheus_reporter_async_update"); + } + return value == null ? String.valueOf(defaultValue) : value.trim(); + } + private static class MetricConfigDescriptorHolder { private static final MetricConfigDescriptor INSTANCE = new MetricConfigDescriptor(); } diff --git a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java index 7846e056f8cdd..0316e4fffe9fe 100644 --- a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java +++ b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java @@ -66,25 +66,43 @@ import java.util.HashMap; import java.util.Map; import java.util.Objects; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; public class PrometheusReporter implements Reporter { private static final Logger LOGGER = LoggerFactory.getLogger(PrometheusReporter.class); private static final MetricConfig METRIC_CONFIG = MetricConfigDescriptor.getInstance().getMetricConfig(); + private static final long PROMETHEUS_DEFAULT_SCRAPE_INTERVAL_SECONDS = 15; private final AbstractMetricManager metricManager; - private DisposableServer httpServer; + private volatile ScheduledExecutorService snapshotUpdateExecutor; + private volatile DisposableServer httpServer; + private volatile String metricsSnapshot = ""; + private volatile ScheduledFuture snapshotUpdateFuture; private static final String REALM = "metrics"; public static final String BASIC_AUTH_PREFIX = "Basic "; public static final char DIVIDER_BETWEEN_USERNAME_AND_DIVIDER = ':'; + /** + * Creates a reporter with a self-managed scheduler for compatibility with standalone users. + * Server-side code should use the constructor accepting the IoTDB thread pool. + */ public PrometheusReporter(AbstractMetricManager metricManager) { + this(metricManager, null); + } + + public PrometheusReporter( + AbstractMetricManager metricManager, ScheduledExecutorService snapshotUpdateExecutor) { this.metricManager = metricManager; + this.snapshotUpdateExecutor = snapshotUpdateExecutor; } @Override @SuppressWarnings("java:S1181") - public boolean start() { + public synchronized boolean start() { if (httpServer != null) { LOGGER.warn(MetricsMessages.PROMETHEUS_REPORTER_ALREADY_START); return false; @@ -104,8 +122,12 @@ public boolean start() { // authenticate not pass return Mono.empty(); } + String metrics = + METRIC_CONFIG.isPrometheusReporterAsyncUpdate() + ? metricsSnapshot + : scrape(); return res.header(HttpHeaderNames.CONTENT_TYPE, "text/plain") - .sendString(Mono.just(scrape())); + .sendString(Mono.just(metrics)); })); if (METRIC_CONFIG.isEnableSSL()) { SslContext sslContext; @@ -122,9 +144,20 @@ public boolean start() { serverTransport = serverTransport.secure(spec -> spec.sslContext(sslContext)); } httpServer = serverTransport.bindNow(); + if (METRIC_CONFIG.isPrometheusReporterAsyncUpdate()) { + startSnapshotUpdater(); + } } catch (Throwable e) { // catch Throwable rather than Exception here because the code above might cause a // NoClassDefFoundError + stopSnapshotUpdater(); + if (httpServer != null) { + try { + httpServer.disposeNow(Duration.ofSeconds(10)); + } catch (Exception ignored) { + // do nothing + } + } httpServer = null; LOGGER.warn(MetricsMessages.PROMETHEUS_REPORTER_START_FAILED, e); return false; @@ -135,6 +168,45 @@ public boolean start() { return true; } + @SuppressWarnings("unsafeThreadSchedule") + private void startSnapshotUpdater() { + // Keep metric collection off Reactor HTTP threads and avoid overlapping scrapes. + if (snapshotUpdateExecutor == null) { + snapshotUpdateExecutor = + Executors.newSingleThreadScheduledExecutor( + runnable -> { + Thread thread = new Thread(runnable, "prometheus-reporter-snapshot-updater"); + thread.setDaemon(true); + return thread; + }); + } + snapshotUpdateFuture = + snapshotUpdateExecutor.scheduleAtFixedRate( + this::updateSnapshot, 0, PROMETHEUS_DEFAULT_SCRAPE_INTERVAL_SECONDS, TimeUnit.SECONDS); + } + + private void updateSnapshot() { + try { + metricsSnapshot = scrape(); + } catch (Throwable t) { + LOGGER.error( + MetricsMessages + .LOG_PROMETHEUSREPORTER_FAILED_TO_UPDATE_METRICS_SNAPSHOT_ASYNCHRONOUSLY_F19FE4E3, + t); + } + } + + private void stopSnapshotUpdater() { + if (snapshotUpdateFuture != null) { + snapshotUpdateFuture.cancel(false); + snapshotUpdateFuture = null; + } + if (snapshotUpdateExecutor != null) { + snapshotUpdateExecutor.shutdownNow(); + snapshotUpdateExecutor = null; + } + } + private boolean authenticate(HttpServerRequest req, HttpServerResponse res) { if (!METRIC_CONFIG.prometheusNeedAuth()) { return true; @@ -318,7 +390,8 @@ private SslContext createSslContext( } @Override - public boolean stop() { + public synchronized boolean stop() { + stopSnapshotUpdater(); if (httpServer != null) { try { httpServer.disposeNow(Duration.ofSeconds(10)); diff --git a/iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/config/MetricConfigTest.java b/iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/config/MetricConfigTest.java index 30d3bd1ff374c..3654010c4f686 100644 --- a/iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/config/MetricConfigTest.java +++ b/iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/config/MetricConfigTest.java @@ -39,6 +39,7 @@ public void testConfigNodeMetricConfig() { properties.setProperty("cn_metric_level", "ALL"); properties.setProperty("cn_metric_async_collect_period", "10"); properties.setProperty("cn_metric_prometheus_reporter_port", "9090"); + properties.setProperty("prometheus_reporter_async_update", "false"); properties.setProperty("cn_metric_iotdb_reporter_host", "0.0.0.0"); properties.setProperty("cn_metric_iotdb_reporter_port", "6669"); properties.setProperty("cn_metric_iotdb_reporter_username", "user"); @@ -55,6 +56,7 @@ public void testConfigNodeMetricConfig() { assertEquals(MetricLevel.ALL, metricConfig.getMetricLevel()); assertEquals(10, (int) metricConfig.getAsyncCollectPeriodInSecond()); assertEquals(9090, (int) metricConfig.getPrometheusReporterPort()); + assertEquals(false, metricConfig.isPrometheusReporterAsyncUpdate()); IoTDBReporterConfig reporterConfig = metricConfig.getIoTDBReporterConfig(); assertEquals("0.0.0.0", reporterConfig.getHost()); @@ -75,6 +77,7 @@ public void testDataNodeMetricConfig() { properties.setProperty("dn_metric_level", "ALL"); properties.setProperty("dn_metric_async_collect_period", "10"); properties.setProperty("dn_metric_prometheus_reporter_port", "9090"); + properties.setProperty("metric_prometheus_reporter_async_update", "true"); properties.setProperty("dn_metric_iotdb_reporter_host", "0.0.0.0"); properties.setProperty("dn_metric_iotdb_reporter_port", "6669"); properties.setProperty("dn_metric_iotdb_reporter_username", "user"); @@ -92,6 +95,7 @@ public void testDataNodeMetricConfig() { assertEquals(MetricLevel.ALL, metricConfig.getMetricLevel()); assertEquals(10, (int) metricConfig.getAsyncCollectPeriodInSecond()); assertEquals(9090, (int) metricConfig.getPrometheusReporterPort()); + assertEquals(true, metricConfig.isPrometheusReporterAsyncUpdate()); IoTDBReporterConfig reporterConfig = metricConfig.getIoTDBReporterConfig(); assertEquals("0.0.0.0", reporterConfig.getHost()); diff --git a/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template b/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template index e325ff4257302..89b1397ba5708 100644 --- a/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template +++ b/iotdb-core/node-commons/src/assembly/resources/conf/iotdb-system.properties.template @@ -382,6 +382,12 @@ metric_prometheus_reporter_username= # Datatype: String metric_prometheus_reporter_password= +# Whether Prometheus metrics are collected asynchronously and served from a cached snapshot. +# The snapshot is refreshed every 15 seconds, matching Prometheus's default scrape interval. +# effectiveMode: restart +# Datatype: boolean +prometheus_reporter_async_update=true + # The reporters of metric module to report metrics # If there are more than one reporter, please separate them by commas ",". # Options: [JMX, PROMETHEUS] diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java index 606a08bde3322..a0b62cd6a62c7 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java @@ -191,6 +191,7 @@ public enum ThreadName { // -------------------------- Metrics -------------------------- SYSTEM_SCHEDULE_METRICS("SystemScheduleMetrics"), RESOURCE_CONTROL_DISK_STATISTIC("ResourceControl-DataRegionDiskStatistics"), + PROMETHEUS_REPORTER_SNAPSHOT_UPDATER("PrometheusReporter-Snapshot-Updater"), PROMETHEUS_REACTOR_HTTP_EPOLL("reactor-http-epoll"), PROMETHEUS_REACTOR_HTTP_NIO("reactor-http-nio"), PROMETHEUS_BOUNDED_ELASTIC("boundedElastic-evictor"), @@ -394,6 +395,7 @@ public enum ThreadName { Arrays.asList( SYSTEM_SCHEDULE_METRICS, RESOURCE_CONTROL_DISK_STATISTIC, + PROMETHEUS_REPORTER_SNAPSHOT_UPDATER, PROMETHEUS_REACTOR_HTTP_EPOLL, PROMETHEUS_REACTOR_HTTP_NIO, PROMETHEUS_REACTOR_HTTP_EPOLL, diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/MetricService.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/MetricService.java index 51357fe8998e3..a54d4ca483d94 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/MetricService.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/MetricService.java @@ -19,6 +19,8 @@ package org.apache.iotdb.commons.service.metric; +import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; +import org.apache.iotdb.commons.concurrent.ThreadName; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.exception.StartupException; import org.apache.iotdb.commons.i18n.ServiceMessages; @@ -74,7 +76,15 @@ protected void loadReporter() { hasJmxReporter = true; break; case PROMETHEUS: - reporter = new PrometheusReporter(metricManager); + if (METRIC_CONFIG.isPrometheusReporterAsyncUpdate()) { + reporter = + new PrometheusReporter( + metricManager, + IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor( + ThreadName.PROMETHEUS_REPORTER_SNAPSHOT_UPDATER.getName())); + } else { + reporter = new PrometheusReporter(metricManager); + } break; case IOTDB: reporter = new IoTDBSessionReporter(metricManager); From 17c11517eb41aaaaa6e85d5aa84b9ef8b8a9934f Mon Sep 17 00:00:00 2001 From: Tian Jiang Date: Tue, 1 Sep 2026 18:11:27 +0800 Subject: [PATCH 2/4] fix compilation --- .../iotdb/metrics/config/MetricConfigDescriptor.java | 12 ++++-------- .../iotdb/metrics/config/MetricConfigTest.java | 12 ++++++++++++ 2 files changed, 16 insertions(+), 8 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 d7a08a09a549b..113bff7ebcc66 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 @@ -132,7 +132,7 @@ private MetricConfig generateFromProperties(Properties properties, String prefix loadConfig.setPrometheusReporterAsyncUpdate( Boolean.parseBoolean( getPrometheusReporterAsyncUpdateProperty( - properties, isConfigNode, loadConfig.isPrometheusReporterAsyncUpdate()))); + properties, prefix, loadConfig.isPrometheusReporterAsyncUpdate()))); loadConfig.setPrometheusReporterUsername( getPropertyWithoutPrefix( @@ -228,21 +228,17 @@ private String getPropertyWithoutPrefix( } private String getPrometheusReporterAsyncUpdateProperty( - Properties properties, boolean isConfigNode, boolean defaultValue) { + Properties properties, String prefix, boolean defaultValue) { String value = properties.getProperty("prometheus_reporter_async_update"); if (value == null) { // Keep accepting the metric-prefixed forms for compatibility with node-specific configs. value = properties.getProperty("metric_prometheus_reporter_async_update"); } if (value == null) { - value = - properties.getProperty( - (isConfigNode ? "cn_" : "dn_") + "prometheus_reporter_async_update"); + value = properties.getProperty(prefix + "prometheus_reporter_async_update"); } if (value == null) { - value = - properties.getProperty( - (isConfigNode ? "cn_" : "dn_") + "metric_prometheus_reporter_async_update"); + value = properties.getProperty(prefix + "metric_prometheus_reporter_async_update"); } return value == null ? String.valueOf(defaultValue) : value.trim(); } diff --git a/iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/config/MetricConfigTest.java b/iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/config/MetricConfigTest.java index 3654010c4f686..3d1d3a0b0cff7 100644 --- a/iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/config/MetricConfigTest.java +++ b/iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/config/MetricConfigTest.java @@ -107,4 +107,16 @@ public void testDataNodeMetricConfig() { assertEquals(5, (int) reporterConfig.getPushPeriodInSecond()); assertEquals(InternalReporterType.IOTDB, metricConfig.getInternalReportType()); } + + @Test + public void testMetricConfigWithCustomNodePrefix() { + Properties properties = new Properties(); + properties.setProperty("sn_metric_prometheus_reporter_async_update", "false"); + + MetricConfigDescriptor.getInstance().loadProps(properties, "sn_"); + + assertEquals( + false, + MetricConfigDescriptor.getInstance().getMetricConfig().isPrometheusReporterAsyncUpdate()); + } } From 8b6d0f65435e3c5749e4578ee34eddd542e97254 Mon Sep 17 00:00:00 2001 From: Tian Jiang Date: Wed, 2 Sep 2026 11:14:09 +0800 Subject: [PATCH 3/4] Fix Prometheus reporter initial snapshot race --- .../prometheus/PrometheusReporter.java | 34 ++++++++++++++++--- 1 file changed, 30 insertions(+), 4 deletions(-) diff --git a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java index 0316e4fffe9fe..fa5695e96213b 100644 --- a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java +++ b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java @@ -79,7 +79,10 @@ public class PrometheusReporter implements Reporter { private final AbstractMetricManager metricManager; private volatile ScheduledExecutorService snapshotUpdateExecutor; private volatile DisposableServer httpServer; - private volatile String metricsSnapshot = ""; + + /** A null snapshot means that no complete scrape has been published yet. */ + private volatile String metricsSnapshot; + private volatile ScheduledFuture snapshotUpdateFuture; private static final String REALM = "metrics"; @@ -107,6 +110,8 @@ public synchronized boolean start() { LOGGER.warn(MetricsMessages.PROMETHEUS_REPORTER_ALREADY_START); return false; } + // A reporter can be started again after its metric manager has been reset. + metricsSnapshot = null; try { HttpServer serverTransport = HttpServer.create() @@ -124,7 +129,7 @@ public synchronized boolean start() { } String metrics = METRIC_CONFIG.isPrometheusReporterAsyncUpdate() - ? metricsSnapshot + ? getMetricsSnapshot() : scrape(); return res.header(HttpHeaderNames.CONTENT_TYPE, "text/plain") .sendString(Mono.just(metrics)); @@ -180,14 +185,24 @@ private void startSnapshotUpdater() { return thread; }); } + // Delay the first background scrape until metric sets have been bound by the metric service. snapshotUpdateFuture = snapshotUpdateExecutor.scheduleAtFixedRate( - this::updateSnapshot, 0, PROMETHEUS_DEFAULT_SCRAPE_INTERVAL_SECONDS, TimeUnit.SECONDS); + this::updateSnapshot, + PROMETHEUS_DEFAULT_SCRAPE_INTERVAL_SECONDS, + PROMETHEUS_DEFAULT_SCRAPE_INTERVAL_SECONDS, + TimeUnit.SECONDS); } private void updateSnapshot() { try { - metricsSnapshot = scrape(); + String snapshot = scrape(); + // Do not publish an empty scrape taken before metric sets are bound. The request path will + // synchronously scrape until the first complete snapshot is available. Empty snapshots are + // published after initialization so removed metrics are not kept in the cache indefinitely. + if (!snapshot.isEmpty() || metricsSnapshot != null) { + metricsSnapshot = snapshot; + } } catch (Throwable t) { LOGGER.error( MetricsMessages @@ -196,6 +211,17 @@ private void updateSnapshot() { } } + private String getMetricsSnapshot() { + String snapshot = metricsSnapshot; + if (snapshot == null) { + snapshot = scrape(); + if (!snapshot.isEmpty()) { + metricsSnapshot = snapshot; + } + } + return snapshot; + } + private void stopSnapshotUpdater() { if (snapshotUpdateFuture != null) { snapshotUpdateFuture.cancel(false); From 63295e6c8bb8d4b974d1f5b868190e22cb538441 Mon Sep 17 00:00:00 2001 From: Tian Jiang Date: Wed, 2 Sep 2026 16:36:19 +0800 Subject: [PATCH 4/4] Fix Prometheus scheduler lifecycle and duplicate allocation --- .../prometheus/PrometheusReporter.java | 43 +++++++---- .../prometheus/PrometheusReporterTest.java | 75 +++++++++++++++++++ .../commons/service/metric/MetricService.java | 7 +- 3 files changed, 108 insertions(+), 17 deletions(-) create mode 100644 iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporterTest.java diff --git a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java index fa5695e96213b..00ae6c6b98810 100644 --- a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java +++ b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporter.java @@ -70,6 +70,7 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; +import java.util.function.Supplier; public class PrometheusReporter implements Reporter { private static final Logger LOGGER = LoggerFactory.getLogger(PrometheusReporter.class); @@ -77,6 +78,7 @@ public class PrometheusReporter implements Reporter { MetricConfigDescriptor.getInstance().getMetricConfig(); private static final long PROMETHEUS_DEFAULT_SCRAPE_INTERVAL_SECONDS = 15; private final AbstractMetricManager metricManager; + private final Supplier snapshotUpdateExecutorSupplier; private volatile ScheduledExecutorService snapshotUpdateExecutor; private volatile DisposableServer httpServer; @@ -91,16 +93,30 @@ public class PrometheusReporter implements Reporter { /** * Creates a reporter with a self-managed scheduler for compatibility with standalone users. - * Server-side code should use the constructor accepting the IoTDB thread pool. + * Server-side code should use the constructor accepting a scheduler factory. */ public PrometheusReporter(AbstractMetricManager metricManager) { - this(metricManager, null); + this(metricManager, PrometheusReporter::newStandaloneSnapshotUpdateExecutor); } + /** + * Creates a reporter with a scheduler factory. The factory is invoked on every start to obtain a + * fresh executor for the reporter lifecycle. + */ public PrometheusReporter( - AbstractMetricManager metricManager, ScheduledExecutorService snapshotUpdateExecutor) { + AbstractMetricManager metricManager, + Supplier snapshotUpdateExecutorSupplier) { this.metricManager = metricManager; - this.snapshotUpdateExecutor = snapshotUpdateExecutor; + this.snapshotUpdateExecutorSupplier = Objects.requireNonNull(snapshotUpdateExecutorSupplier); + } + + private static ScheduledExecutorService newStandaloneSnapshotUpdateExecutor() { + return Executors.newSingleThreadScheduledExecutor( + runnable -> { + Thread thread = new Thread(runnable, "prometheus-reporter-snapshot-updater"); + thread.setDaemon(true); + return thread; + }); } @Override @@ -176,14 +192,10 @@ public synchronized boolean start() { @SuppressWarnings("unsafeThreadSchedule") private void startSnapshotUpdater() { // Keep metric collection off Reactor HTTP threads and avoid overlapping scrapes. - if (snapshotUpdateExecutor == null) { - snapshotUpdateExecutor = - Executors.newSingleThreadScheduledExecutor( - runnable -> { - Thread thread = new Thread(runnable, "prometheus-reporter-snapshot-updater"); - thread.setDaemon(true); - return thread; - }); + if (snapshotUpdateExecutor == null || snapshotUpdateExecutor.isShutdown()) { + // Create a fresh executor for every start so a stopped reporter can be started again with + // the same managed thread-pool factory. + snapshotUpdateExecutor = Objects.requireNonNull(snapshotUpdateExecutorSupplier.get()); } // Delay the first background scrape until metric sets have been bound by the metric service. snapshotUpdateFuture = @@ -227,9 +239,10 @@ private void stopSnapshotUpdater() { snapshotUpdateFuture.cancel(false); snapshotUpdateFuture = null; } - if (snapshotUpdateExecutor != null) { - snapshotUpdateExecutor.shutdownNow(); - snapshotUpdateExecutor = null; + ScheduledExecutorService executor = snapshotUpdateExecutor; + snapshotUpdateExecutor = null; + if (executor != null) { + executor.shutdownNow(); } } diff --git a/iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporterTest.java b/iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporterTest.java new file mode 100644 index 0000000000000..7072777ac6669 --- /dev/null +++ b/iotdb-core/metrics/interface/src/test/java/org/apache/iotdb/metrics/reporter/prometheus/PrometheusReporterTest.java @@ -0,0 +1,75 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iotdb.metrics.reporter.prometheus; + +import org.apache.iotdb.metrics.config.MetricConfig; +import org.apache.iotdb.metrics.config.MetricConfigDescriptor; +import org.apache.iotdb.metrics.impl.DoNothingMetricManager; + +import org.junit.Test; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +public class PrometheusReporterTest { + + @Test + public void testManagedExecutorRecreatedAfterRestart() { + MetricConfig metricConfig = MetricConfigDescriptor.getInstance().getMetricConfig(); + boolean originalAsyncUpdate = metricConfig.isPrometheusReporterAsyncUpdate(); + Integer originalPort = metricConfig.getPrometheusReporterPort(); + metricConfig.setPrometheusReporterAsyncUpdate(true); + metricConfig.setPrometheusReporterPort(0); + + AtomicInteger factoryCalls = new AtomicInteger(); + List executors = new ArrayList<>(); + PrometheusReporter reporter = + new PrometheusReporter( + new DoNothingMetricManager(), + () -> { + factoryCalls.incrementAndGet(); + ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); + executors.add(executor); + return executor; + }); + try { + assertTrue(reporter.start()); + assertEquals(1, factoryCalls.get()); + assertTrue(reporter.stop()); + assertTrue(executors.get(0).isShutdown()); + + assertTrue(reporter.start()); + assertEquals(2, factoryCalls.get()); + assertTrue(reporter.stop()); + assertTrue(executors.get(1).isShutdown()); + } finally { + reporter.stop(); + metricConfig.setPrometheusReporterAsyncUpdate(originalAsyncUpdate); + metricConfig.setPrometheusReporterPort(originalPort); + executors.forEach(ScheduledExecutorService::shutdownNow); + } + } +} diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/MetricService.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/MetricService.java index a54d4ca483d94..4af78e3d9cfc5 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/MetricService.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/MetricService.java @@ -77,11 +77,14 @@ protected void loadReporter() { break; case PROMETHEUS: if (METRIC_CONFIG.isPrometheusReporterAsyncUpdate()) { + // Defer pool creation until start so duplicate reporters rejected below do not + // register an unused pool. reporter = new PrometheusReporter( metricManager, - IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor( - ThreadName.PROMETHEUS_REPORTER_SNAPSHOT_UPDATER.getName())); + () -> + IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor( + ThreadName.PROMETHEUS_REPORTER_SNAPSHOT_UPDATER.getName())); } else { reporter = new PrometheusReporter(metricManager); }