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-11991-f0046a0022028f0b511b2920a8a6cb7ded7d0d8b in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit b4158f01ddf9c3a1dda3a48a559d4ae6aec557d2 Author: hutiefang76 <[email protected]> AuthorDate: Mon Aug 31 17:25:30 2026 +0000 [Fix][Zeta] Avoid metrics failures during coordinator startup (#11991) Signed-off-by: hutiefang76 <[email protected]> --- .../telemetry/metrics/AbstractCollector.java | 19 +++++++ .../metrics/exports/JobMetricExports.java | 5 +- .../exports/RequestSlotOperationExports.java | 8 ++- .../TelemetryCollectorCoordinatorGuardTest.java | 62 ++++++++++++++++++++++ 4 files changed, 92 insertions(+), 2 deletions(-) diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/AbstractCollector.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/AbstractCollector.java index 3ef0203385..92ae69ddc6 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/AbstractCollector.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/AbstractCollector.java @@ -19,6 +19,7 @@ package org.apache.seatunnel.engine.server.telemetry.metrics; import org.apache.seatunnel.shade.com.google.common.collect.Lists; +import org.apache.seatunnel.engine.common.exception.SeaTunnelEngineException; import org.apache.seatunnel.engine.server.CoordinatorService; import org.apache.seatunnel.engine.server.SeaTunnelServer; @@ -69,6 +70,24 @@ public abstract class AbstractCollector extends Collector { return getServer().getCoordinatorService(); } + /** + * Gets the coordinator after a readiness check without failing a metrics scrape during a master + * transition. + * + * @return the active coordinator, or {@code null} when the coordinator is unavailable + */ + protected CoordinatorService getReadyCoordinatorService() { + try { + return getCoordinatorService(); + } catch (SeaTunnelEngineException e) { + getLogger(getClass()) + .fine( + "Coordinator service is unavailable while collecting metrics; skipping this scrape", + e); + return null; + } + } + // Non-blocking coordinator readiness check; call before getCoordinatorService() in collect(). protected boolean isCoordinatorReady() { return getServer().isCoordinatorActive(); diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/JobMetricExports.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/JobMetricExports.java index b197e29bb1..7037487fe2 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/JobMetricExports.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/JobMetricExports.java @@ -38,7 +38,10 @@ public class JobMetricExports extends AbstractCollector { List<MetricFamilySamples> mfs = new ArrayList(); // Report metrics only when the local node is ACTIVE and the master is available. if (isMaster() && isCoordinatorReady()) { - CoordinatorService coordinatorService = getCoordinatorService(); + CoordinatorService coordinatorService = getReadyCoordinatorService(); + if (coordinatorService == null) { + return mfs; + } JobCounter jobCountMetrics = coordinatorService.getJobCountMetrics(); GaugeMetricFamily metricFamily = diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/RequestSlotOperationExports.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/RequestSlotOperationExports.java index 5e36e560a6..b7a9db57d8 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/RequestSlotOperationExports.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/RequestSlotOperationExports.java @@ -17,6 +17,7 @@ package org.apache.seatunnel.engine.server.telemetry.metrics.exports; +import org.apache.seatunnel.engine.server.CoordinatorService; import org.apache.seatunnel.engine.server.resourcemanager.ResourceManager; import org.apache.seatunnel.engine.server.telemetry.metrics.AbstractCollector; import org.apache.seatunnel.engine.server.telemetry.metrics.entity.RequestSlotOperationStats; @@ -41,7 +42,12 @@ public class RequestSlotOperationExports extends AbstractCollector { return mfs; } - ResourceManager resourceManager = getCoordinatorService().getInitializedResourceManager(); + CoordinatorService coordinatorService = getReadyCoordinatorService(); + if (coordinatorService == null) { + return mfs; + } + + ResourceManager resourceManager = coordinatorService.getInitializedResourceManager(); if (resourceManager == null) { return mfs; } diff --git a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/telemetry/metrics/TelemetryCollectorCoordinatorGuardTest.java b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/telemetry/metrics/TelemetryCollectorCoordinatorGuardTest.java index c4e603d042..48342bed86 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/telemetry/metrics/TelemetryCollectorCoordinatorGuardTest.java +++ b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/telemetry/metrics/TelemetryCollectorCoordinatorGuardTest.java @@ -17,6 +17,8 @@ package org.apache.seatunnel.engine.server.telemetry.metrics; +import org.apache.seatunnel.engine.common.exception.SeaTunnelEngineException; +import org.apache.seatunnel.engine.common.exception.SeaTunnelEngineRetryableException; import org.apache.seatunnel.engine.server.CoordinatorService; import org.apache.seatunnel.engine.server.SeaTunnelServer; import org.apache.seatunnel.engine.server.TaskExecutionService; @@ -120,6 +122,36 @@ public class TelemetryCollectorCoordinatorGuardTest { Mockito.verify(mockServer, Mockito.never()).getCoordinatorService(); } + @Test + void testJobMetricExportsReturnsEmptyWhenCoordinatorBecomesUnavailable() { + Mockito.when(mockNode.isMaster()).thenReturn(true); + Mockito.when(mockServer.isCoordinatorActive()).thenReturn(true); + Mockito.when(mockServer.getCoordinatorService()) + .thenThrow(new SeaTunnelEngineRetryableException("coordinator is starting")); + + JobMetricExports exports = new JobMetricExports(mockNode); + List<Collector.MetricFamilySamples> result = exports.collect(); + + Assertions.assertTrue( + result.isEmpty(), + "collect() must not fail a metrics scrape when the coordinator changes state"); + } + + @Test + void testJobMetricExportsReturnsEmptyWhenNodeLosesMasterRole() { + Mockito.when(mockNode.isMaster()).thenReturn(true); + Mockito.when(mockServer.isCoordinatorActive()).thenReturn(true); + Mockito.when(mockServer.getCoordinatorService()) + .thenThrow(new SeaTunnelEngineException("This is not a master node now.")); + + JobMetricExports exports = new JobMetricExports(mockNode); + List<Collector.MetricFamilySamples> result = exports.collect(); + + Assertions.assertTrue( + result.isEmpty(), + "collect() must not fail a metrics scrape when the master changes"); + } + @Test void testJobMetricExportsReturnsMetricsWhenCoordinatorReady() { Mockito.when(mockNode.isMaster()).thenReturn(true); @@ -283,6 +315,36 @@ public class TelemetryCollectorCoordinatorGuardTest { Mockito.verify(mockServer, Mockito.never()).getCoordinatorService(); } + @Test + void testRequestSlotOperationExportsReturnsEmptyWhenCoordinatorBecomesUnavailable() { + Mockito.when(mockNode.isMaster()).thenReturn(true); + Mockito.when(mockServer.isCoordinatorActive()).thenReturn(true); + Mockito.when(mockServer.getCoordinatorService()) + .thenThrow(new SeaTunnelEngineRetryableException("coordinator is starting")); + + RequestSlotOperationExports exports = new RequestSlotOperationExports(mockNode); + List<Collector.MetricFamilySamples> result = exports.collect(); + + Assertions.assertTrue( + result.isEmpty(), + "collect() must not fail a metrics scrape when the coordinator changes state"); + } + + @Test + void testRequestSlotOperationExportsReturnsEmptyWhenNodeLosesMasterRole() { + Mockito.when(mockNode.isMaster()).thenReturn(true); + Mockito.when(mockServer.isCoordinatorActive()).thenReturn(true); + Mockito.when(mockServer.getCoordinatorService()) + .thenThrow(new SeaTunnelEngineException("This is not a master node now.")); + + RequestSlotOperationExports exports = new RequestSlotOperationExports(mockNode); + List<Collector.MetricFamilySamples> result = exports.collect(); + + Assertions.assertTrue( + result.isEmpty(), + "collect() must not fail a metrics scrape when the master changes"); + } + @Test void testRequestSlotOperationExportsDoesNotInitializeResourceManager() { Mockito.when(mockNode.isMaster()).thenReturn(true);
