This is an automated email from the ASF dual-hosted git repository.
smengcl pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/master by this push:
new bd0f23b49ba HDDS-16119. Fix DatanodeStorageMetrics metrics-system
deadlock and mini-cluster source leak (#10988)
bd0f23b49ba is described below
commit bd0f23b49bacdddc2a31f7fc565f9165ca7f40dc
Author: Siyao Meng <[email protected]>
AuthorDate: Tue Aug 11 10:10:33 2026 -0700
HDDS-16119. Fix DatanodeStorageMetrics metrics-system deadlock and
mini-cluster source leak (#10988)
Co-authored-by: Claude Opus 4.8 (1M context) <[email protected]>
---
.../common/volume/DatanodeStorageMetrics.java | 39 ++++++++++++----
.../container/common/volume/MutableVolumeSet.java | 37 +++++++++++++++
.../common/volume/TestDatanodeStorageMetrics.java | 52 +++++++++++++++++++++-
.../container/common/volume/TestVolumeSet.java | 27 +++++++++++
.../dn/TestDatanodeStorageMetricsIntegration.java | 9 +++-
5 files changed, 151 insertions(+), 13 deletions(-)
diff --git
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/volume/DatanodeStorageMetrics.java
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/volume/DatanodeStorageMetrics.java
index fc1edbee11e..b0fdd9036fa 100644
---
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/volume/DatanodeStorageMetrics.java
+++
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/volume/DatanodeStorageMetrics.java
@@ -30,10 +30,13 @@
/**
* Node-level storage totals for a DataNode, aggregated over its HDDS data
volumes only
- * ({@code VolumeType.DATA_VOLUME}) via {@link
MutableVolumeSet#getStorageReport()}.
+ * ({@code VolumeType.DATA_VOLUME}) via {@link
MutableVolumeSet#getStorageReportSnapshot()}.
* This is the same scope as the {@code storageReport} entries produced by
* {@code OzoneContainer.getNodeReport()}; meta and DB volumes are excluded.
* Registered as {@code
Hadoop:service=HddsDatanode,name=DatanodeStorageMetrics}.
+ * In mini-cluster mode, where many datanodes share one metrics system, the
name
+ * is suffixed with the datanode UUID ({@code DatanodeStorageMetrics-<uuid>})
so
+ * registration and unregistration stay unique per datanode.
*/
@Metrics(about = "Ozone DataNode node-level storage totals",
context = OzoneConsts.OZONE)
@@ -53,10 +56,12 @@ public final class DatanodeStorageMetrics implements
MetricsSource {
private final MetricsRegistry registry;
private final MutableVolumeSet volumeSet;
+ private final String sourceName;
- private DatanodeStorageMetrics(MutableVolumeSet volumeSet) {
+ private DatanodeStorageMetrics(MutableVolumeSet volumeSet, String
sourceName) {
this.volumeSet = volumeSet;
- this.registry = new MetricsRegistry(SOURCE_NAME);
+ this.sourceName = sourceName;
+ this.registry = new MetricsRegistry(sourceName);
}
/**
@@ -64,9 +69,20 @@ private DatanodeStorageMetrics(MutableVolumeSet volumeSet) {
* with the default Metrics2 system.
*/
public static DatanodeStorageMetrics create(MutableVolumeSet volumeSet) {
- DatanodeStorageMetrics datanodeStorageMetrics = new
DatanodeStorageMetrics(volumeSet);
- DefaultMetricsSystem.instance().register(
- SOURCE_NAME, "DataNode node-level storage totals",
datanodeStorageMetrics);
+ // In mini-cluster mode many datanodes share one metrics system, so the
+ // metrics system uniquifies the constant source name on registration
+ // (DatanodeStorageMetrics-1, -2, ...). Unregistering by the constant base
+ // name would then leak every source past the first, and each leaked source
+ // pins a shut-down datanode's MutableVolumeSet. Make the name unique per
+ // datanode up front so register and unregister stay symmetric. In
+ // production there is one instance per JVM, so keep the plain name for
+ // stable JMX and Prometheus metric names.
+ String sourceName = DefaultMetricsSystem.inMiniClusterMode()
+ ? SOURCE_NAME + '-' + volumeSet.getDatanodeUuid()
+ : SOURCE_NAME;
+ DatanodeStorageMetrics datanodeStorageMetrics = new
DatanodeStorageMetrics(volumeSet, sourceName);
+ DefaultMetricsSystem.instance().register(sourceName,
+ "DataNode node-level storage totals", datanodeStorageMetrics);
return datanodeStorageMetrics;
}
@@ -74,7 +90,7 @@ public static DatanodeStorageMetrics create(MutableVolumeSet
volumeSet) {
* Unregisters this source from the Metrics2 system.
*/
public void unregister() {
- DefaultMetricsSystem.instance().unregisterSource(SOURCE_NAME);
+ DefaultMetricsSystem.instance().unregisterSource(sourceName);
}
/**
@@ -83,12 +99,17 @@ public void unregister() {
*/
@Override
public void getMetrics(MetricsCollector collector, boolean all) {
- MetricsRecordBuilder builder = collector.addRecord(SOURCE_NAME);
+ MetricsRecordBuilder builder = collector.addRecord(sourceName);
registry.snapshot(builder, all);
long capacity = 0L;
long used = 0L;
- for (StorageLocationReport report : volumeSet.getStorageReport()) {
+ // getMetrics() runs while the DefaultMetricsSystem monitor is held. Read a
+ // lock-free snapshot instead of getStorageReport(), which takes the
+ // volume-set lock: a volume-failure handler holds that lock while
+ // unregistering volume metrics (which needs the same monitor), so locking
+ // here can deadlock the metrics system.
+ for (StorageLocationReport report : volumeSet.getStorageReportSnapshot()) {
capacity = Math.addExact(capacity, report.getCapacity());
used = Math.addExact(used, report.getScmUsed());
}
diff --git
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/volume/MutableVolumeSet.java
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/volume/MutableVolumeSet.java
index a79a06b6541..4d60015a050 100644
---
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/volume/MutableVolumeSet.java
+++
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/volume/MutableVolumeSet.java
@@ -28,6 +28,7 @@
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.function.Function;
+import java.util.stream.Stream;
import org.apache.hadoop.hdds.conf.ConfigurationSource;
import org.apache.hadoop.hdds.fs.SpaceUsageCheckFactory;
import org.apache.hadoop.hdds.scm.ScmConfigKeys;
@@ -404,6 +405,16 @@ public boolean hasEnoughVolumes() {
return hasEnoughVolumes;
}
+ /**
+ * Returns a consistent snapshot of the storage reports under the volume-set
+ * read lock.
+ *
+ * <p>Do not call this from metrics collection (or any caller that may hold
the
+ * {@code DefaultMetricsSystem} monitor): a volume-failure handler holds the
+ * volume-set write lock while unregistering volume metrics under that same
+ * monitor, so taking the read lock here can deadlock the metrics system. Use
+ * {@link #getStorageReportSnapshot()} from those paths instead.
+ */
public StorageLocationReport[] getStorageReport() {
this.readLock();
try {
@@ -421,6 +432,32 @@ public StorageLocationReport[] getStorageReport() {
}
}
+ /**
+ * Lock-free variant of {@link #getStorageReport()}. Both {@code volumeMap}
+ * and {@code failedVolumeMap} are {@link ConcurrentHashMap}s, so this
returns
+ * a weakly-consistent snapshot (mirroring {@link #getVolumesList()}) without
+ * acquiring the volume-set lock.
+ *
+ * <p>Use this from callers that must not block on the volume-set lock. In
+ * particular, metrics sampling runs while the {@code DefaultMetricsSystem}
+ * monitor is held, and a volume-failure handler holds the volume-set write
+ * lock while unregistering volume metrics (which needs that same monitor);
+ * acquiring the volume-set lock from the sampling thread can therefore
+ * deadlock the whole metrics system.
+ */
+ public StorageLocationReport[] getStorageReportSnapshot() {
+ // volumeMap and failedVolumeMap are ConcurrentHashMaps; their value
streams
+ // are weakly consistent, so no lock is needed here (same guarantee as
+ // getVolumesList()).
+ return Stream.concat(volumeMap.values().stream(),
failedVolumeMap.values().stream())
+ .map(StorageVolume::getReport)
+ .toArray(StorageLocationReport[]::new);
+ }
+
+ public String getDatanodeUuid() {
+ return datanodeUuid;
+ }
+
public StorageVolume.VolumeType getVolumeType() {
return volumeType;
}
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/volume/TestDatanodeStorageMetrics.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/volume/TestDatanodeStorageMetrics.java
index d54a8b0f9ea..05d14410a4b 100644
---
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/volume/TestDatanodeStorageMetrics.java
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/volume/TestDatanodeStorageMetrics.java
@@ -21,9 +21,12 @@
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
+import java.util.UUID;
import org.apache.hadoop.metrics2.AbstractMetric;
+import org.apache.hadoop.metrics2.MetricsSystem;
import org.apache.hadoop.metrics2.impl.MetricsCollectorImpl;
import org.apache.hadoop.metrics2.impl.MetricsRecordImpl;
+import org.apache.hadoop.metrics2.lib.DefaultMetricsSystem;
import org.apache.hadoop.ozone.container.common.impl.StorageLocationReport;
import org.junit.jupiter.api.Test;
@@ -51,7 +54,7 @@ void testAggregationAcrossTwoVolumes() {
.build();
MutableVolumeSet volumeSet = mock(MutableVolumeSet.class);
- when(volumeSet.getStorageReport())
+ when(volumeSet.getStorageReportSnapshot())
.thenReturn(new StorageLocationReport[]{vol1, vol2});
DatanodeStorageMetrics metrics = DatanodeStorageMetrics.create(volumeSet);
@@ -78,7 +81,7 @@ void testAggregationAcrossTwoVolumes() {
void testZeroCapacityReturnsZeroPercentage() {
// No volumes → capacity=0, used=0; OzoneUsedPercentage must be 0.0, not
NaN.
MutableVolumeSet volumeSet = mock(MutableVolumeSet.class);
- when(volumeSet.getStorageReport()).thenReturn(new
StorageLocationReport[0]);
+ when(volumeSet.getStorageReportSnapshot()).thenReturn(new
StorageLocationReport[0]);
DatanodeStorageMetrics metrics = DatanodeStorageMetrics.create(volumeSet);
try {
@@ -94,6 +97,51 @@ void testZeroCapacityReturnsZeroPercentage() {
}
}
+ @Test
+ void testNoSourceLeakInMiniClusterMode() {
+ // In mini-cluster mode many datanodes share one metrics system. Register
+ // and unregister must be symmetric so no source (and its pinned volume
set)
+ // leaks. The pre-fix code registered a constant name (uniquified to -N)
but
+ // unregistered the base name, leaking every datanode past the first.
+ boolean prev = DefaultMetricsSystem.inMiniClusterMode();
+ DefaultMetricsSystem.setMiniClusterMode(true);
+ MetricsSystem ms = DefaultMetricsSystem.instance();
+ String uuidA = "dn-" + UUID.randomUUID();
+ String uuidB = "dn-" + UUID.randomUUID();
+ String nameA = DatanodeStorageMetrics.SOURCE_NAME + '-' + uuidA;
+ String nameB = DatanodeStorageMetrics.SOURCE_NAME + '-' + uuidB;
+ DatanodeStorageMetrics a = null;
+ DatanodeStorageMetrics b = null;
+ try {
+ a = DatanodeStorageMetrics.create(mockVolumeSet(uuidA));
+ b = DatanodeStorageMetrics.create(mockVolumeSet(uuidB));
+ assertThat(ms.getSource(nameA)).isNotNull();
+ assertThat(ms.getSource(nameB)).isNotNull();
+
+ a.unregister();
+ b.unregister();
+ assertThat(ms.getSource(nameA)).isNull();
+ assertThat(ms.getSource(nameB)).isNull();
+ } finally {
+ // Do not leak sources into other tests if an assertion above fails.
+ // unregister() is idempotent, so a repeat after the happy path is a
no-op.
+ if (a != null) {
+ a.unregister();
+ }
+ if (b != null) {
+ b.unregister();
+ }
+ DefaultMetricsSystem.setMiniClusterMode(prev);
+ }
+ }
+
+ private static MutableVolumeSet mockVolumeSet(String datanodeUuid) {
+ MutableVolumeSet volumeSet = mock(MutableVolumeSet.class);
+ when(volumeSet.getDatanodeUuid()).thenReturn(datanodeUuid);
+ when(volumeSet.getStorageReportSnapshot()).thenReturn(new
StorageLocationReport[0]);
+ return volumeSet;
+ }
+
private static long findLong(Iterable<AbstractMetric> metrics, String name) {
for (AbstractMetric m : metrics) {
if (name.equals(m.name())) {
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/volume/TestVolumeSet.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/volume/TestVolumeSet.java
index 932101dc526..64e1c00e70e 100644
---
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/volume/TestVolumeSet.java
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/volume/TestVolumeSet.java
@@ -25,6 +25,7 @@
import static org.assertj.core.api.Assumptions.assumeThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.io.File;
@@ -34,6 +35,9 @@
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
import org.apache.commons.io.FileUtils;
import org.apache.hadoop.hdds.HddsConfigKeys;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
@@ -194,6 +198,29 @@ void testFailVolumes(@TempDir File readOnlyVolumePath,
@TempDir File volumePath)
volSet.shutdown();
}
+ @Test
+ public void testStorageReportSnapshotDoesNotBlockOnWriteLock() throws
Exception {
+ // Regression test: DatanodeStorageMetrics.getMetrics() samples the volume
+ // set while the DefaultMetricsSystem monitor is held. It must not block on
+ // the volume-set lock, otherwise it deadlocks with a volume-failure
handler
+ // that holds the write lock while unregistering volume metrics.
+ volumeSet.writeLock();
+ try {
+ CompletableFuture<Integer> snapshot = CompletableFuture.supplyAsync(
+ () -> volumeSet.getStorageReportSnapshot().length);
+ // The snapshot must return without waiting for the write lock.
+ assertEquals(2, snapshot.get(10, TimeUnit.SECONDS));
+
+ // Sanity check that the lock is genuinely held: the locking variant does
+ // block behind the write lock (it must time out here).
+ CompletableFuture<Integer> locked = CompletableFuture.supplyAsync(
+ () -> volumeSet.getStorageReport().length);
+ assertThrows(TimeoutException.class, () -> locked.get(2,
TimeUnit.SECONDS));
+ } finally {
+ volumeSet.writeUnlock();
+ }
+ }
+
@Test
public void testInterrupt() throws Exception {
Method method = this.volumeSet.getClass()
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/TestDatanodeStorageMetricsIntegration.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/TestDatanodeStorageMetricsIntegration.java
index 076889a7e8f..774184e8689 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/TestDatanodeStorageMetricsIntegration.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/TestDatanodeStorageMetricsIntegration.java
@@ -123,7 +123,12 @@ void storageMetricsReflectWrittenData() throws Exception {
* Each call re-reads the underlying storage reports — do not mix values
* from different calls when checking invariants across gauges.
*/
- private static MetricsRecordBuilder storageMetrics() {
- return getMetrics(DatanodeStorageMetrics.SOURCE_NAME);
+ private MetricsRecordBuilder storageMetrics() {
+ // In mini-cluster mode the source name is made unique per datanode (to
keep
+ // metrics registration and unregistration symmetric and avoid the
+ // shared-JVM source leak), so look it up by the per-datanode name.
+ MutableVolumeSet volumeSet = cluster.getHddsDatanodes().get(0)
+ .getDatanodeStateMachine().getContainer().getVolumeSet();
+ return getMetrics(DatanodeStorageMetrics.SOURCE_NAME + '-' +
volumeSet.getDatanodeUuid());
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]