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 06dae2b7ab [Test][Zeta] Stabilize coordinator retryable operation test
(#11425)
06dae2b7ab is described below
commit 06dae2b7abe0bad8ce3b75f931158df004f05f70
Author: Daniel <[email protected]>
AuthorDate: Wed Jul 15 22:32:43 2026 +0800
[Test][Zeta] Stabilize coordinator retryable operation test (#11425)
Co-authored-by: DanielLeens <[email protected]>
---
.../engine/server/CoordinatorServiceTest.java | 40 ++++++++++++++++------
.../operation/ReturnRetryTimesOperation.java | 18 ++++++++++
2 files changed, 48 insertions(+), 10 deletions(-)
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 8c934a9b5b..286ba67fa6 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
@@ -71,6 +71,7 @@ import com.hazelcast.logging.ILogger;
import com.hazelcast.map.IMap;
import com.hazelcast.spi.exception.RetryableHazelcastException;
import com.hazelcast.spi.impl.NodeEngineImpl;
+import com.hazelcast.spi.properties.ClusterProperty;
import lombok.extern.slf4j.Slf4j;
import java.lang.reflect.Field;
@@ -596,26 +597,45 @@ public class CoordinatorServiceTest {
@Test
public void
testSeaTunnelEngineRetryableExceptionOperationCanBeRetryByHazelcast() {
- HazelcastInstanceImpl instance =
- SeaTunnelServerStarter.createHazelcastInstance(
+ int maxRetryCount = 3;
+ ReturnRetryTimesOperation.resetRetryTimes();
+ SeaTunnelConfig seaTunnelConfig =
ConfigProvider.locateAndGetSeaTunnelConfig();
+ seaTunnelConfig
+ .getHazelcastConfig()
+ .setClusterName(
TestUtils.getClusterName(
"CoordinatorServiceTest_testSeaTunnelEngineRetryableExceptionOperationCanBeRetryByHazelcast"));
+ // Keep the production retry contract intact while using a small
test-only budget so CI
+ // can still verify terminal exception propagation without waiting for
250 retries.
+ seaTunnelConfig
+ .getHazelcastConfig()
+ .setProperty(
+ ClusterProperty.INVOCATION_MAX_RETRY_COUNT.getName(),
+ String.valueOf(maxRetryCount));
+ seaTunnelConfig
+ .getHazelcastConfig()
+ .setProperty(ClusterProperty.INVOCATION_RETRY_PAUSE.getName(),
"1");
+ HazelcastInstanceImpl instance =
+
SeaTunnelServerStarter.createHazelcastInstance(seaTunnelConfig);
try {
CompletionException exception =
Assertions.assertThrows(
CompletionException.class,
- () -> {
- NodeEngineUtil.sendOperationToMemberNode(
- instance.node.getNodeEngine(),
- new
ReturnRetryTimesOperation(),
-
instance.getCluster().getLocalMember().getAddress())
- .join();
- });
+ () ->
+ NodeEngineUtil.sendOperationToMemberNode(
+
instance.node.getNodeEngine(),
+ new
ReturnRetryTimesOperation(),
+ instance.getCluster()
+ .getLocalMember()
+ .getAddress())
+ .join());
Assertions.assertTrue(
exception
.getCause()
.getMessage()
- .contains("Retryable exception occurred, retry
times: 250"));
+ .contains(
+ "Retryable exception occurred, retry
times: " + maxRetryCount));
+ Assertions.assertEquals(maxRetryCount,
ReturnRetryTimesOperation.getRetryTimes());
} finally {
instance.shutdown();
}
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/ReturnRetryTimesOperation.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/ReturnRetryTimesOperation.java
index 9f3a1f46fe..d90a4c4165 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/ReturnRetryTimesOperation.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/ReturnRetryTimesOperation.java
@@ -30,6 +30,24 @@ public class ReturnRetryTimesOperation extends Operation
private static final AtomicInteger retryTimes = new AtomicInteger(0);
+ /**
+ * Reset the static retry counter before starting a new retry-focused test.
+ *
+ * <p>The operation is reused across Hazelcast retries in the same test
JVM.
+ */
+ public static void resetRetryTimes() {
+ retryTimes.set(0);
+ }
+
+ /**
+ * Return the number of operation attempts observed by the current test.
+ *
+ * <p>The terminal assertion uses this value to verify retry exhaustion.
+ */
+ public static int getRetryTimes() {
+ return retryTimes.get();
+ }
+
@Override
public void run() {
retryTimes.getAndIncrement();