This is an automated email from the ASF dual-hosted git repository.
zhoujinsong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/amoro.git
The following commit(s) were added to refs/heads/master by this push:
new 4953d9c90 [AMORO-4303][AMS] Stabilize optimizer heartbeat tests (#4307)
4953d9c90 is described below
commit 4953d9c90d378d0f031bad3a1bc879a0dc16512c
Author: Power John <[email protected]>
AuthorDate: Mon Aug 3 20:04:49 2026 +0800
[AMORO-4303][AMS] Stabilize optimizer heartbeat tests (#4307)
---
.../amoro/server/DefaultOptimizingService.java | 4 +-
.../apache/amoro/server/AMSServiceTestBase.java | 12 +++-
.../amoro/server/TestDefaultOptimizingService.java | 74 +++++++++++++++-------
3 files changed, 66 insertions(+), 24 deletions(-)
diff --git
a/amoro-ams/src/main/java/org/apache/amoro/server/DefaultOptimizingService.java
b/amoro-ams/src/main/java/org/apache/amoro/server/DefaultOptimizingService.java
index 6831d57c7..427b6edd5 100644
---
a/amoro-ams/src/main/java/org/apache/amoro/server/DefaultOptimizingService.java
+++
b/amoro-ams/src/main/java/org/apache/amoro/server/DefaultOptimizingService.java
@@ -608,7 +608,9 @@ public class DefaultOptimizingService extends
StatedPersistentBase
}
} catch (InterruptedException ignored) {
} catch (Throwable t) {
- LOG.error("{} has encountered a problem.",
this.getClass().getSimpleName(), t);
+ if (!stopped) {
+ LOG.error("{} has encountered a problem.",
this.getClass().getSimpleName(), t);
+ }
}
}
}
diff --git
a/amoro-ams/src/test/java/org/apache/amoro/server/AMSServiceTestBase.java
b/amoro-ams/src/test/java/org/apache/amoro/server/AMSServiceTestBase.java
index b52b71c21..9c2abe969 100644
--- a/amoro-ams/src/test/java/org/apache/amoro/server/AMSServiceTestBase.java
+++ b/amoro-ams/src/test/java/org/apache/amoro/server/AMSServiceTestBase.java
@@ -33,16 +33,24 @@ import org.junit.BeforeClass;
import java.time.Duration;
public abstract class AMSServiceTestBase extends AMSManagerTestBase {
+ // Normal optimizer tests may block in pollTask for 3 seconds. Keep the
heartbeat timeout well
+ // outside that window so CI scheduling delays do not turn unrelated tests
into expiry tests.
+ private static final Duration DEFAULT_OPTIMIZER_HEARTBEAT_TIMEOUT =
Duration.ofSeconds(10);
+
private static DefaultTableService TABLE_SERVICE = null;
private static DefaultOptimizingService OPTIMIZING_SERVICE = null;
private static ProcessService PROCESS_SERVICE = null;
@BeforeClass
public static void initTableService() {
+ initTableService(DEFAULT_OPTIMIZER_HEARTBEAT_TIMEOUT);
+ }
+
+ protected static void initTableService(Duration optimizerHeartbeatTimeout) {
DefaultTableRuntimeFactory runtimeFactory = new
DefaultTableRuntimeFactory();
try {
Configurations configurations = new Configurations();
- configurations.set(AmoroManagementConf.OPTIMIZER_HB_TIMEOUT,
Duration.ofMillis(800L));
+ configurations.set(AmoroManagementConf.OPTIMIZER_HB_TIMEOUT,
optimizerHeartbeatTimeout);
configurations.set(
AmoroManagementConf.OPTIMIZER_TASK_EXECUTE_TIMEOUT,
Duration.ofMillis(30000L));
// must stay above OPTIMIZER_POLLING_TIMEOUT (3s): a blocking pollTask
waiting out its full
@@ -78,6 +86,8 @@ public abstract class AMSServiceTestBase extends
AMSManagerTestBase {
@AfterClass
public static void disposeTableService() {
TABLE_SERVICE.dispose();
+ OPTIMIZING_SERVICE.dispose();
+ PROCESS_SERVICE.dispose();
MetricManager.dispose();
EventsManager.dispose();
}
diff --git
a/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultOptimizingService.java
b/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultOptimizingService.java
index 365eee02d..ff2c94a42 100644
---
a/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultOptimizingService.java
+++
b/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultOptimizingService.java
@@ -70,6 +70,7 @@ import org.junit.jupiter.api.Assertions;
import org.junit.runner.RunWith;
import org.junit.runners.Parameterized;
+import java.time.Duration;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
@@ -77,9 +78,14 @@ import java.util.Map;
@RunWith(Parameterized.class)
public class TestDefaultOptimizingService extends AMSTableTestBase {
+ private static final Duration EXPIRATION_TEST_HEARTBEAT_TIMEOUT =
Duration.ofMillis(800);
+ private static final long OPTIMIZER_HEARTBEAT_INTERVAL_MS = 100;
+ private static final long ASYNC_WAIT_TIMEOUT_MS = 10000;
+
private final int THREAD_ID = 0;
private String token;
private Toucher toucher;
+ private boolean customHeartbeatTimeout;
@Parameterized.Parameters(name = "{0}, {1}")
public static Object[] parameters() {
@@ -121,6 +127,12 @@ public class TestDefaultOptimizingService extends
AMSTableTestBase {
dropDatabase();
} catch (Exception e) {
// ignore
+ } finally {
+ if (customHeartbeatTimeout) {
+ disposeTableService();
+ initTableService();
+ customHeartbeatTimeout = false;
+ }
}
}
@@ -255,20 +267,22 @@ public class TestDefaultOptimizingService extends
AMSTableTestBase {
@Test
public void testTouchTimeout() throws InterruptedException {
+ rebootWithHeartbeatTimeout(EXPIRATION_TEST_HEARTBEAT_TIMEOUT);
OptimizingTask task = optimizingService().pollTask(token, THREAD_ID);
Assertions.assertNotNull(task);
+ String expiredToken = token;
toucher.stop();
toucher = null;
- Thread.sleep(1000);
- Assertions.assertThrows(PluginRetryAuthException.class, () ->
optimizingService().touch(token));
+ waitForOptimizerExpiration(expiredToken, ASYNC_WAIT_TIMEOUT_MS);
Assertions.assertThrows(
- PluginRetryAuthException.class, () ->
optimizingService().pollTask(token, THREAD_ID));
+ PluginRetryAuthException.class, () ->
optimizingService().touch(expiredToken));
+ Assertions.assertThrows(
+ PluginRetryAuthException.class,
+ () -> optimizingService().pollTask(expiredToken, THREAD_ID));
// After optimizer expires, its tasks are immediately reset to PLANNED
// because unregister happens before task scan in OptimizerKeeper
- assertTaskStatus(TaskRuntime.Status.PLANNED);
- token = optimizingService().authenticate(buildRegisterInfo());
+ waitForTaskStatus(TaskRuntime.Status.PLANNED, ASYNC_WAIT_TIMEOUT_MS);
toucher = new Toucher();
- Thread.sleep(1000);
assertTaskStatus(TaskRuntime.Status.PLANNED);
OptimizingTask task2 = optimizingService().pollTask(token, THREAD_ID);
Assertions.assertEquals(task2.getTaskId(), task.getTaskId());
@@ -283,11 +297,10 @@ public class TestDefaultOptimizingService extends
AMSTableTestBase {
public void testRebootAndPoll() throws InterruptedException {
OptimizingTask task = optimizingService().pollTask(token, THREAD_ID);
Assertions.assertNotNull(task);
- reboot();
+ rebootWithHeartbeatTimeout(EXPIRATION_TEST_HEARTBEAT_TIMEOUT);
// wait for last optimizer expiring
- Thread.sleep(1000);
- assertTaskStatus(TaskRuntime.Status.PLANNED);
+ waitForTaskStatus(TaskRuntime.Status.PLANNED, ASYNC_WAIT_TIMEOUT_MS);
OptimizingTask task2 = optimizingService().pollTask(token, THREAD_ID);
Assertions.assertNotNull(task2);
Assertions.assertEquals(task2.getTaskId(), task.getTaskId());
@@ -329,10 +342,10 @@ public class TestDefaultOptimizingService extends
AMSTableTestBase {
Assertions.assertNotNull(task);
assertTaskStatus(TaskRuntime.Status.SCHEDULED); // polled but NOT acked
- // the optimizer stays alive (Toucher touches every 300ms), so waiting
past the ack timeout hits
+ // the optimizer stays alive, so waiting past the ack timeout hits
// the SCHEDULED + ackTimeout branch rather than the optimizer-expired
branch: the keeper resets
// the task out from under the live optimizer
- waitForTaskStatus(TaskRuntime.Status.PLANNED, 10000);
+ waitForTaskStatus(TaskRuntime.Status.PLANNED, 20000);
// the delayed ack arrives for the now-reset task -> rejected, exactly
like the issue
Assertions.assertThrows(
@@ -351,11 +364,8 @@ public class TestDefaultOptimizingService extends
AMSTableTestBase {
optimizingService().listTasks(defaultResourceGroup().getName()).get(0);
assertTaskStatus(TaskRuntime.Status.ACKED);
- // In this test, OPTIMIZER_TASK_EXECUTE_TIMEOUT is set to 30 seconds, so
after waiting 45
- // seconds the task will be considered suspended and retried
- Thread.sleep(45000);
-
- assertTaskStatus(TaskRuntime.Status.PLANNED);
+ // In this test, OPTIMIZER_TASK_EXECUTE_TIMEOUT is set to 30 seconds.
+ waitForTaskStatus(TaskRuntime.Status.PLANNED, 60000);
OptimizingTask task2 = optimizingService().pollTask(token, THREAD_ID);
Assertions.assertNotNull(task2);
Assertions.assertEquals(task2.getTaskId(), task.getTaskId());
@@ -706,7 +716,9 @@ public class TestDefaultOptimizingService extends
AMSTableTestBase {
private OptimizerRegisterInfo buildRegisterInfo() {
OptimizerRegisterInfo registerInfo = new OptimizerRegisterInfo();
Map<String, String> registerProperties = Maps.newHashMap();
- registerProperties.put(OptimizerProperties.OPTIMIZER_HEART_BEAT_INTERVAL,
"100");
+ registerProperties.put(
+ OptimizerProperties.OPTIMIZER_HEART_BEAT_INTERVAL,
+ String.valueOf(OPTIMIZER_HEARTBEAT_INTERVAL_MS));
registerInfo.setProperties(registerProperties);
registerInfo.setThreadCount(1);
registerInfo.setMemoryMb(1024);
@@ -785,6 +797,23 @@ public class TestDefaultOptimizingService extends
AMSTableTestBase {
assertTaskStatus(expectedStatus);
}
+ private void waitForOptimizerExpiration(String optimizerToken, long
timeoutMs)
+ throws InterruptedException {
+ long deadline = System.currentTimeMillis() + timeoutMs;
+ while (System.currentTimeMillis() < deadline) {
+ boolean optimizerExists =
+ optimizerManager().listOptimizers().stream()
+ .anyMatch(optimizer ->
optimizerToken.equals(optimizer.getToken()));
+ boolean optimizerAuthenticated =
+ optimizingService().getTotalQuota(defaultResourceGroup().getName())
> 0;
+ if (!optimizerExists && !optimizerAuthenticated) {
+ return;
+ }
+ Thread.sleep(100);
+ }
+ Assertions.fail("Optimizer did not expire within " + timeoutMs + " ms");
+ }
+
private void assertTaskCompleted(TaskRuntime<?> taskRuntime) {
if (taskRuntime != null) {
Assertions.assertEquals(TaskRuntime.Status.SUCCESS,
taskRuntime.getStatus());
@@ -798,17 +827,18 @@ public class TestDefaultOptimizingService extends
AMSTableTestBase {
}
protected void reload() {
- disposeTableService();
toucher.suspend();
+ disposeTableService();
initTableService();
toucher.goOn();
}
- protected void reboot() throws InterruptedException {
- disposeTableService();
+ protected void rebootWithHeartbeatTimeout(Duration heartbeatTimeout) throws
InterruptedException {
toucher.stop();
toucher = null;
- initTableService();
+ disposeTableService();
+ customHeartbeatTimeout = true;
+ initTableService(heartbeatTimeout);
toucher = new Toucher();
}
@@ -855,7 +885,7 @@ public class TestDefaultOptimizingService extends
AMSTableTestBase {
public void run() {
while (!stop) {
try {
- Thread.sleep(300);
+ Thread.sleep(OPTIMIZER_HEARTBEAT_INTERVAL_MS);
synchronized (this) {
if (!suspend) {
optimizingService().touch(token);