This is an automated email from the ASF dual-hosted git repository. xxubai pushed a commit to branch 0.9.x in repository https://gitbox.apache.org/repos/asf/amoro.git
commit 0b77ff162b87b1bdd0907cfde2f70474c708f9cf Author: Power John <[email protected]> AuthorDate: Mon Aug 3 20:04:49 2026 +0800 [AMORO-4303][AMS] Stabilize optimizer heartbeat tests (#4307) (cherry picked from commit 4953d9c90d378d0f031bad3a1bc879a0dc16512c) --- .../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 52345d2cc..cafa1e419 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 @@ -591,7 +591,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);
