This is an automated email from the ASF dual-hosted git repository.

zhangshenghang pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new 033612bfd6 [Fix][E2E] Address follow-up review issues for CI 
stabilization (#11097)
033612bfd6 is described below

commit 033612bfd6e59713dda262032031b6aeca927318
Author: David Zollo <[email protected]>
AuthorDate: Sun Aug 23 22:25:35 2026 +0800

    [Fix][E2E] Address follow-up review issues for CI stabilization (#11097)
    
    Co-authored-by: Daniel <[email protected]>
    Co-authored-by: DanielLeens <[email protected]>
    Co-authored-by: Daniel <[email protected]>
    Co-authored-by: zhiwei.niu <[email protected]>
    Co-authored-by: danielnadean 
<[email protected]>
---
 .../engine/e2e/ClusterFailureNoRestoreIT.java      | 35 ++++++++++++++++++----
 .../checkpoint/CheckpointErrorRestoreEndTest.java  |  9 ++++--
 .../engine/server/event/JobStateEventTest.java     |  9 ++----
 3 files changed, 38 insertions(+), 15 deletions(-)

diff --git 
a/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/ClusterFailureNoRestoreIT.java
 
b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/ClusterFailureNoRestoreIT.java
index b0ddf6a0e4..21b07c20f3 100644
--- 
a/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/ClusterFailureNoRestoreIT.java
+++ 
b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/ClusterFailureNoRestoreIT.java
@@ -23,6 +23,7 @@ import org.apache.seatunnel.common.utils.FileUtils;
 import org.apache.seatunnel.engine.client.SeaTunnelClient;
 import org.apache.seatunnel.engine.client.job.ClientJobExecutionEnvironment;
 import org.apache.seatunnel.engine.client.job.ClientJobProxy;
+import 
org.apache.seatunnel.engine.client.job.JobMetricsRunner.JobMetricsSummary;
 import org.apache.seatunnel.engine.common.config.ConfigProvider;
 import org.apache.seatunnel.engine.common.config.JobConfig;
 import org.apache.seatunnel.engine.common.config.SeaTunnelConfig;
@@ -66,11 +67,10 @@ public class ClusterFailureNoRestoreIT {
     private static final long PRE_SHUTDOWN_RUNNING_TIMEOUT_SECONDS = 30L;
 
     /**
-     * Give the running batch topology a short warm-up window before shutting 
down a worker. The
-     * LocalFile sink used by this test commits files transactionally, so 
intermediate file lines
-     * are not a reliable progress signal while the job is still running.
+     * Wait for job metrics to show that the bounded fake source has started 
consuming records but
+     * has not yet drained its full input before shutting a worker down.
      */
-    private static final long PRE_SHUTDOWN_RUNNING_GRACE_SECONDS = 5L;
+    private static final long PRE_SHUTDOWN_PROGRESS_TIMEOUT_SECONDS = 30L;
 
     private static final String DYNAMIC_TEST_CASE_NAME = 
"dynamic_test_case_name";
 
@@ -86,6 +86,7 @@ public class ClusterFailureNoRestoreIT {
         String testClusterName = "ClusterFailureNoRestoreIT_batch_no_restore";
         long testRowNumber = NO_RESTORE_BATCH_ROW_NUM_PER_PARALLELISM;
         int testParallelism = 6;
+        long expectedSourceReadCount = testRowNumber * testParallelism;
 
         HazelcastInstanceImpl node1 = null;
         HazelcastInstanceImpl node2 = null;
@@ -129,8 +130,8 @@ public class ClusterFailureNoRestoreIT {
                             () ->
                                     Assertions.assertEquals(
                                             JobStatus.RUNNING, 
clientJobProxy.getJobStatus()));
-
-            TimeUnit.SECONDS.sleep(PRE_SHUTDOWN_RUNNING_GRACE_SECONDS);
+            awaitSourceProgressBeforeShutdown(
+                    engineClient, clientJobProxy.getJobId(), 
expectedSourceReadCount);
 
             CompletableFuture<JobResult> waitForCompleteFuture =
                     
CompletableFuture.supplyAsync(clientJobProxy::waitForJobCompleteV2);
@@ -176,6 +177,28 @@ public class ClusterFailureNoRestoreIT {
         }
     }
 
+    /**
+     * Wait until the runtime metrics confirm that the batch source is 
actively producing data and
+     * has not already consumed the full bounded input before the worker 
shutdown is injected.
+     */
+    private void awaitSourceProgressBeforeShutdown(
+            SeaTunnelClient engineClient, long jobId, long 
expectedSourceReadCount) {
+        Awaitility.await()
+                .atMost(PRE_SHUTDOWN_PROGRESS_TIMEOUT_SECONDS, 
TimeUnit.SECONDS)
+                .untilAsserted(
+                        () -> {
+                            JobMetricsSummary jobMetricsSummary =
+                                    engineClient.getJobMetricsSummary(jobId);
+                            Assertions.assertTrue(
+                                    jobMetricsSummary.getSourceReadCount() > 0,
+                                    "Expected the fake source to emit records 
before worker shutdown");
+                            Assertions.assertTrue(
+                                    jobMetricsSummary.getSourceReadCount()
+                                            < expectedSourceReadCount,
+                                    "The bounded fake source finished before 
worker shutdown was injected");
+                        });
+    }
+
     private ImmutablePair<String, String> createTestResources(
             @NonNull String testCaseName, long rowNumber, int parallelism) 
throws IOException {
         checkArgument(rowNumber > 0, "rowNumber must greater than 0");
diff --git 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointErrorRestoreEndTest.java
 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointErrorRestoreEndTest.java
index fbdb09e561..2d9ca2d64f 100644
--- 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointErrorRestoreEndTest.java
+++ 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointErrorRestoreEndTest.java
@@ -35,6 +35,11 @@ public class CheckpointErrorRestoreEndTest
         extends AbstractSeaTunnelServerTest<CheckpointErrorRestoreEndTest> {
     public static String STREAM_CONF_WITH_ERROR_PATH =
             "batch_fakesource_to_inmemory_with_commit_error.conf";
+    /**
+     * The commit-error batch job restores three times before reaching the 
final FAILED state, so
+     * the related regression tests share the same 240-second upper bound.
+     */
+    public static final long RESTORE_TO_FAILED_TIMEOUT_SECONDS = 240L;
 
     @Test
     public void testCheckpointRestoreToFailEnd() {
@@ -43,7 +48,7 @@ public class CheckpointErrorRestoreEndTest
 
         JobMaster jobMaster = 
server.getCoordinatorService().getJobMaster(jobId);
         Assertions.assertEquals(1, 
jobMaster.getPhysicalPlan().getPipelineList().size());
-        await().atMost(240, TimeUnit.SECONDS)
+        await().atMost(RESTORE_TO_FAILED_TIMEOUT_SECONDS, TimeUnit.SECONDS)
                 .untilAsserted(
                         () ->
                                 Assertions.assertEquals(
@@ -53,7 +58,7 @@ public class CheckpointErrorRestoreEndTest
                                                 .getPipelineList()
                                                 .get(0)
                                                 .getPipelineRestoreNum()));
-        await().atMost(240, TimeUnit.SECONDS)
+        await().atMost(RESTORE_TO_FAILED_TIMEOUT_SECONDS, TimeUnit.SECONDS)
                 .untilAsserted(
                         () ->
                                 Assertions.assertEquals(
diff --git 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/event/JobStateEventTest.java
 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/event/JobStateEventTest.java
index 00e4b730ab..1a87671bbf 100644
--- 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/event/JobStateEventTest.java
+++ 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/event/JobStateEventTest.java
@@ -33,17 +33,12 @@ import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 import java.util.concurrent.atomic.AtomicReference;
 
+import static 
org.apache.seatunnel.engine.server.checkpoint.CheckpointErrorRestoreEndTest.RESTORE_TO_FAILED_TIMEOUT_SECONDS;
 import static 
org.apache.seatunnel.engine.server.checkpoint.CheckpointErrorRestoreEndTest.STREAM_CONF_WITH_ERROR_PATH;
 import static org.awaitility.Awaitility.await;
 
 class JobStateEventTest extends AbstractSeaTunnelServerTest {
 
-    /**
-     * The commit-error batch job may restore several times before reaching 
the final FAILED state,
-     * so this event test needs the same longer timeout as the checkpoint 
restore regression test.
-     */
-    private static final long FAILED_JOB_EVENT_TIMEOUT_SECONDS = 240L;
-
     @Test
     void testJobStateEvent() {
 
@@ -96,7 +91,7 @@ class JobStateEventTest extends AbstractSeaTunnelServerTest {
 
         long jobIdFailed = System.currentTimeMillis();
         startJob(jobIdFailed, STREAM_CONF_WITH_ERROR_PATH, false);
-        await().atMost(FAILED_JOB_EVENT_TIMEOUT_SECONDS, TimeUnit.SECONDS)
+        await().atMost(RESTORE_TO_FAILED_TIMEOUT_SECONDS, TimeUnit.SECONDS)
                 .untilAsserted(
                         () ->
                                 Assertions.assertEquals(

Reply via email to