j1wonpark commented on code in PR #4349:
URL: https://github.com/apache/amoro/pull/4349#discussion_r3914526660
##########
amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultProcessService.java:
##########
@@ -392,6 +395,66 @@ public void testRecoverProcessFailSafe() {
}
}
+ /**
+ * Verify that a process which keeps failing is retried a bounded number of
times and then
+ * dropped, instead of being resubmitted in an infinite loop. The process
store must be created
+ * with the retry limit rather than the current retry count, otherwise a
FAILED process is treated
+ * as terminal immediately (retryNumber >= maxRetryTime), every
RETRY_REQUESTED transition is
+ * rejected, retryNumber never increases and the process is resubmitted
forever.
+ */
+ @Test(timeout = 60_000)
+ public void testFailedProcessRetryIsBounded() throws InterruptedException {
+ AlwaysFailingExecuteEngine failingEngine = new
AlwaysFailingExecuteEngine();
+ processServiceService().unInstallAllExecuteEngines();
+ processServiceService().installExecuteEngine(failingEngine);
+ processServiceService().unInstallAllActionCoordinators();
+ processServiceService().installActionCoordinator(new
MockActionCoordinator(failingEngine));
+ try {
+ createTable();
+
+ // Wait until the failing process has been submitted at least once.
+ awaitCondition(() -> failingEngine.getSubmitAttempts() >= 1,
WAIT_TIMEOUT_MS, 100L);
+
+ // The process must exhaust its bounded retries and be untracked. With
the infinite
+ // retry loop, the active process map never empties and this wait times
out.
+ awaitCondition(
+ () -> processServiceService().getActiveTableProcess().isEmpty(),
15_000L, 100L);
+
+ Assert.assertEquals(
+ 1 + ActionCoordinatorScheduler.PROCESS_MAX_RETRY_NUMBER,
+ failingEngine.getSubmitAttempts());
+
+ dropTable();
+ } catch (Throwable t) {
+ throw new RuntimeException(t);
+ }
+ }
+
+ /** Execute engine whose processes always fail immediately after submission.
*/
+ private static class AlwaysFailingExecuteEngine extends MockExecuteEngine {
+ private final AtomicInteger submitAttempts = new AtomicInteger();
+ private final Set<String> failedIdentifiers =
ConcurrentHashMap.newKeySet();
+
+ @Override
+ public String submitTableProcess(org.apache.amoro.process.TableProcess
tableProcess) {
+ String identifier = "failing-" + submitAttempts.incrementAndGet();
+ failedIdentifiers.add(identifier);
+ return identifier;
+ }
+
+ @Override
+ public ProcessStatus getStatus(String processIdentifier) {
+ if (failedIdentifiers.contains(processIdentifier)) {
Review Comment:
Done.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]