This is an automated email from the ASF dual-hosted git repository.
xxubai 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 05c0efe34 [AMORO-4348] Fix infinite retry loop for failed table
processes (#4349)
05c0efe34 is described below
commit 05c0efe346c50b438d3e00e2a2636a347fbd4d9a
Author: Jiwon Park <[email protected]>
AuthorDate: Wed Sep 2 23:59:20 2026 +0900
[AMORO-4348] Fix infinite retry loop for failed table processes (#4349)
* [AMORO-4348] Fix infinite retry loop for failed table processes
Signed-off-by: Jiwon Park <[email protected]>
* [AMORO-4348] Make the failing test engine throw from submitTableProcess
Signed-off-by: Jiwon Park <[email protected]>
---
.../amoro/server/process/ProcessService.java | 25 +++++++----
.../amoro/server/TestDefaultProcessService.java | 51 ++++++++++++++++++++++
2 files changed, 67 insertions(+), 9 deletions(-)
diff --git
a/amoro-ams/src/main/java/org/apache/amoro/server/process/ProcessService.java
b/amoro-ams/src/main/java/org/apache/amoro/server/process/ProcessService.java
index e10d81845..bf44c6480 100644
---
a/amoro-ams/src/main/java/org/apache/amoro/server/process/ProcessService.java
+++
b/amoro-ams/src/main/java/org/apache/amoro/server/process/ProcessService.java
@@ -280,7 +280,7 @@ public class ProcessService extends PersistentBase {
tableRuntime,
processMeta,
scheduler.getAction(),
- processMeta.getRetryNumber());
+ ActionCoordinatorScheduler.PROCESS_MAX_RETRY_NUMBER);
try {
TableProcess process = scheduler.recover(tableRuntime, store);
trackTableProcess(tableRuntime.getTableIdentifier(), store, process);
@@ -367,17 +367,24 @@ public class ProcessService extends PersistentBase {
scheduler =
findScheduler(store.getAction().getName(),
process.getTableRuntime().getFormat());
}
+ boolean retryScheduled = false;
if (scheduler != null
&& store.getStatus() == ProcessStatus.FAILED
&& store.getRetryNumber() <
ActionCoordinatorScheduler.PROCESS_MAX_RETRY_NUMBER
&& process.getTableRuntime() != null) {
- store.tryTransitState(
- ProcessStatus.PENDING,
- ProcessEvent.RETRY_REQUESTED,
- store.getExternalProcessIdentifier(),
- "Regular Retry.",
- process.getProcessParameters(),
- process.getSummary());
+ // Only resubmit when the retry transition is accepted;
resubmitting a process whose
+ // transition was rejected (e.g. already terminal) would loop
forever without ever
+ // increasing the retry number.
+ retryScheduled =
+ store.tryTransitState(
+ ProcessStatus.PENDING,
+ ProcessEvent.RETRY_REQUESTED,
+ store.getExternalProcessIdentifier(),
+ "Regular Retry.",
+ process.getProcessParameters(),
+ process.getSummary());
+ }
+ if (retryScheduled) {
executeOrTraceProcess(store, process);
} else {
if (store.getStatus() == ProcessStatus.FAILED &&
process.getTableRuntime() != null) {
@@ -503,7 +510,7 @@ public class ProcessService extends PersistentBase {
process.getTableRuntime(),
processMeta,
process.getAction(),
- processMeta.getRetryNumber());
+ ActionCoordinatorScheduler.PROCESS_MAX_RETRY_NUMBER);
}
/**
diff --git
a/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultProcessService.java
b/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultProcessService.java
index 0fbcea418..2e06a54eb 100644
---
a/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultProcessService.java
+++
b/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultProcessService.java
@@ -30,6 +30,7 @@ import org.apache.amoro.process.TableProcess;
import org.apache.amoro.process.TableProcessStore;
import org.apache.amoro.server.persistence.PersistentBase;
import org.apache.amoro.server.persistence.mapper.TableProcessMapper;
+import org.apache.amoro.server.process.ActionCoordinatorScheduler;
import org.apache.amoro.server.process.MockActionCoordinator;
import org.apache.amoro.server.process.MockExecuteEngine;
import org.apache.amoro.server.process.ProcessService;
@@ -392,6 +393,56 @@ public class TestDefaultProcessService extends
AMSTableTestBase {
}
}
+ /**
+ * 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 submissions always fail. */
+ private static class AlwaysFailingExecuteEngine extends MockExecuteEngine {
+ private final AtomicInteger submitAttempts = new AtomicInteger();
+
+ @Override
+ public String submitTableProcess(org.apache.amoro.process.TableProcess
tableProcess) {
+ submitAttempts.incrementAndGet();
+ throw new IllegalStateException("Submission failure");
+ }
+
+ private int getSubmitAttempts() {
+ return submitAttempts.get();
+ }
+ }
+
// ---------------------- Private helpers ----------------------
/** Return the first available execute engine and validate its presence. */