From b521d506540dec8b54291b1e017e370c6cee1670 Mon Sep 17 00:00:00 2001 From: JackieTien97 Date: Sun, 20 Sep 2026 17:07:10 +0800 Subject: [PATCH 1/4] Fix client manager metric lifecycle and concurrent collection --- iotdb-core/metrics/core/pom.xml | 5 + .../metrics/core/type/IoTDBAutoGauge.java | 5 +- .../metrics/core/type/IoTDBAutoGaugeTest.java | 68 ++++ .../iotdb/commons/client/ClientManager.java | 12 +- .../commons/client/ClientManagerMetrics.java | 200 +++++---- .../client/ClientManagerMetricsTest.java | 384 ++++++++++++++++++ 6 files changed, 580 insertions(+), 94 deletions(-) create mode 100644 iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/type/IoTDBAutoGaugeTest.java create mode 100644 iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java diff --git a/iotdb-core/metrics/core/pom.xml b/iotdb-core/metrics/core/pom.xml index 2fa80bf1b1126..ac7f9ac32f967 100644 --- a/iotdb-core/metrics/core/pom.xml +++ b/iotdb-core/metrics/core/pom.xml @@ -46,6 +46,11 @@ org.slf4j slf4j-api + + junit + junit + test + diff --git a/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/type/IoTDBAutoGauge.java b/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/type/IoTDBAutoGauge.java index dc0a6df902c46..9b6e98fc2704e 100644 --- a/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/type/IoTDBAutoGauge.java +++ b/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/type/IoTDBAutoGauge.java @@ -37,9 +37,10 @@ public IoTDBAutoGauge(T object, ToDoubleFunction mapper) { @Override public double getValue() { - if (refObject.get() == null) { + T object = refObject.get(); + if (object == null) { return 0d; } - return mapper.applyAsDouble(refObject.get()); + return mapper.applyAsDouble(object); } } diff --git a/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/type/IoTDBAutoGaugeTest.java b/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/type/IoTDBAutoGaugeTest.java new file mode 100644 index 0000000000000..36131b3e498c7 --- /dev/null +++ b/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/type/IoTDBAutoGaugeTest.java @@ -0,0 +1,68 @@ +/* + * 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.core.type; + +import org.junit.Test; + +import java.lang.ref.WeakReference; +import java.lang.reflect.Field; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.Assert.assertEquals; + +public class IoTDBAutoGaugeTest { + + @Test + public void testLiveAndClearedReferent() throws Exception { + AtomicInteger value = new AtomicInteger(7); + IoTDBAutoGauge gauge = new IoTDBAutoGauge<>(value, AtomicInteger::get); + assertEquals(7, gauge.getValue(), 0); + value.set(9); + assertEquals(9, gauge.getValue(), 0); + + Field field = IoTDBAutoGauge.class.getDeclaredField("refObject"); + field.setAccessible(true); + ((WeakReference) field.get(gauge)).clear(); + assertEquals(0, gauge.getValue(), 0); + } + + @Test + public void testReferentRemainsAvailableForCurrentSample() throws Exception { + AtomicInteger value = new AtomicInteger(7); + IoTDBAutoGauge gauge = new IoTDBAutoGauge<>(value, AtomicInteger::get); + // Deterministically simulate collection between two weak-reference reads, without relying on + // GC. + WeakReference reference = + new WeakReference(value) { + @Override + public AtomicInteger get() { + AtomicInteger referent = super.get(); + clear(); + return referent; + } + }; + Field field = IoTDBAutoGauge.class.getDeclaredField("refObject"); + field.setAccessible(true); + field.set(gauge, reference); + + assertEquals(7, gauge.getValue(), 0); + assertEquals(0, gauge.getValue(), 0); + } +} diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManager.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManager.java index 147f3375683d0..10c1caf18a1c4 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManager.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManager.java @@ -125,10 +125,14 @@ public void clearAll() { @Override public void close() { - pool.close(); - // we need to release tManagers for AsyncThriftClientFactory - if (pool.getFactory() instanceof AsyncThriftClientFactory) { - ((AsyncThriftClientFactory) pool.getFactory()).close(); + try { + pool.close(); + // we need to release tManagers for AsyncThriftClientFactory + if (pool.getFactory() instanceof AsyncThriftClientFactory) { + ((AsyncThriftClientFactory) pool.getFactory()).close(); + } + } finally { + ClientManagerMetrics.getInstance().unregisterClientManager(pool); } } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManagerMetrics.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManagerMetrics.java index cabc917ddb512..b1bade78d8138 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManagerMetrics.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManagerMetrics.java @@ -29,6 +29,7 @@ import org.apache.commons.pool2.impl.GenericKeyedObjectPool; import java.util.HashMap; +import java.util.Iterator; import java.util.Map; public class ClientManagerMetrics implements IMetricSet { @@ -58,35 +59,50 @@ private ClientManagerMetrics() { // empty constructor } - public void registerClientManager(String poolName, GenericKeyedObjectPool clientPool) { - synchronized (this) { - if (metricService == null) { - poolMap.put(poolName, clientPool); - } else { - if (!poolMap.containsKey(poolName)) { - poolMap.put(poolName, clientPool); - createMetrics(poolName); + public synchronized void registerClientManager( + String poolName, GenericKeyedObjectPool clientPool) { + GenericKeyedObjectPool existingPool = poolMap.get(poolName); + if (metricService != null && existingPool != null && !existingPool.isClosed()) { + return; + } + if (metricService != null && existingPool != null) { + removeMetrics(metricService, poolName); + } + poolMap.put(poolName, clientPool); + if (metricService != null) { + createMetrics(poolName, clientPool); + } + } + + public synchronized void unregisterClientManager(GenericKeyedObjectPool clientPool) { + Iterator>> iterator = + poolMap.entrySet().iterator(); + while (iterator.hasNext()) { + Map.Entry> entry = iterator.next(); + // A pool with the same name may have been registered while the old pool was closing. + if (entry.getValue() == clientPool) { + iterator.remove(); + if (metricService != null) { + removeMetrics(metricService, entry.getKey()); } } } } @Override - public void bindTo(AbstractMetricService metricService) { + public synchronized void bindTo(AbstractMetricService metricService) { this.metricService = metricService; - synchronized (this) { - for (String poolName : poolMap.keySet()) { - createMetrics(poolName); - } + for (Map.Entry> entry : poolMap.entrySet()) { + createMetrics(entry.getKey(), entry.getValue()); } } - private void createMetrics(String poolName) { + private void createMetrics(String poolName, GenericKeyedObjectPool clientPool) { metricService.createAutoGauge( Metric.CLIENT_MANAGER.toString(), MetricLevel.IMPORTANT, - poolMap, - map -> poolMap.get(poolName).getNumActive(), + clientPool, + GenericKeyedObjectPool::getNumActive, Tag.NAME.toString(), CLIENT_MANAGER_NUM_ACTIVE, Tag.TYPE.toString(), @@ -94,8 +110,8 @@ private void createMetrics(String poolName) { metricService.createAutoGauge( Metric.CLIENT_MANAGER.toString(), MetricLevel.IMPORTANT, - poolMap, - map -> poolMap.get(poolName).getNumIdle(), + clientPool, + GenericKeyedObjectPool::getNumIdle, Tag.NAME.toString(), CLIENT_MANAGER_NUM_IDLE, Tag.TYPE.toString(), @@ -103,8 +119,8 @@ private void createMetrics(String poolName) { metricService.createAutoGauge( Metric.CLIENT_MANAGER.toString(), MetricLevel.IMPORTANT, - poolMap, - map -> poolMap.get(poolName).getBorrowedCount(), + clientPool, + GenericKeyedObjectPool::getBorrowedCount, Tag.NAME.toString(), CLIENT_MANAGER_BORROWED_COUNT, Tag.TYPE.toString(), @@ -112,8 +128,8 @@ private void createMetrics(String poolName) { metricService.createAutoGauge( Metric.CLIENT_MANAGER.toString(), MetricLevel.IMPORTANT, - poolMap, - map -> poolMap.get(poolName).getCreatedCount(), + clientPool, + GenericKeyedObjectPool::getCreatedCount, Tag.NAME.toString(), CLIENT_MANAGER_CREATED_COUNT, Tag.TYPE.toString(), @@ -121,8 +137,8 @@ private void createMetrics(String poolName) { metricService.createAutoGauge( Metric.CLIENT_MANAGER.toString(), MetricLevel.IMPORTANT, - poolMap, - map -> poolMap.get(poolName).getDestroyedCount(), + clientPool, + GenericKeyedObjectPool::getDestroyedCount, Tag.NAME.toString(), CLIENT_MANAGER_DESTROYED_COUNT, Tag.TYPE.toString(), @@ -130,8 +146,8 @@ private void createMetrics(String poolName) { metricService.createAutoGauge( Metric.CLIENT_MANAGER.toString(), MetricLevel.IMPORTANT, - poolMap, - map -> poolMap.get(poolName).getMeanActiveTimeMillis(), + clientPool, + GenericKeyedObjectPool::getMeanActiveTimeMillis, Tag.NAME.toString(), MEAN_ACTIVE_TIME_MILLIS, Tag.TYPE.toString(), @@ -139,8 +155,8 @@ private void createMetrics(String poolName) { metricService.createAutoGauge( Metric.CLIENT_MANAGER.toString(), MetricLevel.IMPORTANT, - poolMap, - map -> poolMap.get(poolName).getMeanBorrowWaitTimeMillis(), + clientPool, + GenericKeyedObjectPool::getMeanBorrowWaitTimeMillis, Tag.NAME.toString(), MEAN_BORROW_WAIT_TIME_MILLIS, Tag.TYPE.toString(), @@ -148,8 +164,8 @@ private void createMetrics(String poolName) { metricService.createAutoGauge( Metric.CLIENT_MANAGER.toString(), MetricLevel.IMPORTANT, - poolMap, - map -> poolMap.get(poolName).getMeanIdleTimeMillis(), + clientPool, + GenericKeyedObjectPool::getMeanIdleTimeMillis, Tag.NAME.toString(), MEAN_IDLE_TIME_MILLIS, Tag.TYPE.toString(), @@ -157,65 +173,73 @@ private void createMetrics(String poolName) { } @Override - public void unbindFrom(AbstractMetricService metricService) { + public synchronized void unbindFrom(AbstractMetricService metricService) { + if (this.metricService != metricService) { + return; + } + this.metricService = null; + // Keep live pools registered so a metric service restart can bind them again. for (String poolName : poolMap.keySet()) { - metricService.remove( - MetricType.GAUGE, - Metric.CLIENT_MANAGER.toString(), - Tag.NAME.toString(), - CLIENT_MANAGER_NUM_ACTIVE, - Tag.TYPE.toString(), - poolName); - metricService.remove( - MetricType.GAUGE, - Metric.CLIENT_MANAGER.toString(), - Tag.NAME.toString(), - CLIENT_MANAGER_NUM_IDLE, - Tag.TYPE.toString(), - poolName); - metricService.remove( - MetricType.GAUGE, - Metric.CLIENT_MANAGER.toString(), - Tag.NAME.toString(), - CLIENT_MANAGER_BORROWED_COUNT, - Tag.TYPE.toString(), - poolName); - metricService.remove( - MetricType.GAUGE, - Metric.CLIENT_MANAGER.toString(), - Tag.NAME.toString(), - CLIENT_MANAGER_CREATED_COUNT, - Tag.TYPE.toString(), - poolName); - metricService.remove( - MetricType.GAUGE, - Metric.CLIENT_MANAGER.toString(), - Tag.NAME.toString(), - CLIENT_MANAGER_DESTROYED_COUNT, - Tag.TYPE.toString(), - poolName); - metricService.remove( - MetricType.GAUGE, - Metric.CLIENT_MANAGER.toString(), - Tag.NAME.toString(), - MEAN_ACTIVE_TIME_MILLIS, - Tag.TYPE.toString(), - poolName); - metricService.remove( - MetricType.GAUGE, - Metric.CLIENT_MANAGER.toString(), - Tag.NAME.toString(), - MEAN_BORROW_WAIT_TIME_MILLIS, - Tag.TYPE.toString(), - poolName); - metricService.remove( - MetricType.GAUGE, - Metric.CLIENT_MANAGER.toString(), - Tag.NAME.toString(), - MEAN_IDLE_TIME_MILLIS, - Tag.TYPE.toString(), - poolName); + removeMetrics(metricService, poolName); } - poolMap.clear(); + } + + private void removeMetrics(AbstractMetricService metricService, String poolName) { + metricService.remove( + MetricType.AUTO_GAUGE, + Metric.CLIENT_MANAGER.toString(), + Tag.NAME.toString(), + CLIENT_MANAGER_NUM_ACTIVE, + Tag.TYPE.toString(), + poolName); + metricService.remove( + MetricType.AUTO_GAUGE, + Metric.CLIENT_MANAGER.toString(), + Tag.NAME.toString(), + CLIENT_MANAGER_NUM_IDLE, + Tag.TYPE.toString(), + poolName); + metricService.remove( + MetricType.AUTO_GAUGE, + Metric.CLIENT_MANAGER.toString(), + Tag.NAME.toString(), + CLIENT_MANAGER_BORROWED_COUNT, + Tag.TYPE.toString(), + poolName); + metricService.remove( + MetricType.AUTO_GAUGE, + Metric.CLIENT_MANAGER.toString(), + Tag.NAME.toString(), + CLIENT_MANAGER_CREATED_COUNT, + Tag.TYPE.toString(), + poolName); + metricService.remove( + MetricType.AUTO_GAUGE, + Metric.CLIENT_MANAGER.toString(), + Tag.NAME.toString(), + CLIENT_MANAGER_DESTROYED_COUNT, + Tag.TYPE.toString(), + poolName); + metricService.remove( + MetricType.AUTO_GAUGE, + Metric.CLIENT_MANAGER.toString(), + Tag.NAME.toString(), + MEAN_ACTIVE_TIME_MILLIS, + Tag.TYPE.toString(), + poolName); + metricService.remove( + MetricType.AUTO_GAUGE, + Metric.CLIENT_MANAGER.toString(), + Tag.NAME.toString(), + MEAN_BORROW_WAIT_TIME_MILLIS, + Tag.TYPE.toString(), + poolName); + metricService.remove( + MetricType.AUTO_GAUGE, + Metric.CLIENT_MANAGER.toString(), + Tag.NAME.toString(), + MEAN_IDLE_TIME_MILLIS, + Tag.TYPE.toString(), + poolName); } } diff --git a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java new file mode 100644 index 0000000000000..7d6076a75da9d --- /dev/null +++ b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java @@ -0,0 +1,384 @@ +/* + * 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.commons.client; + +import org.apache.iotdb.commons.service.metric.MetricService; +import org.apache.iotdb.commons.service.metric.enums.Metric; +import org.apache.iotdb.commons.service.metric.enums.Tag; +import org.apache.iotdb.metrics.DoNothingMetricService; +import org.apache.iotdb.metrics.config.MetricConfig; +import org.apache.iotdb.metrics.config.MetricConfigDescriptor; +import org.apache.iotdb.metrics.core.IoTDBMetricManager; +import org.apache.iotdb.metrics.core.reporter.IoTDBJmxReporter; +import org.apache.iotdb.metrics.core.type.IoTDBAutoGauge; +import org.apache.iotdb.metrics.type.AutoGauge; +import org.apache.iotdb.metrics.type.IMetric; +import org.apache.iotdb.metrics.utils.MetricInfo; +import org.apache.iotdb.metrics.utils.MetricLevel; +import org.apache.iotdb.metrics.utils.MetricType; + +import org.apache.commons.pool2.BaseKeyedPooledObjectFactory; +import org.apache.commons.pool2.PooledObject; +import org.apache.commons.pool2.impl.DefaultPooledObject; +import org.apache.commons.pool2.impl.GenericKeyedObjectPool; +import org.apache.commons.pool2.impl.GenericKeyedObjectPoolConfig; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import javax.management.MBeanServer; +import javax.management.ObjectName; + +import java.lang.management.ManagementFactory; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; +import java.util.stream.Collectors; + +import static org.awaitility.Awaitility.await; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +public class ClientManagerMetricsTest { + + private final ClientManagerMetrics metrics = ClientManagerMetrics.getInstance(); + private final List> managers = new ArrayList<>(); + private final MetricConfig config = MetricConfigDescriptor.getInstance().getMetricConfig(); + private TestMetricService service; + private MetricLevel originalLevel; + private String originalReporters; + + @Before + public void setUp() { + originalLevel = config.getMetricLevel(); + originalReporters = + config.getMetricReporterList().stream().map(Enum::name).collect(Collectors.joining(",")); + config.setMetricLevel(MetricLevel.IMPORTANT); + config.setMetricReporterList(""); + service = new TestMetricService(); + service.startService(); + service.addMetricSet(metrics); + } + + @After + public void tearDown() { + managers.forEach(ClientManager::close); + service.removeMetricSet(metrics); + service.stopService(); + config.setMetricLevel(originalLevel); + config.setMetricReporterList(originalReporters); + } + + @Test + public void testInFlightGaugesAndRepeatedRebind() throws Exception { + ClientManager manager = createManager("pool"); + Object client = manager.borrowClient("node"); + manager.returnClient("node", client); + manager.borrowClient("node"); + manager.getPool().addObject("node"); + Map inFlight = new HashMap<>(service.getAllMetrics()); + assertEquals(8, inFlight.size()); + + for (int i = 0; i < 3; i++) { + service.stopService(); + assertTrue(service.getAllMetrics().isEmpty()); + assertMetrics(inFlight, "pool", manager.getPool()); + service.startService(); + assertEquals(8, service.getAllMetrics().size()); + assertMetrics(service.getAllMetrics(), "pool", manager.getPool()); + } + + manager.borrowClient("node"); + assertMetrics(inFlight, "pool", manager.getPool()); + } + + @Test + public void testRegistrationWhileUnbound() { + ClientManager first = createManager("first"); + service.removeMetricSet(metrics); + ClientManager second = createManager("second"); + assertTrue(service.getAllMetrics().isEmpty()); + + service.addMetricSet(metrics); + assertEquals(16, service.getAllMetrics().size()); + assertMetrics(service.getAllMetrics(), "first", first.getPool()); + assertMetrics(service.getAllMetrics(), "second", second.getPool()); + } + + @Test + public void testMetricServiceRestartAndLevelChanges() throws Exception { + ClientManager manager = createManager("pool"); + manager.borrowClient("node"); + service.removeMetricSet(metrics); + MetricService metricService = MetricService.getInstance(); + metricService.startService(); + try { + metricService.addMetricSet(metrics); + for (MetricLevel level : + new MetricLevel[] {MetricLevel.ALL, MetricLevel.CORE, MetricLevel.IMPORTANT}) { + config.setMetricLevel(level); + metricService.restartService(); + if (level == MetricLevel.CORE) { + assertTrue(metricService.getAllMetrics().isEmpty()); + } else { + assertEquals(8, metricService.getAllMetrics().size()); + assertMetrics(metricService.getAllMetrics(), "pool", manager.getPool()); + } + } + } finally { + metricService.removeMetricSet(metrics); + metricService.stopService(); + } + } + + @Test + public void testClosedPoolsAreNotRetainedAcrossRebinds() { + for (int i = 0; i < 10; i++) { + ClientManager manager = createManager("pool" + i); + assertEquals(8, service.getAllMetrics().size()); + manager.close(); + assertTrue(service.getAllMetrics().isEmpty()); + service.stopService(); + service.startService(); + assertTrue(service.getAllMetrics().isEmpty()); + } + } + + @Test + public void testCloseWhileUnboundRemovesOnlyClosedPool() { + ClientManager first = createManager("first"); + ClientManager second = createManager("second"); + service.removeMetricSet(metrics); + first.close(); + service.addMetricSet(metrics); + + assertEquals(8, service.getAllMetrics().size()); + assertMetrics(service.getAllMetrics(), "second", second.getPool()); + } + + @Test + public void testClosingSupersededPoolDoesNotRemoveReplacement() throws Exception { + ClientManager first = createManager("pool"); + first.borrowClient("node"); + Map inFlight = new HashMap<>(service.getAllMetrics()); + service.removeMetricSet(metrics); + ClientManager replacement = createManager("pool"); + replacement.borrowClient("node"); + replacement.borrowClient("node"); + service.addMetricSet(metrics); + + assertMetrics(inFlight, "pool", first.getPool()); + first.close(); + assertEquals(8, service.getAllMetrics().size()); + assertMetrics(service.getAllMetrics(), "pool", replacement.getPool()); + assertMetrics(inFlight, "pool", first.getPool()); + } + + @Test + public void testClosingUnregisteredDuplicateDoesNotRemoveOriginal() throws Exception { + ClientManager first = createManager("pool"); + first.borrowClient("node"); + ClientManager duplicate = createManager("pool"); + duplicate.close(); + + assertEquals(8, service.getAllMetrics().size()); + assertMetrics(service.getAllMetrics(), "pool", first.getPool()); + } + + @Test + public void testClosedPoolCanBeReplacedBeforeUnregistration() throws Exception { + IoTDBJmxReporter reporter = IoTDBJmxReporter.getInstance(); + service.getMetricManager().setBindJmxReporter(reporter); + try { + ClientManager first = createManager("pool"); + first.borrowClient("node"); + ObjectName objectName = + ((IoTDBAutoGauge) + service.getAutoGauge( + Metric.CLIENT_MANAGER.toString(), + MetricLevel.IMPORTANT, + Tag.NAME.toString(), + "client_manager_num_active", + Tag.TYPE.toString(), + "pool")) + .objectName(); + MBeanServer mBeanServer = ManagementFactory.getPlatformMBeanServer(); + assertEquals(1, (double) mBeanServer.getAttribute(objectName, "Value"), 0); + // The pool is closed, but its owner's close has not reached metric unregistration yet. + first.getPool().close(); + ClientManager replacement = createManager("pool"); + replacement.borrowClient("node"); + replacement.borrowClient("node"); + first.close(); + + assertEquals(8, service.getAllMetrics().size()); + assertMetrics(service.getAllMetrics(), "pool", replacement.getPool()); + assertEquals(2, (double) mBeanServer.getAttribute(objectName, "Value"), 0); + } finally { + service.getMetricManager().setBindJmxReporter(null); + reporter.stop(); + } + } + + @Test + public void testRegistrationAndSamplingDuringUnbind() throws Exception { + ClientManager first = createManager("first"); + first.borrowClient("node"); + createManager("second"); + AutoGauge inFlight = + service.getAutoGauge( + Metric.CLIENT_MANAGER.toString(), + MetricLevel.IMPORTANT, + Tag.NAME.toString(), + "client_manager_num_active", + Tag.TYPE.toString(), + "first"); + CountDownLatch removing = new CountDownLatch(1); + CountDownLatch resume = new CountDownLatch(1); + AtomicBoolean pauseOnce = new AtomicBoolean(true); + service.beforeRemove = + () -> { + if (pauseOnce.compareAndSet(true, false)) { + removing.countDown(); + try { + assertTrue(resume.await(10, TimeUnit.SECONDS)); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AssertionError(e); + } + } + }; + ExecutorService executor = Executors.newFixedThreadPool(3); + try { + Future stop = executor.submit(service::stopService); + assertTrue(removing.await(5, TimeUnit.SECONDS)); + AtomicReference registrationThread = new AtomicReference<>(); + CountDownLatch registering = new CountDownLatch(1); + Future> registration = + executor.submit( + () -> { + registrationThread.set(Thread.currentThread()); + registering.countDown(); + return createManager("third"); + }); + assertTrue(registering.await(5, TimeUnit.SECONDS)); + await() + .atMost(5, TimeUnit.SECONDS) + .until( + () -> + registration.isDone() + || registrationThread.get().getState() == Thread.State.BLOCKED); + assertFalse(registration.isDone()); + // Sampling must not acquire the lifecycle lock held by the paused unbind. + assertEquals(1, executor.submit(inFlight::getValue).get(5, TimeUnit.SECONDS), 0); + resume.countDown(); + stop.get(5, TimeUnit.SECONDS); + registration.get(5, TimeUnit.SECONDS); + assertTrue(service.getAllMetrics().isEmpty()); + + service.startService(); + assertEquals(24, service.getAllMetrics().size()); + assertMetrics(service.getAllMetrics(), "first", first.getPool()); + } finally { + resume.countDown(); + executor.shutdownNow(); + assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + private ClientManager createManager(String name) { + ClientManager manager = + new ClientManager<>( + owner -> { + GenericKeyedObjectPoolConfig poolConfig = + new GenericKeyedObjectPoolConfig<>(); + poolConfig.setJmxEnabled(false); + GenericKeyedObjectPool pool = + new GenericKeyedObjectPool<>( + new BaseKeyedPooledObjectFactory() { + @Override + public Object create(String key) { + return new Object(); + } + + @Override + public PooledObject wrap(Object value) { + return new DefaultPooledObject<>(value); + } + }, + poolConfig); + metrics.registerClientManager(name, pool); + return pool; + }); + managers.add(manager); + return manager; + } + + private void assertMetrics( + Map actual, String poolName, GenericKeyedObjectPool pool) { + Map expected = new HashMap<>(); + expected.put("client_manager_num_active", (double) pool.getNumActive()); + expected.put("client_manager_num_idle", (double) pool.getNumIdle()); + expected.put("client_manager_borrowed_count", (double) pool.getBorrowedCount()); + expected.put("client_manager_created_count", (double) pool.getCreatedCount()); + expected.put("client_manager_destroyed_count", (double) pool.getDestroyedCount()); + expected.put("client_manager_mean_active_time", (double) pool.getMeanActiveTimeMillis()); + expected.put( + "client_manager_mean_borrow_wait_time", (double) pool.getMeanBorrowWaitTimeMillis()); + expected.put("client_manager_mean_idle_time", (double) pool.getMeanIdleTimeMillis()); + expected.forEach( + (name, value) -> { + MetricInfo key = + new MetricInfo( + MetricType.AUTO_GAUGE, + Metric.CLIENT_MANAGER.toString(), + Tag.NAME.toString(), + name, + Tag.TYPE.toString(), + poolName); + assertTrue(actual.get(key) instanceof AutoGauge); + assertEquals(value, ((AutoGauge) actual.get(key)).getValue(), 0); + }); + } + + private static class TestMetricService extends DoNothingMetricService { + private Runnable beforeRemove = () -> {}; + + @Override + protected void loadManager() { + metricManager = IoTDBMetricManager.getInstance(); + } + + @Override + public void remove(MetricType type, String metric, String... tags) { + beforeRemove.run(); + super.remove(type, metric, tags); + } + } +} From 7e249c0826f32d4e002174ad82d03d733f653eb9 Mon Sep 17 00:00:00 2001 From: JackieTien97 Date: Sun, 20 Sep 2026 18:32:10 +0800 Subject: [PATCH 2/4] Isolate JMX metrics by labels and registration owner --- .../core/reporter/IoTDBJmxReporter.java | 102 +++++++---- .../core/utils/IoTDBMetricObjNameFactory.java | 38 +--- .../metrics/core/utils/ObjectNameFactory.java | 45 +++++ .../core/reporter/IoTDBJmxReporterTest.java | 169 ++++++++++++++++++ .../utils/IoTDBMetricObjNameFactoryTest.java | 123 +++++++++++++ .../client/ClientManagerMetricsTest.java | 46 +++++ 6 files changed, 454 insertions(+), 69 deletions(-) create mode 100644 iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporterTest.java create mode 100644 iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/utils/IoTDBMetricObjNameFactoryTest.java diff --git a/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporter.java b/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporter.java index 372afdb36f65c..4ee5bc272cda7 100644 --- a/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporter.java +++ b/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporter.java @@ -40,8 +40,9 @@ import javax.management.ObjectName; import java.lang.management.ManagementFactory; +import java.util.HashMap; +import java.util.Iterator; import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; public class IoTDBJmxReporter implements JmxReporter { private static final Logger LOGGER = LoggerFactory.getLogger(IoTDBJmxReporter.class); @@ -55,37 +56,36 @@ public class IoTDBJmxReporter implements JmxReporter { /** The objectNameFactory used to create objectName for metrics */ private final ObjectNameFactory objectNameFactory; - /** The map that stores all registered metrics */ - private final Map registered; + /** Registrations owned by this reporter, guarded by the map's monitor. */ + private final Map registered; /** The JMX MBeanServer */ private final MBeanServer mBeanServer; - private void registerMBean(Object mBean, ObjectName objectName) throws JMException { - if (!mBeanServer.isRegistered(objectName)) { - ObjectInstance objectInstance = mBeanServer.registerMBean(mBean, objectName); - if (objectInstance != null) { - // the websphere mbeanserver rewrites the objectname to include - // cell, node & server info - // make sure we capture the new objectName for unregistration - registered.put(objectName, objectInstance.getObjectName()); - } else { - registered.put(objectName, objectName); + private void registerMBean(IMetric metric, ObjectName objectName) throws JMException { + Registration previous = registered.get(objectName); + if (previous != null) { + if (previous.metric == metric && mBeanServer.isRegistered(previous.actualName)) { + return; } + unregisterMBean(previous); + registered.remove(objectName); + } + if (!mBeanServer.isRegistered(objectName)) { + ObjectInstance objectInstance = mBeanServer.registerMBean(metric, objectName); + // Some MBean servers rewrite ObjectNames. Keep the actual name together with its owner. + registered.put( + objectName, + new Registration( + metric, objectInstance == null ? objectName : objectInstance.getObjectName())); } } - private void unregisterMBean(ObjectName originalObjectName) - throws InstanceNotFoundException, MBeanRegistrationException { - ObjectName storedObjectName = registered.remove(originalObjectName); - if (storedObjectName != null) { - if (mBeanServer.isRegistered(storedObjectName)) { - mBeanServer.unregisterMBean(storedObjectName); - } - } else { - if (mBeanServer.isRegistered(originalObjectName)) { - mBeanServer.unregisterMBean(originalObjectName); - } + private void unregisterMBean(Registration registration) throws MBeanRegistrationException { + try { + mBeanServer.unregisterMBean(registration.actualName); + } catch (InstanceNotFoundException ignored) { + // An externally removed MBean is already unregistered. } } @@ -94,8 +94,10 @@ public void registerMetric(IMetric metric, MetricInfo metricInfo) { String metricName = metric.getClass().getSimpleName(); try { final ObjectName objectName = createName(metricName, metricInfo); - metric.setObjectName(objectName); - registerMBean(metric, objectName); + synchronized (registered) { + metric.setObjectName(objectName); + registerMBean(metric, objectName); + } } catch (Exception e) { LOGGER.warn(MetricsCoreMessages.JMX_REGISTER_FAILED + metricName, e); } @@ -103,12 +105,20 @@ public void registerMetric(IMetric metric, MetricInfo metricInfo) { @Override public void unregisterMetric(IMetric metric, MetricInfo metricInfo) { + if (metric == null) { + return; + } String metricName = metric.getClass().getSimpleName(); try { final ObjectName objectName = createName(metricName, metricInfo); - unregisterMBean(objectName); - } catch (InstanceNotFoundException e) { - LOGGER.debug(MetricsCoreMessages.JMX_UNREGISTER_FAILED, e); + synchronized (registered) { + Registration registration = registered.get(objectName); + // A delayed callback for an old metric must not delete its replacement. + if (registration != null && registration.metric == metric) { + unregisterMBean(registration); + registered.remove(objectName); + } + } } catch (MBeanRegistrationException e) { LOGGER.warn(MetricsCoreMessages.JMX_UNREGISTER_FAILED, e); } @@ -116,31 +126,37 @@ public void unregisterMetric(IMetric metric, MetricInfo metricInfo) { private ObjectName createName(String type, MetricInfo metricInfo) { String name = metricInfo.getName(); - return objectNameFactory.createName(type, DOMAIN, name); + return objectNameFactory.createName(type, DOMAIN, name, metricInfo.getTags()); } - void unregisterAll() throws InstanceNotFoundException, MBeanRegistrationException { - for (ObjectName name : registered.keySet()) { - unregisterMBean(name); + void unregisterAll() throws MBeanRegistrationException { + synchronized (registered) { + Iterator iterator = registered.values().iterator(); + while (iterator.hasNext()) { + unregisterMBean(iterator.next()); + iterator.remove(); + } } - // clear registered - registered.clear(); } - private IoTDBJmxReporter( + IoTDBJmxReporter( AbstractMetricManager metricManager, MBeanServer mBeanServer, ObjectNameFactory objectNameFactory) { this.metricManager = metricManager; this.mBeanServer = mBeanServer; this.objectNameFactory = objectNameFactory; - this.registered = new ConcurrentHashMap<>(); + this.registered = new HashMap<>(); } @Override public boolean start() { try { - if (!registered.isEmpty()) { + boolean alreadyRegistered; + synchronized (registered) { + alreadyRegistered = !registered.isEmpty(); + } + if (alreadyRegistered) { LOGGER.warn(MetricsCoreMessages.JMX_REPORTER_ALREADY_START); return false; } @@ -171,6 +187,16 @@ public ReporterType getReporterType() { return ReporterType.JMX; } + private static class Registration { + private final IMetric metric; + private final ObjectName actualName; + + private Registration(IMetric metric, ObjectName actualName) { + this.metric = metric; + this.actualName = actualName; + } + } + private static class IoTDBJmxReporterHolder { private static final IoTDBJmxReporter INSTANCE = new IoTDBJmxReporter( diff --git a/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/utils/IoTDBMetricObjNameFactory.java b/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/utils/IoTDBMetricObjNameFactory.java index bcc5a31dcbffc..b2c3b914d4acf 100644 --- a/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/utils/IoTDBMetricObjNameFactory.java +++ b/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/utils/IoTDBMetricObjNameFactory.java @@ -30,7 +30,7 @@ import java.util.Hashtable; public class IoTDBMetricObjNameFactory implements ObjectNameFactory { - private static final char[] QUOTABLE_CHARS = new char[] {',', '=', ':', '"'}; + private static final char[] QUOTABLE_CHARS = new char[] {',', '=', ':', '"', '\n', '*', '?'}; private static final Logger LOGGER = LoggerFactory.getLogger(IoTDBMetricObjNameFactory.class); private IoTDBMetricObjNameFactory() { @@ -40,38 +40,14 @@ private IoTDBMetricObjNameFactory() { @Override public ObjectName createName(String type, String domain, String name) { try { - ObjectName objectName; Hashtable properties = new Hashtable<>(); - - properties.put("name", name); - properties.put("type", type); - objectName = new ObjectName(domain, properties); - - /* - * The only way we can find out if we need to quote the properties is by - * checking an ObjectName that we've constructed. - */ - if (objectName.isDomainPattern()) { - domain = ObjectName.quote(domain); - } - if (objectName.isPropertyValuePattern("name") - || shouldQuote(objectName.getKeyProperty("name"))) { - properties.put("name", ObjectName.quote(name)); - } - if (objectName.isPropertyValuePattern("type") - || shouldQuote(objectName.getKeyProperty("type"))) { - properties.put("type", ObjectName.quote(type)); - } - objectName = new ObjectName(domain, properties); - - return objectName; + // Quote before constructing the name; falling back to a name-only MBean loses its type. + properties.put("name", shouldQuote(name) ? ObjectName.quote(name) : name); + properties.put("type", shouldQuote(type) ? ObjectName.quote(type) : type); + return new ObjectName(domain, properties); } catch (MalformedObjectNameException e) { - try { - return new ObjectName(domain, "name", ObjectName.quote(name)); - } catch (MalformedObjectNameException e1) { - LOGGER.warn(MetricsCoreMessages.JMX_UNABLE_TO_REGISTER, type, name, e1); - throw new RuntimeException(e1); - } + LOGGER.warn(MetricsCoreMessages.JMX_UNABLE_TO_REGISTER, type, name, e); + throw new IllegalArgumentException(e); } } diff --git a/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/utils/ObjectNameFactory.java b/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/utils/ObjectNameFactory.java index 1c019781300f5..96a3bb2eb9c89 100644 --- a/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/utils/ObjectNameFactory.java +++ b/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/utils/ObjectNameFactory.java @@ -19,8 +19,12 @@ package org.apache.iotdb.metrics.core.utils; +import javax.management.MalformedObjectNameException; import javax.management.ObjectName; +import java.util.Hashtable; +import java.util.Map; + public interface ObjectNameFactory { /** * Create objectName for a certain metric. @@ -31,4 +35,45 @@ public interface ObjectNameFactory { * @return metric's objectName */ public ObjectName createName(String type, String domain, String name); + + /** + * Include all metric tags in the MBean identity. Tag keys use a separate namespace and reversible + * escaping so they cannot overwrite the metric name/type or collide after sanitization. Tag + * values are quoted to preserve literal wildcard characters. Untagged metrics keep their existing + * names. + */ + default ObjectName createName(String type, String domain, String name, Map tags) { + ObjectName base = createName(type, domain, name); + if (tags.isEmpty()) { + return base; + } + Hashtable properties = base.getKeyPropertyList(); + tags.forEach((key, value) -> properties.put(encodeTagKey(key), ObjectName.quote(value))); + try { + return new ObjectName(base.getDomain(), properties); + } catch (MalformedObjectNameException e) { + throw new IllegalArgumentException(e); + } + } + + private static String encodeTagKey(String key) { + StringBuilder encoded = new StringBuilder("tag."); + for (int i = 0; i < key.length(); i++) { + char character = key.charAt(i); + if ((character >= 'a' && character <= 'z') + || (character >= 'A' && character <= 'Z') + || (character >= '0' && character <= '9') + || character == '_' + || character == '-' + || character == '.') { + encoded.append(character); + } else { + encoded.append('%'); + for (int shift = 12; shift >= 0; shift -= 4) { + encoded.append(Character.forDigit((character >> shift) & 0xf, 16)); + } + } + } + return encoded.toString(); + } } diff --git a/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporterTest.java b/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporterTest.java new file mode 100644 index 0000000000000..d247ef0b4ebe6 --- /dev/null +++ b/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporterTest.java @@ -0,0 +1,169 @@ +/* + * 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.core.reporter; + +import org.apache.iotdb.metrics.core.type.IoTDBAutoGauge; +import org.apache.iotdb.metrics.core.type.IoTDBCounter; +import org.apache.iotdb.metrics.core.utils.IoTDBMetricObjNameFactory; +import org.apache.iotdb.metrics.impl.DoNothingMetricManager; +import org.apache.iotdb.metrics.type.AutoGauge; +import org.apache.iotdb.metrics.type.Counter; +import org.apache.iotdb.metrics.utils.MetricInfo; +import org.apache.iotdb.metrics.utils.MetricLevel; +import org.apache.iotdb.metrics.utils.MetricType; + +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import javax.management.MBeanServer; +import javax.management.MBeanServerFactory; +import javax.management.ObjectName; + +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.ToDoubleFunction; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotEquals; +import static org.junit.Assert.assertTrue; + +public class IoTDBJmxReporterTest { + private TestMetricManager manager; + private MBeanServer server; + private IoTDBJmxReporter reporter; + + @Before + public void setUp() { + manager = new TestMetricManager(); + server = MBeanServerFactory.newMBeanServer(); + reporter = new IoTDBJmxReporter(manager, server, IoTDBMetricObjNameFactory.getInstance()); + manager.setBindJmxReporter(reporter); + assertTrue(reporter.start()); + } + + @After + public void tearDown() { + assertTrue(reporter.stop()); + } + + @Test + public void testTaggedMetricsAreIndependent() throws Exception { + AtomicInteger firstValue = new AtomicInteger(1); + AtomicInteger secondValue = new AtomicInteger(2); + IoTDBAutoGauge first = gauge(firstValue, "first"); + IoTDBAutoGauge second = gauge(secondValue, "second"); + assertNotEquals(first.objectName(), second.objectName()); + assertEquals(1, (double) server.getAttribute(first.objectName(), "Value"), 0); + assertEquals(2, (double) server.getAttribute(second.objectName(), "Value"), 0); + + manager.remove(MetricType.AUTO_GAUGE, "client_manager", "name", "num_active", "type", "second"); + assertFalse(server.isRegistered(second.objectName())); + assertEquals(1, (double) server.getAttribute(first.objectName(), "Value"), 0); + } + + @Test + public void testDelayedUnregisterPreservesReplacement() throws Exception { + AtomicInteger firstValue = new AtomicInteger(1); + AtomicInteger nextValue = new AtomicInteger(2); + IoTDBAutoGauge first = gauge(firstValue, "pool"); + IoTDBAutoGauge next = gauge(nextValue, "pool"); + assertEquals(first.objectName(), next.objectName()); + reporter.unregisterMetric(first, info("pool")); + assertEquals(2, (double) server.getAttribute(next.objectName(), "Value"), 0); + + assertTrue(reporter.stop()); + assertTrue(reporter.start()); + reporter.unregisterMetric(first, info("pool")); + assertEquals(2, (double) server.getAttribute(next.objectName(), "Value"), 0); + } + + @Test + public void testUnregisterDoesNotRemoveUnownedMBean() throws Exception { + AtomicInteger value = new AtomicInteger(3); + IoTDBAutoGauge external = new IoTDBAutoGauge<>(value, AtomicInteger::get); + ObjectName name = + IoTDBMetricObjNameFactory.getInstance() + .createName( + "IoTDBAutoGauge", + "org.apache.iotdb.metrics", + "client_manager", + info("pool").getTags()); + server.registerMBean(external, name); + reporter.unregisterMetric(external, info("pool")); + assertEquals(3, (double) server.getAttribute(name, "Value"), 0); + assertTrue(reporter.stop()); + assertTrue(server.isRegistered(name)); + } + + @Test + public void testRepeatedAndMissingUnregistration() { + AtomicInteger value = new AtomicInteger(1); + IoTDBAutoGauge gauge = gauge(value, "pool"); + manager.remove(MetricType.AUTO_GAUGE, "client_manager", "name", "num_active", "type", "pool"); + reporter.unregisterMetric(gauge, info("pool")); + reporter.unregisterMetric(null, info("pool")); + assertFalse(server.isRegistered(gauge.objectName())); + } + + @Test + public void testCounterRegistration() throws Exception { + IoTDBCounter counter = + (IoTDBCounter) + manager.getOrCreateCounter("requests", MetricLevel.IMPORTANT, "name", "client"); + counter.inc(7); + assertEquals(7L, server.getAttribute(counter.objectName(), "Count")); + } + + private IoTDBAutoGauge gauge(AtomicInteger value, String pool) { + return (IoTDBAutoGauge) + manager.createAutoGauge( + "client_manager", + MetricLevel.IMPORTANT, + value, + AtomicInteger::get, + "name", + "num_active", + "type", + pool); + } + + private MetricInfo info(String pool) { + return new MetricInfo( + MetricType.AUTO_GAUGE, "client_manager", "name", "num_active", "type", pool); + } + + private static class TestMetricManager extends DoNothingMetricManager { + @Override + public boolean isEnableMetricInGivenLevel(MetricLevel level) { + return true; + } + + @Override + public AutoGauge createAutoGauge(T object, ToDoubleFunction mapper) { + return new IoTDBAutoGauge<>(object, mapper); + } + + @Override + public Counter createCounter() { + return new IoTDBCounter(); + } + } +} diff --git a/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/utils/IoTDBMetricObjNameFactoryTest.java b/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/utils/IoTDBMetricObjNameFactoryTest.java new file mode 100644 index 0000000000000..9bd8f5195b973 --- /dev/null +++ b/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/utils/IoTDBMetricObjNameFactoryTest.java @@ -0,0 +1,123 @@ +/* + * 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.core.utils; + +import org.junit.Test; + +import javax.management.ObjectName; + +import java.util.Collections; +import java.util.HashSet; +import java.util.LinkedHashMap; +import java.util.Map; +import java.util.Set; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotEquals; + +public class IoTDBMetricObjNameFactoryTest { + private final ObjectNameFactory factory = IoTDBMetricObjNameFactory.getInstance(); + + @Test + public void testUntaggedNamesStayCompatible() throws Exception { + ObjectName expected = new ObjectName("org.apache.iotdb.metrics:name=plain,type=IoTDBAutoGauge"); + assertEquals( + expected, factory.createName("IoTDBAutoGauge", "org.apache.iotdb.metrics", "plain")); + assertEquals( + expected, + factory.createName( + "IoTDBAutoGauge", "org.apache.iotdb.metrics", "plain", Collections.emptyMap())); + } + + @Test + public void testTagsCannotOverwriteMetricNameAndType() { + ObjectName name = + factory.createName( + "IoTDBAutoGauge", + "org.apache.iotdb.metrics", + "client_manager", + Map.of("name", "num_active", "type", "first", "tag.name", "nested")); + assertEquals("client_manager", name.getKeyProperty("name")); + assertEquals("IoTDBAutoGauge", name.getKeyProperty("type")); + assertEquals("num_active", ObjectName.unquote(name.getKeyProperty("tag.name"))); + assertEquals("first", ObjectName.unquote(name.getKeyProperty("tag.type"))); + assertEquals("nested", ObjectName.unquote(name.getKeyProperty("tag.tag.name"))); + assertEquals(5, name.getKeyPropertyList().size()); + } + + @Test + public void testTagOrderDoesNotChangeIdentity() { + Map first = new LinkedHashMap<>(); + first.put("name", "num_active"); + first.put("type", "pool"); + Map second = new LinkedHashMap<>(); + second.put("type", "pool"); + second.put("name", "num_active"); + ObjectName a = + factory.createName("IoTDBAutoGauge", "org.apache.iotdb.metrics", "client_manager", first); + ObjectName b = + factory.createName("IoTDBAutoGauge", "org.apache.iotdb.metrics", "client_manager", second); + assertEquals(a, b); + assertEquals(a.getCanonicalName(), b.getCanonicalName()); + assertEquals(2, first.size()); + second.put("type", "other"); + assertNotEquals( + a, + factory.createName("IoTDBAutoGauge", "org.apache.iotdb.metrics", "client_manager", second)); + } + + @Test + public void testEscapingIsLiteralAndCollisionFree() { + String value = "value,=:*?\"\\\n"; + String[] keys = {"", "a:b", "a,b", "a=b", "a?b", "a*b", "a\nb", "a%b", "a.b", "a%003ab", "标签"}; + Set names = new HashSet<>(); + Map tags = new LinkedHashMap<>(); + for (String key : keys) { + ObjectName name = + factory.createName( + "IoTDBAutoGauge", "org.apache.iotdb.metrics", "client_manager", Map.of(key, value)); + assertFalse(name.isPattern()); + names.add(name); + tags.put(key, value); + } + assertEquals(keys.length, names.size()); + ObjectName all = + factory.createName("IoTDBAutoGauge", "org.apache.iotdb.metrics", "client_manager", tags); + assertEquals(keys.length + 2, all.getKeyPropertyList().size()); + all.getKeyPropertyList() + .forEach( + (key, actual) -> { + if (key.startsWith("tag.")) { + assertEquals(value, ObjectName.unquote(actual)); + } + }); + } + + @Test + public void testSpecialMetricNamesRetainType() { + ObjectName name = + factory.createName( + "Gauge:*", "org.apache.iotdb.metrics", "metric,=:\"\n?", Map.of("type", "pool")); + assertEquals("Gauge:*", ObjectName.unquote(name.getKeyProperty("type"))); + assertEquals("metric,=:\"\n?", ObjectName.unquote(name.getKeyProperty("name"))); + assertFalse(name.isPattern()); + } +} diff --git a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java index 7d6076a75da9d..3c6f223272bb0 100644 --- a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java +++ b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java @@ -214,6 +214,7 @@ public void testClosingUnregisteredDuplicateDoesNotRemoveOriginal() throws Excep @Test public void testClosedPoolCanBeReplacedBeforeUnregistration() throws Exception { IoTDBJmxReporter reporter = IoTDBJmxReporter.getInstance(); + assertTrue(reporter.start()); service.getMetricManager().setBindJmxReporter(reporter); try { ClientManager first = createManager("pool"); @@ -246,6 +247,51 @@ public void testClosedPoolCanBeReplacedBeforeUnregistration() throws Exception { } } + @Test + public void testClosingAnotherPoolPreservesAllLivePoolMBeans() throws Exception { + IoTDBJmxReporter reporter = IoTDBJmxReporter.getInstance(); + assertTrue(reporter.start()); + service.getMetricManager().setBindJmxReporter(reporter); + MBeanServer server = ManagementFactory.getPlatformMBeanServer(); + ObjectName pattern = + new ObjectName("org.apache.iotdb.metrics:name=client_manager,type=IoTDBAutoGauge,*"); + try { + ClientManager first = createManager("first"); + first.borrowClient("node"); + ClientManager second = createManager("second"); + second.borrowClient("node"); + second.borrowClient("node"); + assertEquals(16, server.queryNames(pattern, null).size()); + assertJmxMetrics(server); + + second.close(); + assertEquals(8, server.queryNames(pattern, null).size()); + assertMetrics(service.getAllMetrics(), "first", first.getPool()); + assertJmxMetrics(server); + + ClientManager replacement = createManager("second"); + replacement.borrowClient("node"); + second.close(); + assertEquals(16, server.queryNames(pattern, null).size()); + assertJmxMetrics(server); + replacement.close(); + assertEquals(8, server.queryNames(pattern, null).size()); + assertJmxMetrics(server); + first.close(); + assertTrue(server.queryNames(pattern, null).isEmpty()); + } finally { + service.getMetricManager().setBindJmxReporter(null); + reporter.stop(); + } + } + + private void assertJmxMetrics(MBeanServer server) throws Exception { + for (IMetric metric : service.getAllMetrics().values()) { + IoTDBAutoGauge gauge = (IoTDBAutoGauge) metric; + assertEquals(gauge.getValue(), (double) server.getAttribute(gauge.objectName(), "Value"), 0); + } + } + @Test public void testRegistrationAndSamplingDuringUnbind() throws Exception { ClientManager first = createManager("first"); From 414d7bae82c7deeeb201d2ea02c5ae9dcc45d9b6 Mon Sep 17 00:00:00 2001 From: JackieTien97 Date: Sun, 20 Sep 2026 18:32:10 +0800 Subject: [PATCH 3/4] Serialize metric registry lifecycle and isolate close cleanup failures --- .../core/reporter/IoTDBJmxReporter.java | 19 +- .../core/MetricManagerLifecycleTest.java | 339 ++++++++++++++++++ .../core/reporter/IoTDBJmxReporterTest.java | 36 ++ .../iotdb/metrics/AbstractMetricManager.java | 209 +++++++---- .../iotdb/commons/i18n/ClientMessages.java | 3 + .../iotdb/commons/i18n/ClientMessages.java | 3 + .../iotdb/commons/client/ClientManager.java | 9 +- .../client/ClientManagerMetricsTest.java | 152 +++++++- 8 files changed, 689 insertions(+), 81 deletions(-) create mode 100644 iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/MetricManagerLifecycleTest.java diff --git a/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporter.java b/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporter.java index 4ee5bc272cda7..cb750ac64a287 100644 --- a/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporter.java +++ b/iotdb-core/metrics/core/src/main/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporter.java @@ -59,6 +59,8 @@ public class IoTDBJmxReporter implements JmxReporter { /** Registrations owned by this reporter, guarded by the map's monitor. */ private final Map registered; + private boolean started; + /** The JMX MBeanServer */ private final MBeanServer mBeanServer; @@ -95,6 +97,10 @@ public void registerMetric(IMetric metric, MetricInfo metricInfo) { try { final ObjectName objectName = createName(metricName, metricInfo); synchronized (registered) { + // Ignore callbacks from a stopped reporter or a superseded registry entry. + if (!started || metricManager.getAllMetrics().get(metricInfo) != metric) { + return; + } metric.setObjectName(objectName); registerMBean(metric, objectName); } @@ -152,17 +158,19 @@ void unregisterAll() throws MBeanRegistrationException { @Override public boolean start() { try { - boolean alreadyRegistered; + boolean alreadyStarted; synchronized (registered) { - alreadyRegistered = !registered.isEmpty(); + alreadyStarted = started; + started = true; } - if (alreadyRegistered) { + if (alreadyStarted) { LOGGER.warn(MetricsCoreMessages.JMX_REPORTER_ALREADY_START); return false; } // register all existed metrics into JmxReporter metricManager.getAllMetrics().forEach((key, value) -> registerMetric(value, key)); } catch (Exception e) { + stop(); LOGGER.warn(MetricsCoreMessages.JMX_REPORTER_START_FAILED, e); return false; } @@ -173,7 +181,10 @@ public boolean start() { @Override public boolean stop() { try { - unregisterAll(); + synchronized (registered) { + started = false; + unregisterAll(); + } } catch (Exception e) { LOGGER.warn(MetricsCoreMessages.JMX_REPORTER_STOP_FAILED, e); return false; diff --git a/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/MetricManagerLifecycleTest.java b/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/MetricManagerLifecycleTest.java new file mode 100644 index 0000000000000..c80478b1e01fc --- /dev/null +++ b/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/MetricManagerLifecycleTest.java @@ -0,0 +1,339 @@ +/* + * 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.core; + +import org.apache.iotdb.metrics.core.type.IoTDBAutoGauge; +import org.apache.iotdb.metrics.core.type.IoTDBCounter; +import org.apache.iotdb.metrics.impl.DoNothingMetricManager; +import org.apache.iotdb.metrics.reporter.JmxReporter; +import org.apache.iotdb.metrics.type.AutoGauge; +import org.apache.iotdb.metrics.type.Counter; +import org.apache.iotdb.metrics.type.IMetric; +import org.apache.iotdb.metrics.utils.MetricInfo; +import org.apache.iotdb.metrics.utils.MetricLevel; +import org.apache.iotdb.metrics.utils.MetricType; +import org.apache.iotdb.metrics.utils.ReporterType; + +import org.junit.Test; + +import java.lang.management.ManagementFactory; +import java.lang.management.ThreadInfo; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; +import java.util.concurrent.locks.LockSupport; +import java.util.function.ToDoubleFunction; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; + +public class MetricManagerLifecycleTest { + @Test + public void testResetWaitsForRemovalButCachedMetricsRemainAccessible() throws Exception { + TestMetricManager manager = new TestMetricManager(); + RecordingReporter reporter = new RecordingReporter(manager); + manager.setBindJmxReporter(reporter); + AtomicInteger value = new AtomicInteger(3); + AutoGauge closing = + manager.createAutoGauge("closing", MetricLevel.IMPORTANT, value, AtomicInteger::get); + Object[] cached = cachedMetrics(manager); + CountDownLatch removed = new CountDownLatch(1); + CountDownLatch resume = new CountDownLatch(1); + Map original = + new ConcurrentHashMap(manager.getAllMetrics()) { + @Override + public IMetric remove(Object key) { + IMetric metric = super.remove(key); + if (metric != null && ((MetricInfo) key).getName().equals("closing")) { + removed.countDown(); + await(resume); + } + return metric; + } + }; + manager.useRegistry(original); + ExecutorService executor = Executors.newFixedThreadPool(3); + AtomicReference removingThread = new AtomicReference<>(); + AtomicReference resettingThread = new AtomicReference<>(); + CountDownLatch resetting = new CountDownLatch(1); + try { + Future removal = + executor.submit( + () -> { + removingThread.set(Thread.currentThread()); + manager.remove(MetricType.AUTO_GAUGE, "closing"); + }); + assertTrue(removed.await(5, TimeUnit.SECONDS)); + Future reset = + executor.submit( + () -> { + resettingThread.set(Thread.currentThread()); + resetting.countDown(); + manager.reset(); + return manager.getOrCreateCounter( + "closing", MetricLevel.IMPORTANT, "generation", "new"); + }); + assertTrue(resetting.await(5, TimeUnit.SECONDS)); + assertBlockedBy(resettingThread.get(), removingThread.get(), reset); + assertSame(original, manager.getAllMetrics()); + Object[] actual = executor.submit(() -> cachedMetrics(manager)).get(5, TimeUnit.SECONDS); + for (int i = 0; i < cached.length; i++) { + assertSame(cached[i], actual[i]); + } + ((Counter) actual[0]).inc(); + assertEquals(1, ((Counter) cached[0]).getCount()); + assertEquals(3, executor.submit(closing::getValue).get(5, TimeUnit.SECONDS), 0); + resume.countDown(); + removal.get(5, TimeUnit.SECONDS); + Counter replacement = reset.get(5, TimeUnit.SECONDS); + assertEquals(1, reporter.removals.get()); + assertSame(closing, reporter.removed.get()); + assertEquals(1, manager.getAllMetrics().size()); + assertSame( + replacement, + manager.getOrCreateCounter("closing", MetricLevel.IMPORTANT, "generation", "new")); + assertTrue(manager.hasMetadata("closing")); + } finally { + resume.countDown(); + executor.shutdownNow(); + assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + @Test + public void testConcurrentRemovalsNotifyOnlyOnce() throws Exception { + TestMetricManager manager = new TestMetricManager(); + RecordingReporter reporter = new RecordingReporter(manager); + manager.setBindJmxReporter(reporter); + Counter counter = manager.getOrCreateCounter("counter", MetricLevel.IMPORTANT); + ExecutorService executor = Executors.newFixedThreadPool(2); + CountDownLatch start = new CountDownLatch(1); + try { + List> removals = new ArrayList<>(); + for (int i = 0; i < 2; i++) { + removals.add( + executor.submit( + () -> { + await(start); + manager.remove(MetricType.COUNTER, "counter"); + })); + } + start.countDown(); + for (Future removal : removals) { + removal.get(5, TimeUnit.SECONDS); + } + manager.remove(MetricType.COUNTER, "counter"); + assertEquals(1, reporter.removals.get()); + assertSame(counter, reporter.removed.get()); + } finally { + start.countDown(); + executor.shutdownNow(); + assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + @Test + public void testConcurrentCreationPublishesOneMetricBeforeReporting() throws Exception { + TestMetricManager manager = new TestMetricManager(); + RecordingReporter reporter = new RecordingReporter(manager); + manager.setBindJmxReporter(reporter); + ExecutorService executor = Executors.newFixedThreadPool(4); + CountDownLatch start = new CountDownLatch(1); + try { + List> counters = new ArrayList<>(); + for (int i = 0; i < 4; i++) { + counters.add( + executor.submit( + () -> { + await(start); + return manager.getOrCreateCounter("counter", MetricLevel.IMPORTANT); + })); + } + start.countDown(); + Counter first = counters.get(0).get(5, TimeUnit.SECONDS); + for (Future counter : counters) { + assertSame(first, counter.get(5, TimeUnit.SECONDS)); + } + assertEquals(1, manager.counterCreations.get()); + assertEquals(1, reporter.registrations.get()); + } finally { + start.countDown(); + executor.shutdownNow(); + assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + @Test + public void testReporterCallbackDoesNotHoldRegistryLock() throws Exception { + TestMetricManager manager = new TestMetricManager(); + RecordingReporter reporter = new RecordingReporter(manager); + manager.setBindJmxReporter(reporter); + Counter old = manager.getOrCreateCounter("counter", MetricLevel.IMPORTANT); + CountDownLatch callback = new CountDownLatch(1); + CountDownLatch resume = new CountDownLatch(1); + reporter.beforeRemoval = + () -> { + callback.countDown(); + await(resume); + }; + ExecutorService executor = Executors.newFixedThreadPool(2); + try { + Future removal = executor.submit(() -> manager.remove(MetricType.COUNTER, "counter")); + assertTrue(callback.await(5, TimeUnit.SECONDS)); + // A reporter/MBean-server callback can run after reset without pinning the registry lock. + Counter next = + executor + .submit( + () -> { + manager.reset(); + return manager.getOrCreateCounter( + "counter", MetricLevel.IMPORTANT, "generation", "new"); + }) + .get(5, TimeUnit.SECONDS); + resume.countDown(); + removal.get(5, TimeUnit.SECONDS); + assertSame(old, reporter.removed.get()); + assertSame( + next, manager.getOrCreateCounter("counter", MetricLevel.IMPORTANT, "generation", "new")); + assertTrue(manager.hasMetadata("counter")); + } finally { + resume.countDown(); + executor.shutdownNow(); + assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + private static Object[] cachedMetrics(TestMetricManager manager) { + return new Object[] { + manager.getOrCreateCounter("cached_counter", MetricLevel.IMPORTANT), + manager.getOrCreateGauge("cached_gauge", MetricLevel.IMPORTANT), + manager.getOrCreateRate("cached_rate", MetricLevel.IMPORTANT), + manager.getOrCreateHistogram("cached_histogram", MetricLevel.IMPORTANT), + manager.getOrCreateTimer("cached_timer", MetricLevel.IMPORTANT) + }; + } + + private static void await(CountDownLatch latch) { + try { + assertTrue(latch.await(10, TimeUnit.SECONDS)); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AssertionError(e); + } + } + + private static void assertBlockedBy(Thread thread, Thread owner, Future task) { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5); + while (!task.isDone() && System.nanoTime() < deadline) { + ThreadInfo info = ManagementFactory.getThreadMXBean().getThreadInfo(thread.getId()); + if (info != null + && info.getThreadState() == Thread.State.BLOCKED + && info.getLockOwnerId() == owner.getId()) { + return; + } + LockSupport.parkNanos(TimeUnit.MILLISECONDS.toNanos(1)); + } + fail("Registry reset did not wait for the active removal"); + } + + private static class TestMetricManager extends DoNothingMetricManager { + private final AtomicInteger counterCreations = new AtomicInteger(); + + @Override + public boolean isEnableMetricInGivenLevel(MetricLevel level) { + return true; + } + + @Override + public Counter createCounter() { + counterCreations.incrementAndGet(); + return new IoTDBCounter(); + } + + @Override + public AutoGauge createAutoGauge(T object, ToDoubleFunction mapper) { + return new IoTDBAutoGauge<>(object, mapper); + } + + void useRegistry(Map registry) { + metrics = registry; + } + + void reset() { + stop(); + } + + boolean hasMetadata(String name) { + return nameToMetaInfo.containsKey(name); + } + } + + private static class RecordingReporter implements JmxReporter { + private final TestMetricManager manager; + private final AtomicInteger registrations = new AtomicInteger(); + private final AtomicInteger removals = new AtomicInteger(); + private final AtomicReference removed = new AtomicReference<>(); + private Runnable beforeRemoval = () -> {}; + + private RecordingReporter(TestMetricManager manager) { + this.manager = manager; + } + + @Override + public void registerMetric(IMetric metric, MetricInfo info) { + assertSame(metric, manager.getAllMetrics().get(info)); + registrations.incrementAndGet(); + } + + @Override + public void unregisterMetric(IMetric metric, MetricInfo info) { + assertNotNull(metric); + beforeRemoval.run(); + removed.set(metric); + removals.incrementAndGet(); + } + + @Override + public boolean start() { + return true; + } + + @Override + public boolean stop() { + return true; + } + + @Override + public ReporterType getReporterType() { + return ReporterType.JMX; + } + } +} diff --git a/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporterTest.java b/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporterTest.java index d247ef0b4ebe6..9bc6bd4fe0996 100644 --- a/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporterTest.java +++ b/iotdb-core/metrics/core/src/test/java/org/apache/iotdb/metrics/core/reporter/IoTDBJmxReporterTest.java @@ -37,6 +37,8 @@ import javax.management.MBeanServerFactory; import javax.management.ObjectName; +import java.util.ArrayList; +import java.util.List; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.ToDoubleFunction; @@ -49,6 +51,7 @@ public class IoTDBJmxReporterTest { private TestMetricManager manager; private MBeanServer server; private IoTDBJmxReporter reporter; + private final List values = new ArrayList<>(); @Before public void setUp() { @@ -98,6 +101,7 @@ public void testDelayedUnregisterPreservesReplacement() throws Exception { @Test public void testUnregisterDoesNotRemoveUnownedMBean() throws Exception { AtomicInteger value = new AtomicInteger(3); + values.add(value); IoTDBAutoGauge external = new IoTDBAutoGauge<>(value, AtomicInteger::get); ObjectName name = IoTDBMetricObjNameFactory.getInstance() @@ -132,7 +136,39 @@ public void testCounterRegistration() throws Exception { assertEquals(7L, server.getAttribute(counter.objectName(), "Count")); } + @Test + public void testDelayedRegistrationCannotResurrectRemovedMetric() { + IoTDBAutoGauge gauge = gauge(new AtomicInteger(1), "pool"); + manager.remove(MetricType.AUTO_GAUGE, "client_manager", "name", "num_active", "type", "pool"); + reporter.registerMetric(gauge, info("pool")); + assertFalse(server.isRegistered(gauge.objectName())); + } + + @Test + public void testDelayedRegistrationCannotReplaceNewMetric() throws Exception { + IoTDBAutoGauge old = gauge(new AtomicInteger(1), "pool"); + IoTDBAutoGauge next = gauge(new AtomicInteger(2), "pool"); + reporter.registerMetric(old, info("pool")); + assertEquals(2, (double) server.getAttribute(next.objectName(), "Value"), 0); + } + + @Test + public void testRegistrationWhileStoppedIsDeferred() throws Exception { + assertTrue(reporter.stop()); + IoTDBAutoGauge gauge = gauge(new AtomicInteger(1), "pool"); + ObjectName pattern = new ObjectName("org.apache.iotdb.metrics:*"); + assertTrue(server.queryNames(pattern, null).isEmpty()); + for (int i = 0; i < 3; i++) { + assertTrue(reporter.start()); + assertEquals(1, server.queryNames(pattern, null).size()); + assertEquals(1, (double) server.getAttribute(gauge.objectName(), "Value"), 0); + assertTrue(reporter.stop()); + assertTrue(server.queryNames(pattern, null).isEmpty()); + } + } + private IoTDBAutoGauge gauge(AtomicInteger value, String pool) { + values.add(value); return (IoTDBAutoGauge) manager.createAutoGauge( "client_manager", diff --git a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/AbstractMetricManager.java b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/AbstractMetricManager.java index cf5df80f36a85..b58b52fd5b116 100644 --- a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/AbstractMetricManager.java +++ b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/AbstractMetricManager.java @@ -21,7 +21,6 @@ import org.apache.iotdb.metrics.config.MetricConfig; import org.apache.iotdb.metrics.config.MetricConfigDescriptor; -import org.apache.iotdb.metrics.i18n.MetricsMessages; import org.apache.iotdb.metrics.impl.DoNothingMetricManager; import org.apache.iotdb.metrics.reporter.JmxReporter; import org.apache.iotdb.metrics.type.AutoGauge; @@ -53,13 +52,18 @@ public abstract class AbstractMetricManager { private static final String ALREADY_EXISTS = " is already used for a different type of name"; /** The map from metric name to metric metaInfo. */ - protected Map nameToMetaInfo; + protected volatile Map nameToMetaInfo; /** The map from metricInfo to metric. */ - protected Map metrics; + protected volatile Map metrics; /** The bind IoTDBJmxReporter */ - protected JmxReporter bindJmxReporter = null; + protected volatile JmxReporter bindJmxReporter = null; + + /** + * Serializes creation/removal with registry resets; existing metric reads do not take this lock. + */ + private final Object metricLifecycleLock = new Object(); protected AbstractMetricManager() { nameToMetaInfo = new ConcurrentHashMap<>(); @@ -72,9 +76,9 @@ protected AbstractMetricManager() { * @param metric the created metric * @param metricInfo the created metric info */ - private void notifyReporterOnAdd(IMetric metric, MetricInfo metricInfo) { + private void notifyReporterOnAdd(IMetric metric, MetricInfo metricInfo, JmxReporter reporter) { // if the reporter type is JMX, register the new metric - Optional.ofNullable(bindJmxReporter) + Optional.ofNullable(reporter) .ifPresent(x -> bindJmxReporter.registerMetric(metric, metricInfo)); } @@ -84,9 +88,9 @@ private void notifyReporterOnAdd(IMetric metric, MetricInfo metricInfo) { * @param metric the removed metric * @param metricInfo the removed metric info */ - private void notifyReporterOnRemove(IMetric metric, MetricInfo metricInfo) { + private void notifyReporterOnRemove(IMetric metric, MetricInfo metricInfo, JmxReporter reporter) { // if the reporter type is JMX, unregister the new metric - Optional.ofNullable(bindJmxReporter) + Optional.ofNullable(reporter) .ifPresent(x -> bindJmxReporter.unregisterMetric(metric, metricInfo)); } @@ -103,15 +107,23 @@ public Counter getOrCreateCounter(String name, MetricLevel metricLevel, String.. return DoNothingMetricManager.DO_NOTHING_COUNTER; } MetricInfo metricInfo = new MetricInfo(MetricType.COUNTER, name, tags); - IMetric metric = - metrics.computeIfAbsent( - metricInfo, - key -> { - Counter counter = createCounter(); - nameToMetaInfo.put(name, metricInfo.getMetaInfo()); - notifyReporterOnAdd(counter, metricInfo); - return counter; - }); + IMetric metric = metrics.get(metricInfo); + if (metric == null) { + JmxReporter reporter = null; + synchronized (metricLifecycleLock) { + if (invalid(metricLevel, name, tags)) { + return DoNothingMetricManager.DO_NOTHING_COUNTER; + } + metric = metrics.get(metricInfo); + if (metric == null) { + metric = createCounter(); + nameToMetaInfo.put(name, metricInfo.getMetaInfo()); + metrics.put(metricInfo, metric); + reporter = bindJmxReporter; + } + } + notifyReporterOnAdd(metric, metricInfo, reporter); + } if (metric instanceof Counter) { return (Counter) metric; } @@ -134,14 +146,24 @@ public Counter getOrCreateCounter(String name, MetricLevel metricLevel, String.. */ public AutoGauge createAutoGauge( String name, MetricLevel metricLevel, T obj, ToDoubleFunction mapper, String... tags) { - if (invalid(metricLevel, name, tags)) { - return DoNothingMetricManager.DO_NOTHING_AUTO_GAUGE; - } MetricInfo metricInfo = new MetricInfo(MetricType.AUTO_GAUGE, name, tags); - AutoGauge gauge = createAutoGauge(obj, mapper); - nameToMetaInfo.put(name, metricInfo.getMetaInfo()); - metrics.put(metricInfo, gauge); - notifyReporterOnAdd(gauge, metricInfo); + AutoGauge gauge; + IMetric previous; + JmxReporter reporter; + synchronized (metricLifecycleLock) { + if (invalid(metricLevel, name, tags)) { + return DoNothingMetricManager.DO_NOTHING_AUTO_GAUGE; + } + gauge = createAutoGauge(obj, mapper); + nameToMetaInfo.put(name, metricInfo.getMetaInfo()); + previous = metrics.put(metricInfo, gauge); + reporter = bindJmxReporter; + } + // Publish before notifying, and never hold the registry lock across MBean-server callbacks. + if (previous != null) { + notifyReporterOnRemove(previous, metricInfo, reporter); + } + notifyReporterOnAdd(gauge, metricInfo, reporter); return gauge; } @@ -187,15 +209,23 @@ public Gauge getOrCreateGauge(String name, MetricLevel metricLevel, String... ta return DoNothingMetricManager.DO_NOTHING_GAUGE; } MetricInfo metricInfo = new MetricInfo(MetricType.GAUGE, name, tags); - IMetric metric = - metrics.computeIfAbsent( - metricInfo, - key -> { - Gauge gauge = createGauge(); - nameToMetaInfo.put(name, metricInfo.getMetaInfo()); - notifyReporterOnAdd(gauge, metricInfo); - return gauge; - }); + IMetric metric = metrics.get(metricInfo); + if (metric == null) { + JmxReporter reporter = null; + synchronized (metricLifecycleLock) { + if (invalid(metricLevel, name, tags)) { + return DoNothingMetricManager.DO_NOTHING_GAUGE; + } + metric = metrics.get(metricInfo); + if (metric == null) { + metric = createGauge(); + nameToMetaInfo.put(name, metricInfo.getMetaInfo()); + metrics.put(metricInfo, metric); + reporter = bindJmxReporter; + } + } + notifyReporterOnAdd(metric, metricInfo, reporter); + } if (metric instanceof Gauge) { return (Gauge) metric; } @@ -218,15 +248,23 @@ public Rate getOrCreateRate(String name, MetricLevel metricLevel, String... tags return DoNothingMetricManager.DO_NOTHING_RATE; } MetricInfo metricInfo = new MetricInfo(MetricType.RATE, name, tags); - IMetric metric = - metrics.computeIfAbsent( - metricInfo, - key -> { - Rate rate = createRate(); - nameToMetaInfo.put(name, metricInfo.getMetaInfo()); - notifyReporterOnAdd(rate, metricInfo); - return rate; - }); + IMetric metric = metrics.get(metricInfo); + if (metric == null) { + JmxReporter reporter = null; + synchronized (metricLifecycleLock) { + if (invalid(metricLevel, name, tags)) { + return DoNothingMetricManager.DO_NOTHING_RATE; + } + metric = metrics.get(metricInfo); + if (metric == null) { + metric = createRate(); + nameToMetaInfo.put(name, metricInfo.getMetaInfo()); + metrics.put(metricInfo, metric); + reporter = bindJmxReporter; + } + } + notifyReporterOnAdd(metric, metricInfo, reporter); + } if (metric instanceof Rate) { return (Rate) metric; } @@ -249,15 +287,23 @@ public Histogram getOrCreateHistogram(String name, MetricLevel metricLevel, Stri return DoNothingMetricManager.DO_NOTHING_HISTOGRAM; } MetricInfo metricInfo = new MetricInfo(MetricType.HISTOGRAM, name, tags); - IMetric metric = - metrics.computeIfAbsent( - metricInfo, - key -> { - Histogram histogram = createHistogram(); - nameToMetaInfo.put(name, metricInfo.getMetaInfo()); - notifyReporterOnAdd(histogram, metricInfo); - return histogram; - }); + IMetric metric = metrics.get(metricInfo); + if (metric == null) { + JmxReporter reporter = null; + synchronized (metricLifecycleLock) { + if (invalid(metricLevel, name, tags)) { + return DoNothingMetricManager.DO_NOTHING_HISTOGRAM; + } + metric = metrics.get(metricInfo); + if (metric == null) { + metric = createHistogram(); + nameToMetaInfo.put(name, metricInfo.getMetaInfo()); + metrics.put(metricInfo, metric); + reporter = bindJmxReporter; + } + } + notifyReporterOnAdd(metric, metricInfo, reporter); + } if (metric instanceof Histogram) { return (Histogram) metric; } @@ -280,15 +326,23 @@ public Timer getOrCreateTimer(String name, MetricLevel metricLevel, String... ta return DoNothingMetricManager.DO_NOTHING_TIMER; } MetricInfo metricInfo = new MetricInfo(MetricType.TIMER, name, tags); - IMetric metric = - metrics.computeIfAbsent( - metricInfo, - key -> { - Timer timer = createTimer(); - nameToMetaInfo.put(name, metricInfo.getMetaInfo()); - notifyReporterOnAdd(timer, metricInfo); - return timer; - }); + IMetric metric = metrics.get(metricInfo); + if (metric == null) { + JmxReporter reporter = null; + synchronized (metricLifecycleLock) { + if (invalid(metricLevel, name, tags)) { + return DoNothingMetricManager.DO_NOTHING_TIMER; + } + metric = metrics.get(metricInfo); + if (metric == null) { + metric = createTimer(); + nameToMetaInfo.put(name, metricInfo.getMetaInfo()); + metrics.put(metricInfo, metric); + reporter = bindJmxReporter; + } + } + notifyReporterOnAdd(metric, metricInfo, reporter); + } if (metric instanceof Timer) { return (Timer) metric; } @@ -419,26 +473,27 @@ public Map getMetricsByType(MetricType metricType) { // region remove metric /** - * remove name. + * Remove a metric. Removing an already absent metric is a no-op. * * @param type the type of name * @param name the name of name * @param tags string pairs, like sg="ln" will be "sg", "ln" - * @throws IllegalArgumentException when there has different type metric with same name */ public void remove(MetricType type, String name, String... tags) { MetricInfo metricInfo = new MetricInfo(type, name, tags); - if (metrics.containsKey(metricInfo)) { - if (type == metricInfo.getMetaInfo().getType()) { - notifyReporterOnRemove(metrics.get(metricInfo), metricInfo); - nameToMetaInfo.remove(metricInfo.getName()); - metrics.remove(metricInfo); - removeMetric(type, metricInfo); - } else { - throw new IllegalArgumentException( - metricInfo + MetricsMessages.EXCEPTION_FAILED_REMOVE_BECAUSE_MISMATCH_TYPE_044E55F6); + IMetric removed; + JmxReporter reporter; + synchronized (metricLifecycleLock) { + removed = metrics.remove(metricInfo); + if (removed == null) { + return; } + nameToMetaInfo.remove(metricInfo.getName()); + removeMetric(type, metricInfo); + reporter = bindJmxReporter; } + // The removed instance remains valid even if stop() replaces the registry before this callback. + notifyReporterOnRemove(removed, metricInfo, reporter); } protected abstract void removeMetric(MetricType type, MetricInfo metricInfo); @@ -451,14 +506,18 @@ public boolean isEnableMetricInGivenLevel(MetricLevel metricLevel) { } public void setBindJmxReporter(JmxReporter reporter) { - this.bindJmxReporter = reporter; + synchronized (metricLifecycleLock) { + this.bindJmxReporter = reporter; + } } /** Stop and clear metric manager. */ protected boolean stop() { - metrics = new ConcurrentHashMap<>(); - nameToMetaInfo = new ConcurrentHashMap<>(); - return stopFramework(); + synchronized (metricLifecycleLock) { + metrics = new ConcurrentHashMap<>(); + nameToMetaInfo = new ConcurrentHashMap<>(); + return stopFramework(); + } } protected abstract boolean stopFramework(); diff --git a/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/ClientMessages.java b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/ClientMessages.java index a4a3a8e6feb9f..ac17ee9a866ff 100644 --- a/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/ClientMessages.java +++ b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/ClientMessages.java @@ -30,6 +30,9 @@ public final class ClientMessages { public static final String CLEAR_CLIENT_POOL_FAILED = "Clear all client in pool for node {} failed."; + public static final String LOG_FAILED_TO_UNREGISTER_CLIENT_POOL_METRICS_WHILE_CLOSING_CLIENT_MANAGER_101A9751 = + "Failed to unregister client pool metrics while closing client manager"; + // ThriftClient public static final String EXCEPTION_LEVEL_DETAIL = "level-{} Exception class {}, message {}"; diff --git a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/ClientMessages.java b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/ClientMessages.java index 1f4ad1070a841..06715a051d8a3 100644 --- a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/ClientMessages.java +++ b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/ClientMessages.java @@ -30,6 +30,9 @@ public final class ClientMessages { public static final String CLEAR_CLIENT_POOL_FAILED = "清除节点 {} 的所有连接池客户端失败。"; + public static final String LOG_FAILED_TO_UNREGISTER_CLIENT_POOL_METRICS_WHILE_CLOSING_CLIENT_MANAGER_101A9751 = + "关闭客户端管理器时,注销客户端连接池指标失败"; + // ThriftClient public static final String EXCEPTION_LEVEL_DETAIL = "第 {} 层异常,类名 {},消息 {}"; diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManager.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManager.java index 10c1caf18a1c4..336658fd5dd95 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManager.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/client/ClientManager.java @@ -132,7 +132,14 @@ public void close() { ((AsyncThriftClientFactory) pool.getFactory()).close(); } } finally { - ClientManagerMetrics.getInstance().unregisterClientManager(pool); + try { + ClientManagerMetrics.getInstance().unregisterClientManager(pool); + } catch (RuntimeException e) { + LOGGER.warn( + ClientMessages + .LOG_FAILED_TO_UNREGISTER_CLIENT_POOL_METRICS_WHILE_CLOSING_CLIENT_MANAGER_101A9751, + e); + } } } } diff --git a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java index 3c6f223272bb0..e3eda335b1644 100644 --- a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java +++ b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/client/ClientManagerMetricsTest.java @@ -22,6 +22,7 @@ import org.apache.iotdb.commons.service.metric.MetricService; import org.apache.iotdb.commons.service.metric.enums.Metric; import org.apache.iotdb.commons.service.metric.enums.Tag; +import org.apache.iotdb.metrics.AbstractMetricManager; import org.apache.iotdb.metrics.DoNothingMetricService; import org.apache.iotdb.metrics.config.MetricConfig; import org.apache.iotdb.metrics.config.MetricConfigDescriptor; @@ -47,10 +48,13 @@ import javax.management.ObjectName; import java.lang.management.ManagementFactory; +import java.lang.management.ThreadInfo; +import java.lang.reflect.Field; import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -63,7 +67,12 @@ import static org.awaitility.Awaitility.await; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; +import static org.mockito.Mockito.doNothing; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; public class ClientManagerMetricsTest { @@ -286,12 +295,153 @@ public void testClosingAnotherPoolPreservesAllLivePoolMBeans() throws Exception } private void assertJmxMetrics(MBeanServer server) throws Exception { - for (IMetric metric : service.getAllMetrics().values()) { + assertJmxMetrics(server, service.getAllMetrics()); + } + + private void assertJmxMetrics(MBeanServer server, Map allMetrics) + throws Exception { + for (IMetric metric : allMetrics.values()) { IoTDBAutoGauge gauge = (IoTDBAutoGauge) metric; assertEquals(gauge.getValue(), (double) server.getAttribute(gauge.objectName(), "Value"), 0); } } + @Test + public void testCloseAndCoreRestartAreSerialized() throws Exception { + service.removeMetricSet(metrics); + config.setMetricReporterList("JMX"); + MetricService realService = MetricService.getInstance(); + realService.startService(); + CountDownLatch removed = new CountDownLatch(1); + CountDownLatch resume = new CountDownLatch(1); + ExecutorService executor = Executors.newFixedThreadPool(3); + try { + realService.addMetricSet(metrics); + ClientManager closing = createManager("closing"); + closing.borrowClient("node"); + ClientManager live = createManager("live"); + live.borrowClient("node"); + AutoGauge liveGauge = + realService.getAutoGauge( + Metric.CLIENT_MANAGER.toString(), + MetricLevel.IMPORTANT, + Tag.NAME.toString(), + "client_manager_num_active", + Tag.TYPE.toString(), + "live"); + AbstractMetricManager manager = realService.getMetricManager(); + AtomicBoolean pauseOnce = new AtomicBoolean(true); + Map registry = + new ConcurrentHashMap(manager.getAllMetrics()) { + @Override + public IMetric remove(Object key) { + IMetric metric = super.remove(key); + if (metric != null + && "closing".equals(((MetricInfo) key).getTags().get(Tag.TYPE.toString())) + && pauseOnce.compareAndSet(true, false)) { + removed.countDown(); + try { + assertTrue(resume.await(10, TimeUnit.SECONDS)); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AssertionError(e); + } + } + return metric; + } + }; + Field field = AbstractMetricManager.class.getDeclaredField("metrics"); + field.setAccessible(true); + field.set(manager, registry); + AtomicReference closingThread = new AtomicReference<>(); + AtomicReference restartingThread = new AtomicReference<>(); + CountDownLatch restarting = new CountDownLatch(1); + Future close = + executor.submit( + () -> { + closingThread.set(Thread.currentThread()); + closing.close(); + }); + assertTrue(removed.await(5, TimeUnit.SECONDS)); + Future restart = + executor.submit( + () -> { + restartingThread.set(Thread.currentThread()); + restarting.countDown(); + realService.restartService(); + }); + assertTrue(restarting.await(5, TimeUnit.SECONDS)); + await() + .atMost(5, TimeUnit.SECONDS) + .until( + () -> { + ThreadInfo info = + ManagementFactory.getThreadMXBean() + .getThreadInfo(restartingThread.get().getId()); + return restart.isDone() + || (info != null + && info.getThreadState() == Thread.State.BLOCKED + && info.getLockOwnerId() == closingThread.get().getId()); + }); + assertFalse(restart.isDone()); + // The core registry must not change while removal is paused, even before metric-set rebind. + assertSame(registry, manager.getAllMetrics()); + assertEquals(1, executor.submit(liveGauge::getValue).get(5, TimeUnit.SECONDS), 0); + resume.countDown(); + close.get(5, TimeUnit.SECONDS); + restart.get(5, TimeUnit.SECONDS); + assertEquals(8, realService.getAllMetrics().size()); + assertMetrics(realService.getAllMetrics(), "live", live.getPool()); + assertJmxMetrics(ManagementFactory.getPlatformMBeanServer(), realService.getAllMetrics()); + } finally { + resume.countDown(); + executor.shutdownNow(); + assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + realService.removeMetricSet(metrics); + realService.stopService(); + } + } + + @Test + public void testMetricCleanupFailureDoesNotFailClose() { + ClientManager manager = createManager("pool"); + service.beforeRemove = + () -> { + throw new IllegalStateException("metric cleanup failure"); + }; + try { + manager.close(); + assertTrue(manager.getPool().isClosed()); + } finally { + service.beforeRemove = () -> {}; + } + } + + @Test + public void testMetricCleanupFailureDoesNotMaskPoolCloseFailure() { + @SuppressWarnings("unchecked") + GenericKeyedObjectPool pool = mock(GenericKeyedObjectPool.class); + ClientManager manager = + new ClientManager<>( + owner -> { + metrics.registerClientManager("pool", pool); + return pool; + }); + managers.add(manager); + IllegalStateException poolFailure = new IllegalStateException("pool close failure"); + doThrow(poolFailure).when(pool).close(); + service.beforeRemove = + () -> { + throw new IllegalStateException("metric cleanup failure"); + }; + try { + assertSame(poolFailure, assertThrows(IllegalStateException.class, manager::close)); + } finally { + doNothing().when(pool).close(); + service.beforeRemove = () -> {}; + } + } + @Test public void testRegistrationAndSamplingDuringUnbind() throws Exception { ClientManager first = createManager("first"); From 87aa3dd9155c2f028dca35975f9c324adbef5779 Mon Sep 17 00:00:00 2001 From: JackieTien97 Date: Sun, 20 Sep 2026 19:41:59 +0800 Subject: [PATCH 4/4] Use captured JMX reporters for metric notifications --- .../iotdb/metrics/AbstractMetricManager.java | 14 +- .../metric/MetricReporterSwitchTest.java | 166 ++++++++++++++++++ 2 files changed, 173 insertions(+), 7 deletions(-) create mode 100644 iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/service/metric/MetricReporterSwitchTest.java diff --git a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/AbstractMetricManager.java b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/AbstractMetricManager.java index b58b52fd5b116..ce9330ca5ef04 100644 --- a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/AbstractMetricManager.java +++ b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/AbstractMetricManager.java @@ -76,10 +76,10 @@ protected AbstractMetricManager() { * @param metric the created metric * @param metricInfo the created metric info */ - private void notifyReporterOnAdd(IMetric metric, MetricInfo metricInfo, JmxReporter reporter) { + private static void notifyReporterOnAdd( + IMetric metric, MetricInfo metricInfo, JmxReporter reporter) { // if the reporter type is JMX, register the new metric - Optional.ofNullable(reporter) - .ifPresent(x -> bindJmxReporter.registerMetric(metric, metricInfo)); + Optional.ofNullable(reporter).ifPresent(x -> x.registerMetric(metric, metricInfo)); } /** @@ -88,10 +88,10 @@ private void notifyReporterOnAdd(IMetric metric, MetricInfo metricInfo, JmxRepor * @param metric the removed metric * @param metricInfo the removed metric info */ - private void notifyReporterOnRemove(IMetric metric, MetricInfo metricInfo, JmxReporter reporter) { - // if the reporter type is JMX, unregister the new metric - Optional.ofNullable(reporter) - .ifPresent(x -> bindJmxReporter.unregisterMetric(metric, metricInfo)); + private static void notifyReporterOnRemove( + IMetric metric, MetricInfo metricInfo, JmxReporter reporter) { + // Use the captured reporter even if the current binding has changed. + Optional.ofNullable(reporter).ifPresent(x -> x.unregisterMetric(metric, metricInfo)); } /** diff --git a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/service/metric/MetricReporterSwitchTest.java b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/service/metric/MetricReporterSwitchTest.java new file mode 100644 index 0000000000000..b2da6ca97085e --- /dev/null +++ b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/service/metric/MetricReporterSwitchTest.java @@ -0,0 +1,166 @@ +/* + * 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.commons.service.metric; + +import org.apache.iotdb.metrics.AbstractMetricManager; +import org.apache.iotdb.metrics.config.MetricConfig; +import org.apache.iotdb.metrics.config.MetricConfigDescriptor; +import org.apache.iotdb.metrics.config.ReloadLevel; +import org.apache.iotdb.metrics.core.reporter.IoTDBJmxReporter; +import org.apache.iotdb.metrics.reporter.JmxReporter; +import org.apache.iotdb.metrics.type.AutoGauge; +import org.apache.iotdb.metrics.type.IMetric; +import org.apache.iotdb.metrics.utils.MetricInfo; +import org.apache.iotdb.metrics.utils.MetricLevel; +import org.apache.iotdb.metrics.utils.ReporterType; + +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +import javax.management.MBeanServer; +import javax.management.ObjectName; + +import java.lang.management.ManagementFactory; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +public class MetricReporterSwitchTest { + private static final String METRIC = "reporter_switch_value"; + private final MetricConfig config = MetricConfigDescriptor.getInstance().getMetricConfig(); + private final MetricService service = MetricService.getInstance(); + private final AtomicInteger oldValue = new AtomicInteger(1); + private final AtomicInteger newValue = new AtomicInteger(2); + private MetricLevel previousLevel; + private String previousReporters; + + @Before + public void setUp() { + previousLevel = config.getMetricLevel(); + previousReporters = + config.getMetricReporterList().stream().map(Enum::name).collect(Collectors.joining(",")); + config.setMetricLevel(MetricLevel.IMPORTANT); + config.setMetricReporterList("JMX"); + service.startService(); + } + + @After + public void tearDown() { + service.getMetricManager().setBindJmxReporter(null); + service.stopService(); + config.setMetricLevel(previousLevel); + config.setMetricReporterList(previousReporters); + } + + @Test + public void testDisablingReporterDuringGaugeReplacement() throws Exception { + replaceGaugeWhileReloading(""); + } + + @Test + public void testReloadingReporterKeepsCapturedNotificationTarget() throws Exception { + replaceGaugeWhileReloading("JMX"); + } + + private void replaceGaugeWhileReloading(String reporters) throws Exception { + AbstractMetricManager manager = service.getMetricManager(); + IoTDBJmxReporter delegate = IoTDBJmxReporter.getInstance(); + CountDownLatch removalCallback = new CountDownLatch(1); + CountDownLatch resume = new CountDownLatch(1); + AtomicInteger registrations = new AtomicInteger(); + AtomicInteger removals = new AtomicInteger(); + // Preserve real JMX behavior while pausing between replacement's remove/add notifications. + manager.setBindJmxReporter( + new JmxReporter() { + @Override + public void registerMetric(IMetric metric, MetricInfo info) { + registrations.incrementAndGet(); + delegate.registerMetric(metric, info); + } + + @Override + public void unregisterMetric(IMetric metric, MetricInfo info) { + removals.incrementAndGet(); + delegate.unregisterMetric(metric, info); + removalCallback.countDown(); + try { + assertTrue(resume.await(10, TimeUnit.SECONDS)); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AssertionError(e); + } + } + + @Override + public boolean start() { + return delegate.start(); + } + + @Override + public boolean stop() { + return delegate.stop(); + } + + @Override + public ReporterType getReporterType() { + return ReporterType.JMX; + } + }); + manager.createAutoGauge(METRIC, MetricLevel.IMPORTANT, oldValue, AtomicInteger::get); + ExecutorService executor = Executors.newSingleThreadExecutor(); + try { + Future replacement = + executor.submit( + () -> + manager.createAutoGauge( + METRIC, MetricLevel.IMPORTANT, newValue, AtomicInteger::get)); + assertTrue(removalCallback.await(5, TimeUnit.SECONDS)); + config.setMetricReporterList(reporters); + service.reloadService(ReloadLevel.RESTART_REPORTER); + resume.countDown(); + assertEquals(2, replacement.get(5, TimeUnit.SECONDS).getValue(), 0); + // Reload either clears the binding or replaces our wrapper with the real JMX reporter. + // Both notifications for the in-flight replacement must still use the captured wrapper. + assertEquals(2, registrations.get()); + assertEquals(1, removals.get()); + MBeanServer server = ManagementFactory.getPlatformMBeanServer(); + ObjectName name = + new ObjectName("org.apache.iotdb.metrics:name=" + METRIC + ",type=IoTDBAutoGauge"); + if (reporters.isEmpty()) { + assertFalse(server.isRegistered(name)); + } else { + assertEquals(2, (double) server.getAttribute(name, "Value"), 0); + } + } finally { + resume.countDown(); + executor.shutdownNow(); + assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS)); + } + } +}