This is an automated email from the ASF dual-hosted git repository.
davidzollo 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 f35a56d0ca [Improve][Engine] fix engine UT flaky test (#11000)
f35a56d0ca is described below
commit f35a56d0ca02a04a0f75ca66c3be730a3ae9de91
Author: zhiwei.niu <[email protected]>
AuthorDate: Sat Jun 13 17:27:04 2026 +0800
[Improve][Engine] fix engine UT flaky test (#11000)
---
.../seatunnel/api/tracing/MDCTracerTest.java | 247 +++++++++++----------
.../command/ServerExecuteCommandTest.java | 61 +++--
.../client/SeaTunnelEngineClusterRoleTest.java | 45 ++--
.../engine/server/CoordinatorServiceTest.java | 141 ++++++------
...CoordinatorServiceWithCancelPendingJobTest.java | 16 +-
.../server/checkpoint/CheckpointStorageTest.java | 21 +-
.../server/checkpoint/CheckpointTimeOutTest.java | 2 +-
.../master/FollowerRunningJobsFilterTest.java | 4 +-
.../server/master/JobHistoryServiceTest.java | 3 +-
.../engine/server/rest/BaseServletTest.java | 3 +-
.../engine/server/rest/RestApiHttpBasicTest.java | 4 +-
.../server/rest/RestApiHttpsForTruststoreTest.java | 3 +-
.../engine/server/rest/RestApiHttpsTest.java | 13 +-
.../exports/EngineStateStoreMetricExportsTest.java | 4 +-
14 files changed, 322 insertions(+), 245 deletions(-)
diff --git
a/seatunnel-api/src/test/java/org/apache/seatunnel/api/tracing/MDCTracerTest.java
b/seatunnel-api/src/test/java/org/apache/seatunnel/api/tracing/MDCTracerTest.java
index 694bac8293..63c8154865 100644
---
a/seatunnel-api/src/test/java/org/apache/seatunnel/api/tracing/MDCTracerTest.java
+++
b/seatunnel-api/src/test/java/org/apache/seatunnel/api/tracing/MDCTracerTest.java
@@ -23,7 +23,9 @@ import org.slf4j.MDC;
import java.util.concurrent.Callable;
import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Supplier;
@@ -127,129 +129,138 @@ public class MDCTracerTest {
public void testMDCTracedExecutorService() throws Exception {
MDCContext mdcContext = MDCContext.of(1, 2, 3);
- MDCExecutorService tracedExecutorService =
- MDCTracer.tracing(mdcContext,
Executors.newSingleThreadExecutor());
-
- Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
- Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
- tracedExecutorService
- .submit(
- new Runnable() {
- @Override
- public void run() {
- Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
- Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
- }
- })
- .get();
- Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
- Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
-
- tracedExecutorService
- .submit(
- new Callable<Void>() {
- @Override
- public Void call() throws Exception {
- Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
- Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
- return null;
- }
- })
- .get();
- Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
- Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
-
- MDCScheduledExecutorService tracedScheduledExecutorService =
- MDCTracer.tracing(mdcContext,
Executors.newSingleThreadScheduledExecutor());
- Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
- Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
+ ExecutorService rawExecutor = Executors.newSingleThreadExecutor();
+ ScheduledExecutorService rawScheduledExecutor =
+ Executors.newSingleThreadScheduledExecutor();
+ try {
+ MDCExecutorService tracedExecutorService =
MDCTracer.tracing(mdcContext, rawExecutor);
+
+ Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
+ Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
+ tracedExecutorService
+ .submit(
+ new Runnable() {
+ @Override
+ public void run() {
+ Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
+ Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
+ }
+ })
+ .get();
+ Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
+ Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
+
+ tracedExecutorService
+ .submit(
+ new Callable<Void>() {
+ @Override
+ public Void call() throws Exception {
+ Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
+ Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
+ return null;
+ }
+ })
+ .get();
+ Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
+ Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
+
+ MDCScheduledExecutorService tracedScheduledExecutorService =
+ MDCTracer.tracing(mdcContext, rawScheduledExecutor);
+ Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
+ Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
+
+ tracedScheduledExecutorService
+ .schedule(
+ new Runnable() {
+ @Override
+ public void run() {
+ Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
+ Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
+ }
+ },
+ 1,
+ TimeUnit.SECONDS)
+ .get();
+ Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
+ Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
+
+ tracedScheduledExecutorService
+ .schedule(
+ new Callable<Object>() {
+ @Override
+ public Object call() {
+ Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
+ Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
+ return null;
+ }
+ },
+ 1,
+ TimeUnit.SECONDS)
+ .get();
+ Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
+ Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
- tracedScheduledExecutorService
- .schedule(
- new Runnable() {
- @Override
- public void run() {
- Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
- Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
- }
- },
- 1,
- TimeUnit.SECONDS)
- .get();
- Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
- Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
+ CompletableFuture<Boolean> futureWithScheduleAtFixedRate = new
CompletableFuture<>();
+ tracedScheduledExecutorService.scheduleAtFixedRate(
+ new Runnable() {
+ AtomicInteger executeCount = new AtomicInteger(0);
- tracedScheduledExecutorService
- .schedule(
- new Callable<Object>() {
- @Override
- public Object call() {
- Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
- Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
- return null;
+ @Override
+ public void run() {
+ Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
+ Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
+ executeCount.incrementAndGet();
+ if (executeCount.get() > 10
+ &&
!futureWithScheduleAtFixedRate.isDone()) {
+ futureWithScheduleAtFixedRate.complete(true);
}
- },
- 1,
- TimeUnit.SECONDS)
- .get();
- Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
- Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
-
- CompletableFuture<Boolean> futureWithScheduleAtFixedRate = new
CompletableFuture<>();
- tracedScheduledExecutorService.scheduleAtFixedRate(
- new Runnable() {
- AtomicInteger executeCount = new AtomicInteger(0);
-
- @Override
- public void run() {
- Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
- Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
- executeCount.incrementAndGet();
- if (executeCount.get() > 10 &&
!futureWithScheduleAtFixedRate.isDone()) {
- futureWithScheduleAtFixedRate.complete(true);
}
- }
- },
- 0,
- 10,
- TimeUnit.MILLISECONDS);
- futureWithScheduleAtFixedRate.join();
-
- CompletableFuture<Boolean> futureWithScheduleAtFixedDelay = new
CompletableFuture<>();
- tracedScheduledExecutorService.scheduleWithFixedDelay(
- new Runnable() {
- AtomicInteger executeCount = new AtomicInteger(0);
-
- @Override
- public void run() {
- Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
- Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
- executeCount.incrementAndGet();
- if (executeCount.get() > 10 &&
!futureWithScheduleAtFixedDelay.isDone()) {
- futureWithScheduleAtFixedDelay.complete(true);
+ },
+ 0,
+ 10,
+ TimeUnit.MILLISECONDS);
+ futureWithScheduleAtFixedRate.get(30, TimeUnit.SECONDS);
+
+ CompletableFuture<Boolean> futureWithScheduleAtFixedDelay = new
CompletableFuture<>();
+ tracedScheduledExecutorService.scheduleWithFixedDelay(
+ new Runnable() {
+ AtomicInteger executeCount = new AtomicInteger(0);
+
+ @Override
+ public void run() {
+ Assertions.assertEquals("1",
MDC.get(MDCContext.JOB_ID));
+ Assertions.assertEquals("2",
MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertEquals("3",
MDC.get(MDCContext.TASK_ID));
+ executeCount.incrementAndGet();
+ if (executeCount.get() > 10
+ &&
!futureWithScheduleAtFixedDelay.isDone()) {
+ futureWithScheduleAtFixedDelay.complete(true);
+ }
}
- }
- },
- 0,
- 10,
- TimeUnit.MILLISECONDS);
- futureWithScheduleAtFixedDelay.join();
-
- Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
- Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
- Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
+ },
+ 0,
+ 10,
+ TimeUnit.MILLISECONDS);
+ futureWithScheduleAtFixedDelay.get(30, TimeUnit.SECONDS);
+
+ Assertions.assertNull(MDC.get(MDCContext.JOB_ID));
+ Assertions.assertNull(MDC.get(MDCContext.PIPELINE_ID));
+ Assertions.assertNull(MDC.get(MDCContext.TASK_ID));
+ } finally {
+ rawExecutor.shutdownNow();
+ rawScheduledExecutor.shutdownNow();
+ }
}
@Test
diff --git
a/seatunnel-core/seatunnel-starter/src/test/java/org/apache/seatunnel/core/starter/seatunnel/command/ServerExecuteCommandTest.java
b/seatunnel-core/seatunnel-starter/src/test/java/org/apache/seatunnel/core/starter/seatunnel/command/ServerExecuteCommandTest.java
index 16f76aef74..78b7d06b7f 100644
---
a/seatunnel-core/seatunnel-starter/src/test/java/org/apache/seatunnel/core/starter/seatunnel/command/ServerExecuteCommandTest.java
+++
b/seatunnel-core/seatunnel-starter/src/test/java/org/apache/seatunnel/core/starter/seatunnel/command/ServerExecuteCommandTest.java
@@ -28,7 +28,10 @@ import org.junit.jupiter.api.condition.DisabledOnJre;
import org.junit.jupiter.api.condition.JRE;
import com.hazelcast.cluster.Member;
+import com.hazelcast.instance.impl.HazelcastInstanceImpl;
+import java.util.ArrayList;
+import java.util.List;
import java.util.Set;
public class ServerExecuteCommandTest {
@@ -37,33 +40,57 @@ public class ServerExecuteCommandTest {
@DisabledOnJre(value = JRE.JAVA_11, disabledReason = "the test case only
works on Java 8")
public void testJavaVersionCheck() {
String realVersion = System.getProperty("java.version");
- System.setProperty("java.version", "1.8.0_191");
-
Assertions.assertFalse(ServerExecuteCommand.isAllocatingThreadGetName());
- System.setProperty("java.version", "1.8.0_60");
-
Assertions.assertTrue(ServerExecuteCommand.isAllocatingThreadGetName());
- System.setProperty("java.version", realVersion);
+ try {
+ System.setProperty("java.version", "1.8.0_191");
+
Assertions.assertFalse(ServerExecuteCommand.isAllocatingThreadGetName());
+ System.setProperty("java.version", "1.8.0_60");
+
Assertions.assertTrue(ServerExecuteCommand.isAllocatingThreadGetName());
+ } finally {
+ System.setProperty("java.version", realVersion);
+ }
}
@Test
- public void testMemberList() {
+ public void testMemberList() throws InterruptedException {
String clusterName = getClusterName("ServerExecuteCommandTest");
SeaTunnelConfig seaTunnelConfig =
ConfigProvider.locateAndGetSeaTunnelConfig();
seaTunnelConfig.getHazelcastConfig().setClusterName(clusterName);
seaTunnelConfig.getEngineConfig().getHttpConfig().setEnableDynamicPort(true);
- SeaTunnelServerStarter.createMasterHazelcastInstance(seaTunnelConfig);
- SeaTunnelServerStarter.createMasterHazelcastInstance(seaTunnelConfig);
- SeaTunnelServerStarter.createWorkerHazelcastInstance(seaTunnelConfig);
- SeaTunnelServerStarter.createWorkerHazelcastInstance(seaTunnelConfig);
- SeaTunnelServerStarter.createWorkerHazelcastInstance(seaTunnelConfig);
+ List<HazelcastInstanceImpl> instances = new ArrayList<>();
+ try {
+
instances.add(SeaTunnelServerStarter.createMasterHazelcastInstance(seaTunnelConfig));
+
instances.add(SeaTunnelServerStarter.createMasterHazelcastInstance(seaTunnelConfig));
+
instances.add(SeaTunnelServerStarter.createWorkerHazelcastInstance(seaTunnelConfig));
+
instances.add(SeaTunnelServerStarter.createWorkerHazelcastInstance(seaTunnelConfig));
+
instances.add(SeaTunnelServerStarter.createWorkerHazelcastInstance(seaTunnelConfig));
- ServerCommandArgs serverCommandArgs = new ServerCommandArgs();
- serverCommandArgs.setClusterName(clusterName);
- serverCommandArgs.setShowClusterMembers(true);
+ HazelcastInstanceImpl firstInstance = instances.get(0);
+ long deadline = System.currentTimeMillis() + 30_000;
+ while (firstInstance.getCluster().getMembers().size() < 5) {
+ if (System.currentTimeMillis() > deadline) {
+ Assertions.fail(
+ "Cluster did not form within 30s, members: "
+ +
firstInstance.getCluster().getMembers().size());
+ }
+ Thread.sleep(500);
+ }
- ServerExecuteCommand serverExecuteCommand = new
ServerExecuteCommand(serverCommandArgs);
- Set<Member> members = serverExecuteCommand.showClusterMembers();
- Assertions.assertEquals(5, members.size());
+ ServerCommandArgs serverCommandArgs = new ServerCommandArgs();
+ serverCommandArgs.setClusterName(clusterName);
+ serverCommandArgs.setShowClusterMembers(true);
+
+ ServerExecuteCommand serverExecuteCommand = new
ServerExecuteCommand(serverCommandArgs);
+ Set<Member> members = serverExecuteCommand.showClusterMembers();
+ Assertions.assertEquals(5, members.size());
+ } finally {
+ for (HazelcastInstanceImpl inst : instances) {
+ try {
+ inst.shutdown();
+ } catch (Exception ignored) {
+ }
+ }
+ }
}
public static String getClusterName(String testClassName) {
diff --git
a/seatunnel-engine/seatunnel-engine-client/src/test/java/org/apache/seatunnel/engine/client/SeaTunnelEngineClusterRoleTest.java
b/seatunnel-engine/seatunnel-engine-client/src/test/java/org/apache/seatunnel/engine/client/SeaTunnelEngineClusterRoleTest.java
index 0dee6aac4c..4da8347206 100644
---
a/seatunnel-engine/seatunnel-engine-client/src/test/java/org/apache/seatunnel/engine/client/SeaTunnelEngineClusterRoleTest.java
+++
b/seatunnel-engine/seatunnel-engine-client/src/test/java/org/apache/seatunnel/engine/client/SeaTunnelEngineClusterRoleTest.java
@@ -80,7 +80,7 @@ public class SeaTunnelEngineClusterRoleTest {
masterNode =
SeaTunnelServerStarter.createMasterHazelcastInstance(seaTunnelConfig);
HazelcastInstanceImpl finalMasterNode = masterNode;
Awaitility.await()
- .atMost(10000, TimeUnit.MILLISECONDS)
+ .atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
@@ -91,7 +91,7 @@ public class SeaTunnelEngineClusterRoleTest {
HazelcastInstanceImpl finalWorkerNode = workerNode1;
Awaitility.await()
- .atMost(10000, TimeUnit.MILLISECONDS)
+ .atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
@@ -149,7 +149,7 @@ public class SeaTunnelEngineClusterRoleTest {
HazelcastInstanceImpl finalMasterNode = masterNode;
Awaitility.await()
- .atMost(10000, TimeUnit.MILLISECONDS)
+ .atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
@@ -189,6 +189,8 @@ public class SeaTunnelEngineClusterRoleTest {
@Test
public void enterPendingWhenResourcesNotEnough() {
HazelcastInstanceImpl masterNode = null;
+ HazelcastInstanceImpl workerNode1 = null;
+ HazelcastInstanceImpl workerNode2 = null;
String testClusterName = "Test_enterPendingWhenResourcesNotEnough";
SeaTunnelClient seaTunnelClient = null;
@@ -214,7 +216,7 @@ public class SeaTunnelEngineClusterRoleTest {
HazelcastInstanceImpl finalMasterNode = masterNode;
Awaitility.await()
- .atMost(10000, TimeUnit.MILLISECONDS)
+ .atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
@@ -226,7 +228,7 @@ public class SeaTunnelEngineClusterRoleTest {
seaTunnelClient.createExecutionContext(filePath,
jobConfig, seaTunnelConfig);
final ClientJobProxy clientJobProxy = jobExecutionEnv.execute();
Awaitility.await()
- .atMost(10000, TimeUnit.MILLISECONDS)
+ .atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
@@ -235,8 +237,8 @@ public class SeaTunnelEngineClusterRoleTest {
status.contains("PENDING");
// start two worker nodes
-
SeaTunnelServerStarter.createWorkerHazelcastInstance(seaTunnelConfig);
-
SeaTunnelServerStarter.createWorkerHazelcastInstance(seaTunnelConfig);
+ workerNode1 =
SeaTunnelServerStarter.createWorkerHazelcastInstance(seaTunnelConfig);
+ workerNode2 =
SeaTunnelServerStarter.createWorkerHazelcastInstance(seaTunnelConfig);
// There are already resources available, wait for job enter
running or complete
Awaitility.await()
@@ -251,6 +253,12 @@ public class SeaTunnelEngineClusterRoleTest {
if (seaTunnelClient != null) {
seaTunnelClient.close();
}
+ if (workerNode1 != null) {
+ workerNode1.shutdown();
+ }
+ if (workerNode2 != null) {
+ workerNode2.shutdown();
+ }
if (masterNode != null) {
masterNode.shutdown();
}
@@ -291,7 +299,7 @@ public class SeaTunnelEngineClusterRoleTest {
seaTunnelClient.createExecutionContext(filePath,
jobConfig, seaTunnelConfig);
final ClientJobProxy clientJobProxy = jobExecutionEnv.execute();
Awaitility.await()
- .atMost(10000, TimeUnit.MILLISECONDS)
+ .atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
@@ -382,7 +390,7 @@ public class SeaTunnelEngineClusterRoleTest {
masterNode1 =
SeaTunnelServerStarter.createMasterHazelcastInstance(seaTunnelConfig);
HazelcastInstanceImpl finalMasterNode1 = masterNode1;
Awaitility.await()
- .atMost(10000, TimeUnit.MILLISECONDS)
+ .atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
@@ -398,14 +406,14 @@ public class SeaTunnelEngineClusterRoleTest {
masterNode2 =
SeaTunnelServerStarter.createMasterHazelcastInstance(seaTunnelConfig2);
HazelcastInstanceImpl finalWorkerNode = workerNode1;
Awaitility.await()
- .atMost(10000, TimeUnit.MILLISECONDS)
+ .atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
4,
finalWorkerNode.getCluster().getMembers().size()));
masterNode1.shutdown();
Awaitility.await()
- .atMost(10000, TimeUnit.MILLISECONDS)
+ .atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
@@ -424,7 +432,7 @@ public class SeaTunnelEngineClusterRoleTest {
.client;
HazelcastClientInstanceImpl finalHazelcastClient = hazelcastClient;
Awaitility.await()
- .atMost(10000, TimeUnit.MILLISECONDS)
+ .atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() -> {
UUID masterUuid =
@@ -454,12 +462,15 @@ public class SeaTunnelEngineClusterRoleTest {
.listJobStatus(true)
.contains("RUNNING")));
jobClient.cancelJob(jobId);
- await().pollDelay(10000, TimeUnit.MILLISECONDS)
- .atMost(60000, TimeUnit.MILLISECONDS)
+ await().atMost(120000, TimeUnit.MILLISECONDS)
.untilAsserted(
- () ->
- Assertions.assertEquals(
- "CANCELED",
jobClient.getJobStatus(jobId)));
+ () -> {
+ String status = jobClient.getJobStatus(jobId);
+ Assertions.assertEquals(
+ "CANCELED",
+ status,
+ "Expected terminal state but was: " +
status);
+ });
} finally {
if (hazelcastClient != null) {
hazelcastClient.shutdown();
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
index 18a83020e9..a2ab949248 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
@@ -109,7 +109,7 @@ public class CoordinatorServiceTest {
SeaTunnelServer server1 =
instance1.node.getNodeEngine().getService(SeaTunnelServer.SERVICE_NAME);
- await().atMost(10, TimeUnit.SECONDS)
+ await().atMost(60, TimeUnit.SECONDS)
.untilAsserted(
() -> {
Assertions.assertTrue(server1.isMasterNode());
@@ -123,7 +123,7 @@ public class CoordinatorServiceTest {
SeaTunnelServer server2 =
instance2.node.getNodeEngine().getService(SeaTunnelServer.SERVICE_NAME);
- await().atMost(10, TimeUnit.SECONDS)
+ await().atMost(60, TimeUnit.SECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
@@ -345,10 +345,10 @@ public class CoordinatorServiceTest {
masterFlag.set(true);
invokeCheckNewActiveMaster(coordinatorService);
- await().atMost(5, TimeUnit.SECONDS)
+ await().atMost(30, TimeUnit.SECONDS)
.untilAsserted(
() ->
Assertions.assertTrue(coordinatorService.isCoordinatorActive()));
- await().atMost(5, TimeUnit.SECONDS)
+ await().atMost(30, TimeUnit.SECONDS)
.untilAsserted(() -> Assertions.assertEquals(0L,
runLatch.getCount()));
Mockito.verify(jobMaster, Mockito.times(1)).run();
} finally {
@@ -396,10 +396,10 @@ public class CoordinatorServiceTest {
invokeCheckNewActiveMaster(newCoordinator);
- await().atMost(5, TimeUnit.SECONDS)
+ await().atMost(30, TimeUnit.SECONDS)
.untilAsserted(
() ->
Assertions.assertTrue(newCoordinator.isCoordinatorActive()));
- await().atMost(5, TimeUnit.SECONDS)
+ await().atMost(30, TimeUnit.SECONDS)
.untilAsserted(() -> Assertions.assertEquals(0L,
newRunLatch.getCount()));
Mockito.verify(newPendingJob, Mockito.times(1)).run();
Assertions.assertFalse(newCoordinator.getPendingJobQueue().contains(30001L));
@@ -424,10 +424,10 @@ public class CoordinatorServiceTest {
JobMaster jobMaster = enqueueMockPendingJob(coordinatorService,
40001L, runLatch);
invokeCheckNewActiveMaster(coordinatorService);
- await().atMost(5, TimeUnit.SECONDS)
+ await().atMost(30, TimeUnit.SECONDS)
.untilAsserted(
() ->
Assertions.assertTrue(coordinatorService.isCoordinatorActive()));
- await().atMost(5, TimeUnit.SECONDS)
+ await().atMost(30, TimeUnit.SECONDS)
.untilAsserted(() -> Assertions.assertEquals(0L,
runLatch.getCount()));
invokeCheckNewActiveMaster(coordinatorService);
@@ -454,7 +454,7 @@ public class CoordinatorServiceTest {
enqueueMockPendingJob(coordinatorService, 50001L,
runLatch, false);
invokeCheckNewActiveMaster(coordinatorService);
- await().atMost(5, TimeUnit.SECONDS)
+ await().atMost(30, TimeUnit.SECONDS)
.untilAsserted(
() ->
Assertions.assertTrue(coordinatorService.isCoordinatorActive()));
await().during(1, TimeUnit.SECONDS)
@@ -483,10 +483,10 @@ public class CoordinatorServiceTest {
enqueueMockPendingJob(coordinatorService, 60001L,
runLatch, false);
invokeCheckNewActiveMaster(coordinatorService);
- await().atMost(5, TimeUnit.SECONDS)
+ await().atMost(30, TimeUnit.SECONDS)
.untilAsserted(
() ->
Assertions.assertTrue(coordinatorService.isCoordinatorActive()));
- await().atMost(5, TimeUnit.SECONDS)
+ await().atMost(30, TimeUnit.SECONDS)
.untilAsserted(
() ->
Assertions.assertFalse(
@@ -669,32 +669,33 @@ public class CoordinatorServiceTest {
SeaTunnelServerStarter.createHazelcastInstance(
TestUtils.getClusterName(
"CoordinatorServiceTest_testInvocationFutureUseCompletableFutureExecutor"));
-
- NodeEngineUtil.sendOperationToMemberNode(
- instance.node.getNodeEngine(),
- new PrintMessageOperation("hello"),
- instance.getCluster().getLocalMember().getAddress())
- .whenComplete(
- (aVoid, error) -> {
- Assertions.assertTrue(
- Thread.currentThread()
- .getName()
-
.startsWith("SeaTunnel-CompletableFuture-Thread"));
- })
- .join();
-
- NodeEngineUtil.sendOperationToMasterNode(
- instance.node.getNodeEngine(), new
PrintMessageOperation("hello"))
- .whenCompleteAsync(
- (aVoid, error) -> {
- Assertions.assertTrue(
- Thread.currentThread()
- .getName()
-
.startsWith("SeaTunnel-CompletableFuture-Thread"));
- })
- .join();
-
- instance.shutdown();
+ try {
+ NodeEngineUtil.sendOperationToMemberNode(
+ instance.node.getNodeEngine(),
+ new PrintMessageOperation("hello"),
+
instance.getCluster().getLocalMember().getAddress())
+ .whenComplete(
+ (aVoid, error) -> {
+ Assertions.assertTrue(
+ Thread.currentThread()
+ .getName()
+
.startsWith("SeaTunnel-CompletableFuture-Thread"));
+ })
+ .join();
+
+ NodeEngineUtil.sendOperationToMasterNode(
+ instance.node.getNodeEngine(), new
PrintMessageOperation("hello"))
+ .whenCompleteAsync(
+ (aVoid, error) -> {
+ Assertions.assertTrue(
+ Thread.currentThread()
+ .getName()
+
.startsWith("SeaTunnel-CompletableFuture-Thread"));
+ })
+ .join();
+ } finally {
+ instance.shutdown();
+ }
}
private static final class BlockingEventProcessor implements
EventProcessor {
@@ -770,7 +771,7 @@ public class CoordinatorServiceTest {
"test_force_stop_running_job");
CoordinatorService coordinatorService =
jobInformation.coordinatorService;
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() -> {
Assertions.assertEquals(
@@ -806,7 +807,7 @@ public class CoordinatorServiceTest {
"test_force_stop_abnormal_savepoint_job");
CoordinatorService coordinatorService =
jobInformation.coordinatorService;
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() -> {
Assertions.assertEquals(
@@ -853,7 +854,7 @@ public class CoordinatorServiceTest {
.getPendingJobQueue()
.contains(jobInformation.jobId));
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertFalse(
@@ -879,7 +880,7 @@ public class CoordinatorServiceTest {
IMap<Object, Object> runningJobStateIMap =
coordinatorService.getJobMaster(jobInformation.jobId).getRunningJobStateIMap();
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() -> {
Assertions.assertEquals(
@@ -894,7 +895,7 @@ public class CoordinatorServiceTest {
.containsKey(jobInformation.jobId));
});
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() -> {
Assertions.assertEquals(
@@ -921,9 +922,9 @@ public class CoordinatorServiceTest {
CoordinatorService coordinatorService =
jobInformation.coordinatorService;
IMap<Long, HashMap<TaskLocation, SeaTunnelMetricsContext>> metricsImap
=
coordinatorService.getMetricsImap();
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(() ->
Assertions.assertFalse(metricsImap.isEmpty()));
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(() ->
Assertions.assertTrue(metricsImap.isEmpty()));
jobInformation.coordinatorService.clearCoordinatorService();
@@ -942,9 +943,9 @@ public class CoordinatorServiceTest {
CoordinatorService coordinatorService =
jobInformation.coordinatorService;
IMap<Long, HashMap<TaskLocation, SeaTunnelMetricsContext>> metricsImap
=
coordinatorService.getMetricsImap();
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(() ->
Assertions.assertFalse(metricsImap.isEmpty()));
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(() ->
Assertions.assertTrue(metricsImap.isEmpty()));
jobInformation.coordinatorService.clearCoordinatorService();
@@ -987,7 +988,7 @@ public class CoordinatorServiceTest {
throw new CompletionException(e);
}
});
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(() -> Assertions.assertEquals(10,
metricsImap.size()));
} finally {
instance1.shutdown();
@@ -1002,14 +1003,19 @@ public class CoordinatorServiceTest {
"CoordinatorServiceTest_testCleanPendingJobMasterMap",
"batch_fake_to_inmemory.conf",
"test_clean_pending_jobmastermap");
- CoordinatorService coordinatorService =
jobInformation.coordinatorService;
- await().atMost(20000, TimeUnit.MILLISECONDS)
- .untilAsserted(
- () ->
- Assertions.assertFalse(
- coordinatorService
- .getPendingJobQueue()
-
.contains(jobInformation.jobId)));
+ try {
+ CoordinatorService coordinatorService =
jobInformation.coordinatorService;
+ await().atMost(20000, TimeUnit.MILLISECONDS)
+ .untilAsserted(
+ () ->
+ Assertions.assertFalse(
+ coordinatorService
+ .getPendingJobQueue()
+
.contains(jobInformation.jobId)));
+ } finally {
+ jobInformation.coordinatorService.clearCoordinatorService();
+ jobInformation.coordinatorServiceTest.shutdown();
+ }
}
@Test
@@ -1051,7 +1057,7 @@ public class CoordinatorServiceTest {
jobImmutableInformation.isStartWithSavePoint()))
.join());
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertNotEquals(
@@ -1200,18 +1206,22 @@ public class CoordinatorServiceTest {
Long jobId = jobInformation.jobId;
HazelcastInstanceImpl coordinatorServiceTest =
jobInformation.coordinatorServiceTest;
- // waiting for job status turn to running
- await().atMost(10000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
JobStatus.RUNNING,
coordinatorService.getJobStatus(jobId)));
- try {
- Thread.sleep(5000);
- } catch (InterruptedException e) {
- throw new RuntimeException(e);
- }
+ await().atMost(30, TimeUnit.SECONDS)
+ .untilAsserted(
+ () ->
+ Assertions.assertTrue(
+
Thread.getAllStackTraces().keySet().stream()
+ .anyMatch(
+ thread ->
+
thread.getName()
+
.startsWith(
+
"pending-job-schedule-runner"))));
int scheduleRunnerThreadCount =
(int)
@@ -1250,7 +1260,7 @@ public class CoordinatorServiceTest {
SeaTunnelServer server =
instance.node.getNodeEngine().getService(SeaTunnelServer.SERVICE_NAME);
CoordinatorService coordinatorService =
server.getCoordinatorService();
- await().atMost(10, TimeUnit.SECONDS)
+ await().atMost(60, TimeUnit.SECONDS)
.untilAsserted(
() ->
Assertions.assertTrue(coordinatorService.isCoordinatorActive()));
@@ -1296,7 +1306,7 @@ public class CoordinatorServiceTest {
coordinatorService, "runningJobInfoIMap",
runningJobInfoIMap);
}
- await().atMost(10, TimeUnit.SECONDS)
+ await().atMost(60, TimeUnit.SECONDS)
.untilAsserted(
() ->
Assertions.assertTrue(
@@ -1478,6 +1488,7 @@ public class CoordinatorServiceTest {
executor.awaitTermination(30, TimeUnit.SECONDS);
instance1.shutdown();
instance2.shutdown();
+ instance3.shutdown();
}
}
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceWithCancelPendingJobTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceWithCancelPendingJobTest.java
index 2639ef2280..2544b71ef0 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceWithCancelPendingJobTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceWithCancelPendingJobTest.java
@@ -142,8 +142,7 @@ public class CoordinatorServiceWithCancelPendingJobTest
extends AbstractSeaTunne
nodeEngine.getHazelcastInstance().getMap(Constant.IMAP_STATE_TIMESTAMPS);
// Verify if the final status of the task is cancelled
- await().pollDelay(3, TimeUnit.SECONDS)
- .atMost(120, TimeUnit.SECONDS)
+ await().atMost(120, TimeUnit.SECONDS)
.untilAsserted(
() -> {
Assertions.assertEquals(
@@ -185,14 +184,15 @@ public class CoordinatorServiceWithCancelPendingJobTest
extends AbstractSeaTunne
JobMaster jobMaster =
server.getCoordinatorService().getJobMaster(jobId);
- // waiting for job status turn to running
await().atMost(120, TimeUnit.SECONDS)
.untilAsserted(
- () -> Assertions.assertEquals(JobStatus.PENDING,
jobMaster.getJobStatus()));
-
- // Because handleCheckpointTimeout is an async method, so we need
sleep 5s to waiting job
- // status become running again
- Thread.sleep(5000);
+ () -> {
+ JobStatus status = jobMaster.getJobStatus();
+ Assertions.assertTrue(
+ JobStatus.PENDING.equals(status)
+ ||
JobStatus.RUNNING.equals(status),
+ "Expected PENDING or RUNNING but was: " +
status);
+ });
return jobMaster;
}
}
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointStorageTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointStorageTest.java
index b95435523f..d6b29834db 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointStorageTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointStorageTest.java
@@ -33,6 +33,8 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.DisabledOnOs;
import org.junit.jupiter.api.condition.OS;
+import java.io.FileNotFoundException;
+import java.nio.file.NoSuchFileException;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.TimeUnit;
@@ -124,9 +126,22 @@ public class CheckpointStorageTest extends
AbstractSeaTunnelServerTest {
Assertions.assertEquals(
JobStatus.FINISHED,
server.getCoordinatorService().getJobStatus(jobId)));
- List<PipelineState> allCheckpoints =
- checkpointStorage.getAllCheckpoints(String.valueOf(jobId));
- Assertions.assertEquals(0, allCheckpoints.size());
+ await().atMost(30000, TimeUnit.MILLISECONDS)
+ .untilAsserted(
+ () -> {
+ try {
+ List<PipelineState> allCheckpoints =
+
checkpointStorage.getAllCheckpoints(String.valueOf(jobId));
+ Assertions.assertEquals(0,
allCheckpoints.size());
+ } catch (CheckpointStorageException e) {
+ Throwable cause = e.getCause();
+ if (cause instanceof FileNotFoundException
+ || cause instanceof
NoSuchFileException) {
+ return;
+ }
+ throw e;
+ }
+ });
}
@Test
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointTimeOutTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointTimeOutTest.java
index af118329ae..79c3115c54 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointTimeOutTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointTimeOutTest.java
@@ -43,7 +43,7 @@ public class CheckpointTimeOutTest extends
AbstractSeaTunnelServerTest {
@Test
public void testJobLevelCheckpointTimeOut() {
long jobId = System.currentTimeMillis();
- startJob(System.currentTimeMillis(), CONF_PATH);
+ startJob(jobId, CONF_PATH);
await().atMost(120000, TimeUnit.MILLISECONDS)
.untilAsserted(
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/master/FollowerRunningJobsFilterTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/master/FollowerRunningJobsFilterTest.java
index 846e9fd24e..266eaed4da 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/master/FollowerRunningJobsFilterTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/master/FollowerRunningJobsFilterTest.java
@@ -60,7 +60,7 @@ class FollowerRunningJobsFilterTest
follower =
SeaTunnelServerStarter.createHazelcastInstance(seaTunnelConfig);
Awaitility.await()
- .atMost(10, TimeUnit.SECONDS)
+ .atMost(60, TimeUnit.SECONDS)
.untilAsserted(
() ->
Assertions.assertEquals(
@@ -80,7 +80,7 @@ class FollowerRunningJobsFilterTest
JobInfoService followerJobInfoService =
new JobInfoService((NodeEngineImpl)
follower.node.nodeEngine);
Awaitility.await()
- .atMost(10, TimeUnit.SECONDS)
+ .atMost(60, TimeUnit.SECONDS)
.untilAsserted(
() ->
Assertions.assertTrue(
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/master/JobHistoryServiceTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/master/JobHistoryServiceTest.java
index 50e5fcb317..a9169b02c9 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/master/JobHistoryServiceTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/master/JobHistoryServiceTest.java
@@ -71,8 +71,7 @@ class JobHistoryServiceTest extends
AbstractSeaTunnelServerTest {
});
// waiting for JOB_1 status turn to FINISHED
- await().pollDelay(5, TimeUnit.SECONDS)
- .atMost(60000, TimeUnit.MILLISECONDS)
+ await().atMost(60000, TimeUnit.MILLISECONDS)
.untilAsserted(
() -> {
List<JobStatusData> jobStatusData = listJob();
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/BaseServletTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/BaseServletTest.java
index 8cbf7bc7f3..3d3d734cfc 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/BaseServletTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/BaseServletTest.java
@@ -46,8 +46,9 @@ class BaseServletTest extends AbstractSeaTunnelServerTest {
private static final Long JOB_1 = System.currentTimeMillis() + 1L;
+ @Override
@BeforeAll
- void setUp() {
+ public void before() {
String name = this.getClass().getName();
Config hazelcastConfig = Config.loadFromString(getHazelcastConfig());
hazelcastConfig.setClusterName(TestUtils.getClusterName("RestApiServletTest_" +
name));
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpBasicTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpBasicTest.java
index 0117325913..235ae0278e 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpBasicTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpBasicTest.java
@@ -64,8 +64,9 @@ class RestApiHttpBasicTest extends
AbstractSeaTunnelServerTest {
private static final String AUTHORIZATION_HEADER = "Authorization";
private static final String BASIC_PREFIX = "Basic ";
+ @Override
@BeforeAll
- void setUp() {
+ public void before() {
String name = this.getClass().getName();
Config hazelcastConfig = Config.loadFromString(getHazelcastConfig());
hazelcastConfig.setClusterName(
@@ -99,6 +100,7 @@ class RestApiHttpBasicTest extends
AbstractSeaTunnelServerTest {
httpConfig.setEnableBasicAuth(Boolean.FALSE);
httpConfig.setBasicAuthUsername("");
httpConfig.setBasicAuthPassword("");
+ super.after();
}
@Test
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpsForTruststoreTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpsForTruststoreTest.java
index f7e848f0c6..54e0f5246c 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpsForTruststoreTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpsForTruststoreTest.java
@@ -54,8 +54,9 @@ public class RestApiHttpsForTruststoreTest extends
AbstractSeaTunnelServerTest {
private static final String CLIENT_KEYSTORE_PASSWORD =
"client_keystore_password";
private static final String CLIENT_TRUSTSTORE_PASSWORD =
"client_truststore_password";
+ @Override
@BeforeAll
- public void setUp() {
+ public void before() {
String name = this.getClass().getName();
Config hazelcastConfig = Config.loadFromString(getHazelcastConfig());
hazelcastConfig.setClusterName(
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpsTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpsTest.java
index 29c3b90b51..f9d5ebae7e 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpsTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/RestApiHttpsTest.java
@@ -67,8 +67,9 @@ public class RestApiHttpsTest extends
AbstractSeaTunnelServerTest {
private static final String SERVER_KEYSTORE_PASSWORD =
"server_keystore_password";
private static final String CLIENT_KEYSTORE_PASSWORD =
"client_keystore_password";
+ @Override
@BeforeAll
- public void setUp() {
+ public void before() {
String name = this.getClass().getName();
Config hazelcastConfig = Config.loadFromString(getHazelcastConfig());
hazelcastConfig.setClusterName(TestUtils.getClusterName("RestApiHttpsTest_" +
name));
@@ -149,9 +150,8 @@ public class RestApiHttpsTest extends
AbstractSeaTunnelServerTest {
}
// wait until all jobs are finished
- await().pollDelay(5, TimeUnit.SECONDS)
- .atMost(30, TimeUnit.SECONDS)
- .pollInterval(100, TimeUnit.MILLISECONDS)
+ await().atMost(60, TimeUnit.SECONDS)
+ .pollInterval(500, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
assertEquals(
@@ -255,9 +255,8 @@ public class RestApiHttpsTest extends
AbstractSeaTunnelServerTest {
}
// wait until all jobs are finished
- await().pollDelay(5, TimeUnit.SECONDS)
- .atMost(30, TimeUnit.SECONDS)
- .pollInterval(100, TimeUnit.MILLISECONDS)
+ await().atMost(60, TimeUnit.SECONDS)
+ .pollInterval(500, TimeUnit.MILLISECONDS)
.untilAsserted(
() ->
assertEquals(
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/EngineStateStoreMetricExportsTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/EngineStateStoreMetricExportsTest.java
index 71c070031a..f9390dd993 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/EngineStateStoreMetricExportsTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/telemetry/metrics/exports/EngineStateStoreMetricExportsTest.java
@@ -57,7 +57,7 @@ class EngineStateStoreMetricExportsTest {
instance =
SeaTunnelServerStarter.createHazelcastInstance(
TestUtils.getClusterName("EngineStateStoreMetricExportsTest_localMetrics"));
- await().atMost(10, TimeUnit.SECONDS)
+ await().atMost(60, TimeUnit.SECONDS)
.untilAsserted(() ->
Assertions.assertTrue(instance.node.isMaster()));
instance.getMap(Constant.IMAP_RUNNING_JOB_INFO).put(1L, "job-info");
@@ -92,7 +92,7 @@ class EngineStateStoreMetricExportsTest {
SeaTunnelServerStarter.createHazelcastInstance(
TestUtils.getClusterName(
"EngineStateStoreMetricExportsTest_allStateStores"));
- await().atMost(10, TimeUnit.SECONDS)
+ await().atMost(60, TimeUnit.SECONDS)
.untilAsserted(() ->
Assertions.assertTrue(instance.node.isMaster()));
List<MetricFamilySamples> metrics =