From 9abf72efcb37dde7f4ed678ebd337c56c6ad7b92 Mon Sep 17 00:00:00 2001 From: Swati Gupta Date: Fri, 19 Jun 2026 00:16:02 +0530 Subject: [PATCH] [FLINK-39925][Autoscaler] Skip scaling decisions when vertex IO metrics are incomplete to prevent incorrect scale down --- .../autoscaler/ScalingMetricCollector.java | 28 ++-- .../ScalingMetricCollectorTest.java | 130 ++++++++++++++++++ 2 files changed, 148 insertions(+), 10 deletions(-) diff --git a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/ScalingMetricCollector.java b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/ScalingMetricCollector.java index d538ab63d9..70c12f201c 100644 --- a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/ScalingMetricCollector.java +++ b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/ScalingMetricCollector.java @@ -249,16 +249,24 @@ protected JobTopology getJobTopology(JobDetailsInfo jobDetailsInfo) { var metrics = new HashMap(); var finished = new HashSet(); - jobDetailsInfo - .getJobVertexInfos() - .forEach( - d -> { - if (d.getExecutionState() == ExecutionState.FINISHED) { - finished.add(d.getJobVertexID()); - } - metrics.put( - d.getJobVertexID(), IOMetrics.from(d.getJobVertexMetrics())); - }); + for (var d : jobDetailsInfo.getJobVertexInfos()) { + if (d.getExecutionState() == ExecutionState.FINISHED) { + finished.add(d.getJobVertexID()); + } + var ioMetricsInfo = d.getJobVertexMetrics(); + if (!ioMetricsInfo.isRecordsReadComplete() + || !ioMetricsInfo.isRecordsWrittenComplete()) { + throw new NotReadyException( + String.format( + "Vertex %s has incomplete IO metrics (read-records-complete=%s," + + " write-records-complete=%s). Skipping metric" + + " collection until metrics are complete.", + d.getJobVertexID(), + ioMetricsInfo.isRecordsReadComplete(), + ioMetricsInfo.isRecordsWrittenComplete())); + } + metrics.put(d.getJobVertexID(), IOMetrics.from(ioMetricsInfo)); + } return JobTopology.fromJsonPlan( json, slotSharingGroupIdMap, maxParallelismMap, metrics, finished); diff --git a/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/ScalingMetricCollectorTest.java b/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/ScalingMetricCollectorTest.java index 70e3f539f6..1f21344cee 100644 --- a/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/ScalingMetricCollectorTest.java +++ b/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/ScalingMetricCollectorTest.java @@ -354,6 +354,136 @@ public void testJobTopologyParsingThrowsNotReadyException() throws Exception { NotReadyException.class, () -> metricsCollector.getJobTopology(jobDetailsInfo)); } + @Test + public void testIncompleteIoMetricsThrowsNotReadyException() throws Exception { + String s = + "{\n" + + " \"jid\": \"bb8f15efbb37f2ce519f55cdc0e049bf\",\n" + + " \"name\": \"State machine job\",\n" + + " \"isStoppable\": false,\n" + + " \"state\": \"RUNNING\",\n" + + " \"start-time\": 1707893512027,\n" + + " \"end-time\": -1,\n" + + " \"duration\": 214716,\n" + + " \"maxParallelism\": -1,\n" + + " \"now\": 1707893726743,\n" + + " \"timestamps\": {\n" + + " \"SUSPENDED\": 0,\n" + + " \"CREATED\": 1707893512139,\n" + + " \"FAILING\": 0,\n" + + " \"FAILED\": 0,\n" + + " \"INITIALIZING\": 1707893512027,\n" + + " \"RECONCILING\": 0,\n" + + " \"RUNNING\": 1707893512217,\n" + + " \"RESTARTING\": 0,\n" + + " \"CANCELLING\": 0,\n" + + " \"FINISHED\": 0,\n" + + " \"CANCELED\": 0\n" + + " },\n" + + " \"vertices\": [\n" + + " {\n" + + " \"id\": \"bc764cd8ddf7a0cff126f51c16239658\",\n" + + " \"name\": \"Source: Custom Source\",\n" + + " \"maxParallelism\": 128,\n" + + " \"parallelism\": 2,\n" + + " \"status\": \"RUNNING\",\n" + + " \"start-time\": 1707893517277,\n" + + " \"end-time\": -1,\n" + + " \"duration\": 209466,\n" + + " \"tasks\": {\n" + + " \"DEPLOYING\": 0,\n" + + " \"INITIALIZING\": 0,\n" + + " \"SCHEDULED\": 0,\n" + + " \"CANCELING\": 0,\n" + + " \"CANCELED\": 0,\n" + + " \"RECONCILING\": 0,\n" + + " \"RUNNING\": 2,\n" + + " \"FAILED\": 0,\n" + + " \"CREATED\": 0,\n" + + " \"FINISHED\": 0\n" + + " },\n" + + " \"metrics\": {\n" + + " \"read-bytes\": 0,\n" + + " \"read-bytes-complete\": true,\n" + + " \"write-bytes\": 0,\n" + + " \"write-bytes-complete\": false,\n" + + " \"read-records\": 0,\n" + + " \"read-records-complete\": true,\n" + + " \"write-records\": 0,\n" + + " \"write-records-complete\": false,\n" + + " \"accumulated-backpressured-time\": 0,\n" + + " \"accumulated-idle-time\": 0,\n" + + " \"accumulated-busy-time\": 0\n" + + " }\n" + + " },\n" + + " {\n" + + " \"id\": \"20ba6b65f97481d5570070de90e4e791\",\n" + + " \"name\": \"Flat Map -> Sink: Print to Std. Out\",\n" + + " \"maxParallelism\": 128,\n" + + " \"parallelism\": 2,\n" + + " \"status\": \"RUNNING\",\n" + + " \"start-time\": 1707893517280,\n" + + " \"end-time\": -1,\n" + + " \"duration\": 209463,\n" + + " \"tasks\": {\n" + + " \"DEPLOYING\": 0,\n" + + " \"INITIALIZING\": 0,\n" + + " \"SCHEDULED\": 0,\n" + + " \"CANCELING\": 0,\n" + + " \"CANCELED\": 0,\n" + + " \"RECONCILING\": 0,\n" + + " \"RUNNING\": 2,\n" + + " \"FAILED\": 0,\n" + + " \"CREATED\": 0,\n" + + " \"FINISHED\": 0\n" + + " },\n" + + " \"metrics\": {\n" + + " \"read-bytes\": 0,\n" + + " \"read-bytes-complete\": false,\n" + + " \"write-bytes\": 0,\n" + + " \"write-bytes-complete\": true,\n" + + " \"read-records\": 0,\n" + + " \"read-records-complete\": false,\n" + + " \"write-records\": 0,\n" + + " \"write-records-complete\": true,\n" + + " \"accumulated-backpressured-time\": 0,\n" + + " \"accumulated-idle-time\": 0,\n" + + " \"accumulated-busy-time\": 1\n" + + " }\n" + + " }\n" + + " ],\n" + + " \"status-counts\": {\n" + + " \"DEPLOYING\": 0,\n" + + " \"INITIALIZING\": 0,\n" + + " \"SCHEDULED\": 0,\n" + + " \"CANCELING\": 0,\n" + + " \"CANCELED\": 0,\n" + + " \"RECONCILING\": 0,\n" + + " \"RUNNING\": 2,\n" + + " \"FAILED\": 0,\n" + + " \"CREATED\": 0,\n" + + " \"FINISHED\": 0\n" + + " },\n" + + " \"plan\": {\n" + + " \"jid\": \"bb8f15efbb37f2ce519f55cdc0e049bf\",\n" + + " \"name\": \"State machine job\",\n" + + " \"type\": \"STREAMING\",\n" + + " \"nodes\": [\n" + + " ]\n" + + " }\n" + + "}"; + JobDetailsInfo jobDetailsInfo = new ObjectMapper().readValue(s, JobDetailsInfo.class); + + var metricsCollector = new RestApiMetricsCollector<>(); + var ex = + assertThrows( + NotReadyException.class, + () -> metricsCollector.getJobTopology(jobDetailsInfo)); + assertTrue( + ex.getMessage().contains("incomplete IO metrics"), + "Exception message should indicate incomplete IO metrics, got: " + ex.getMessage()); + } + @Test public void testJobTopologyParsingFromJobDetailsWithSlotSharingGroup() throws Exception { String s =