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. */

Reply via email to