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(