This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-11492-9ec9c4fe8aa62c628c638ba5a625d2624b8a4849 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 2a739a94a9ea6ab05ab8c86d0e1a739bf16c881e Author: Doyeon Kim <[email protected]> AuthorDate: Sat Aug 29 02:11:08 2026 +0000 [Fix][Zeta] Isolate node metrics from global registry (#11492) --- .../seatunnel/engine/server/NodeExtension.java | 31 ++++++++++++- .../server/rest/RestHttpGetCommandProcessor.java | 4 +- .../engine/server/rest/servlet/MetricsServlet.java | 8 ++-- .../metrics/ExportsInstanceInitializer.java | 51 +++++++++++----------- .../engine/server/metrics/MetricsApiTest.java | 15 ++++++- 5 files changed, 73 insertions(+), 36 deletions(-) diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/NodeExtension.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/NodeExtension.java index 6808648d8f..0e3b12a74d 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/NodeExtension.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/NodeExtension.java @@ -28,10 +28,12 @@ import com.hazelcast.instance.impl.DefaultNodeExtension; import com.hazelcast.instance.impl.Node; import com.hazelcast.internal.ascii.TextCommandService; import com.hazelcast.internal.ascii.TextCommandServiceImpl; +import io.prometheus.client.Collector.MetricFamilySamples; import io.prometheus.client.CollectorRegistry; import lombok.Getter; import lombok.NonNull; +import java.util.Enumeration; import java.util.Map; import static com.hazelcast.internal.ascii.TextCommandConstants.TextCommandType.HTTP_GET; @@ -46,7 +48,14 @@ public class NodeExtension extends DefaultNodeExtension { super(node); seaTunnelServer = new SeaTunnelServer(seaTunnelConfig); extCommon = new NodeExtensionCommon(node, seaTunnelServer); - collectorRegistry = CollectorRegistry.defaultRegistry; + collectorRegistry = new CollectorRegistry(true); + } + + /** Returns process-wide metrics followed by metrics for this node. */ + public Enumeration<MetricFamilySamples> getMetricFamilySamples() { + return new ConcatenatedEnumeration<>( + CollectorRegistry.defaultRegistry.metricFamilySamples(), + collectorRegistry.metricFamilySamples()); } @Override @@ -95,4 +104,24 @@ public class NodeExtension extends DefaultNodeExtension { public void printNodeInfo() { extCommon.printNodeInfo(systemLogger); } + + private static final class ConcatenatedEnumeration<T> implements Enumeration<T> { + private final Enumeration<T> first; + private final Enumeration<T> second; + + private ConcatenatedEnumeration(Enumeration<T> first, Enumeration<T> second) { + this.first = first; + this.second = second; + } + + @Override + public boolean hasMoreElements() { + return first.hasMoreElements() || second.hasMoreElements(); + } + + @Override + public T nextElement() { + return first.hasMoreElements() ? first.nextElement() : second.nextElement(); + } + } } diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestHttpGetCommandProcessor.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestHttpGetCommandProcessor.java index 60d40e6257..4d470cf430 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestHttpGetCommandProcessor.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestHttpGetCommandProcessor.java @@ -286,9 +286,7 @@ public class RestHttpGetCommandProcessor extends HttpCommandProcessor<HttpGetCom (NodeExtension) textCommandService.getNode().getNodeExtension(); try { TextFormat.writeFormat( - contentType, - stringWriter, - nodeExtension.getCollectorRegistry().metricFamilySamples()); + contentType, stringWriter, nodeExtension.getMetricFamilySamples()); this.prepareResponse(httpGetCommand, stringWriter.toString()); } catch (IOException e) { httpGetCommand.send400(); diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/MetricsServlet.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/MetricsServlet.java index 016d604067..50d7fb601a 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/MetricsServlet.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/MetricsServlet.java @@ -21,7 +21,6 @@ import org.apache.seatunnel.engine.server.NodeExtension; import org.apache.seatunnel.engine.server.rest.RestConstant; import com.hazelcast.spi.impl.NodeEngineImpl; -import io.prometheus.client.CollectorRegistry; import io.prometheus.client.exporter.common.TextFormat; import javax.servlet.ServletException; @@ -33,12 +32,11 @@ import java.io.StringWriter; public class MetricsServlet extends BaseServlet { - private final CollectorRegistry collectorRegistry; + private final NodeExtension nodeExtension; public MetricsServlet(NodeEngineImpl nodeEngine) { super(nodeEngine); - NodeExtension nodeExtension = (NodeExtension) nodeEngine.getNode().getNodeExtension(); - collectorRegistry = nodeExtension.getCollectorRegistry(); + nodeExtension = (NodeExtension) nodeEngine.getNode().getNodeExtension(); } @Override @@ -57,7 +55,7 @@ public class MetricsServlet extends BaseServlet { } try (StringWriter stringWriter = new StringWriter()) { TextFormat.writeFormat( - contentType, stringWriter, collectorRegistry.metricFamilySamples()); + contentType, stringWriter, nodeExtension.getMetricFamilySamples()); write(resp, stringWriter.toString()); } } diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/ExportsInstanceInitializer.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/ExportsInstanceInitializer.java index 920aa63036..1e17418e1d 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/ExportsInstanceInitializer.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/ExportsInstanceInitializer.java @@ -17,6 +17,7 @@ package org.apache.seatunnel.engine.server.telemetry.metrics; +import org.apache.seatunnel.engine.server.NodeExtension; import org.apache.seatunnel.engine.server.telemetry.metrics.exports.ClusterMetricExports; import org.apache.seatunnel.engine.server.telemetry.metrics.exports.EngineStateStoreLogicalMetricExports; import org.apache.seatunnel.engine.server.telemetry.metrics.exports.EngineStateStoreMetricExports; @@ -32,34 +33,32 @@ import io.prometheus.client.hotspot.DefaultExports; public final class ExportsInstanceInitializer { - private static boolean initialized = false; - private ExportsInstanceInitializer() {} - public static synchronized void init(Node node) { - if (!initialized) { - // initialize jvm collector - DefaultExports.initialize(); + /** Initializes process-wide JVM metrics and registers SeaTunnel metrics for the node. */ + public static void init(Node node) { + NodeExtension nodeExtension = (NodeExtension) node.getNodeExtension(); + CollectorRegistry collectorRegistry = nodeExtension.getCollectorRegistry(); + + // initialize process-wide JVM collectors once + DefaultExports.initialize(); - // register collectors - CollectorRegistry collectorRegistry = CollectorRegistry.defaultRegistry; - // Job info detail - new JobMetricExports(node).register(collectorRegistry); - // Thread pool status - new JobThreadPoolStatusExports(node).register(collectorRegistry); - // Node metrics - new NodeMetricExports(node).register(collectorRegistry); - // ReportMetricsOperation metrics - new ReportMetricsOperationExports(node).register(collectorRegistry); - // RequestSlotOperation metrics - new RequestSlotOperationExports(node).register(collectorRegistry); - // Engine state store metrics - new EngineStateStoreMetricExports(node).register(collectorRegistry); - // Engine state store logical metrics - new EngineStateStoreLogicalMetricExports(node).register(collectorRegistry); - // Cluster metrics - new ClusterMetricExports(node).register(collectorRegistry); - initialized = true; - } + // register collectors + // Job info detail + new JobMetricExports(node).register(collectorRegistry); + // Thread pool status + new JobThreadPoolStatusExports(node).register(collectorRegistry); + // Node metrics + new NodeMetricExports(node).register(collectorRegistry); + // ReportMetricsOperation metrics + new ReportMetricsOperationExports(node).register(collectorRegistry); + // RequestSlotOperation metrics + new RequestSlotOperationExports(node).register(collectorRegistry); + // Engine state store metrics + new EngineStateStoreMetricExports(node).register(collectorRegistry); + // Engine state store logical metrics + new EngineStateStoreLogicalMetricExports(node).register(collectorRegistry); + // Cluster metrics + new ClusterMetricExports(node).register(collectorRegistry); } } diff --git a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/metrics/MetricsApiTest.java b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/metrics/MetricsApiTest.java index 833c0a7a03..807f328567 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/metrics/MetricsApiTest.java +++ b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/metrics/MetricsApiTest.java @@ -74,16 +74,29 @@ public class MetricsApiTest { @BeforeAll public static void before() { + instance = createHazelcastInstance(); + } + + private static HazelcastInstanceImpl createHazelcastInstance() { SeaTunnelConfig seaTunnelConfig = ConfigProvider.locateAndGetSeaTunnelConfig(); seaTunnelConfig.getEngineConfig().getTelemetryConfig().getMetric().setEnabled(true); seaTunnelConfig.getEngineConfig().getHttpConfig().setEnabled(true); seaTunnelConfig.getEngineConfig().getHttpConfig().setPort(8080); seaTunnelConfig.getEngineConfig().setMode(ExecutionMode.LOCAL); - instance = SeaTunnelServerStarter.createHazelcastInstance(seaTunnelConfig); + return SeaTunnelServerStarter.createHazelcastInstance(seaTunnelConfig); } @Test public void metricsApiTest() { + assertMetricsApi(); + + instance.shutdown(); + instance = createHazelcastInstance(); + + assertMetricsApi(); + } + + private static void assertMetricsApi() { // The HTTP listener accepts requests as soon as the member reaches STARTED, but the // registered collectors read coordinator-owned state that is still being wired up at that // moment. Querying immediately made CI observe a transient 500 from a collector that ran
