This is an automated email from the ASF dual-hosted git repository. czy006 pushed a commit to branch jira/JSPT-2568 in repository https://gitbox.apache.org/repos/asf/amoro.git
commit 0b264bb0d6adc4df436021de38a271e5c10f70d3 Author: ConradJam <[email protected]> AuthorDate: Thu Jul 2 14:47:38 2026 +0800 [hotfix][process] 标记恢复失败的 process 为(FAILD)状态位 --- .../amoro/server/process/ProcessService.java | 44 +++++++++++++++++++++- .../process/TestProcessServiceRetryScheduling.java | 32 ++++++++++++++++ 2 files changed, 75 insertions(+), 1 deletion(-) 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 18dbbc72b..ef9a9da0d 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 @@ -41,6 +41,8 @@ import org.apache.amoro.shade.guava32.com.google.common.annotations.VisibleForTe import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.io.PrintWriter; +import java.io.StringWriter; import java.util.Collections; import java.util.List; import java.util.Map; @@ -162,8 +164,9 @@ public class ProcessService extends PersistentBase { ActionCoordinatorScheduler scheduler = actionCoordinators.get(processMeta.getProcessType()); if (tableRuntime != null && scheduler != null) { + DefaultTableProcessStore store = null; try { - DefaultTableProcessStore store = + store = new DefaultTableProcessStore( processMeta.getProcessId(), tableRuntime, @@ -182,11 +185,50 @@ public class ProcessService extends PersistentBase { processMeta.getProcessType(), processMeta.getExecutionEngine(), e); + if (store != null) { + markRecoverFailed(store, processMeta, e); + } } } }); } + private void markRecoverFailed( + DefaultTableProcessStore store, TableProcessMeta processMeta, Exception exception) { + boolean transitioned = + store.tryTransitState( + ProcessStatus.FAILED, + ProcessEvent.COMPLETE_FAILED, + processMeta.getExternalProcessIdentifier(), + buildRecoverFailureMessage(processMeta, exception), + store.getProcessParameters(), + store.getSummary()); + if (!transitioned) { + LOG.warn( + "Failed to mark recovered table process {} as FAILED after recovery exception." + + " tableId={}, processType={}, executionEngine={}, externalProcessIdentifier={}", + processMeta.getProcessId(), + processMeta.getTableId(), + processMeta.getProcessType(), + processMeta.getExecutionEngine(), + processMeta.getExternalProcessIdentifier()); + } + } + + private String buildRecoverFailureMessage(TableProcessMeta processMeta, Exception exception) { + StringWriter sw = new StringWriter(); + exception.printStackTrace(new PrintWriter(sw)); + return String.format( + "Failed to recover active table process %d, tableId=%d, processType=%s," + + " executionEngine=%s, externalProcessIdentifier=%s.%n%s", + processMeta.getProcessId(), + processMeta.getTableId(), + processMeta.getProcessType(), + processMeta.getExecutionEngine(), + processMeta.getExternalProcessIdentifier(), + sw); + } + /** * Execute the process or trace it if not executable. * diff --git a/amoro-ams/src/test/java/org/apache/amoro/server/process/TestProcessServiceRetryScheduling.java b/amoro-ams/src/test/java/org/apache/amoro/server/process/TestProcessServiceRetryScheduling.java index e2a8822f8..13e57d7fa 100644 --- a/amoro-ams/src/test/java/org/apache/amoro/server/process/TestProcessServiceRetryScheduling.java +++ b/amoro-ams/src/test/java/org/apache/amoro/server/process/TestProcessServiceRetryScheduling.java @@ -32,6 +32,8 @@ import org.apache.amoro.process.ProcessStatus; import org.apache.amoro.process.TableProcess; import org.apache.amoro.process.TableProcessStore; import org.apache.amoro.server.AMSManagerTestBase; +import org.apache.amoro.server.persistence.PersistentBase; +import org.apache.amoro.server.persistence.mapper.TableProcessMapper; import org.apache.amoro.table.StateKey; import org.junit.Assert; import org.junit.Test; @@ -105,6 +107,24 @@ public class TestProcessServiceRetryScheduling extends AMSManagerTestBase { processService.recoverProcesses(Collections.singletonList(runtime)); Assert.assertTrue(processService.getActiveTableProcess().isEmpty()); + TableProcessMeta recoveredMeta = getProcessMeta(holder.getStore().getProcessId()); + Assert.assertEquals(ProcessStatus.FAILED, recoveredMeta.getStatus()); + Assert.assertTrue(recoveredMeta.getFinishTime() > 0); + Assert.assertTrue(recoveredMeta.getFailMessage().contains("recover failed for test")); + Assert.assertTrue( + recoveredMeta + .getFailMessage() + .contains(String.valueOf(holder.getStore().getProcessId()))); + Assert.assertTrue( + recoveredMeta + .getFailMessage() + .contains(String.valueOf(runtime.getTableIdentifier().getId()))); + Assert.assertTrue(recoveredMeta.getFailMessage().contains(RECOVER_FAIL_ACTION.getName())); + Assert.assertTrue(recoveredMeta.getFailMessage().contains(executeEngine.name())); + Assert.assertTrue( + recoveredMeta + .getFailMessage() + .contains(holder.getStore().getExternalProcessIdentifier())); } finally { shutdownExecutor(processService, "processExecutionPool"); shutdownExecutor(processService, "retrySchedulingPool"); @@ -149,6 +169,8 @@ public class TestProcessServiceRetryScheduling extends AMSManagerTestBase { Assert.assertTrue( processService.getTableProcessInstances(failedRuntime.getTableIdentifier()).isEmpty()); + Assert.assertEquals( + ProcessStatus.FAILED, getProcessMeta(failedHolder.getStore().getProcessId()).getStatus()); Assert.assertEquals( 1, processService.getTableProcessInstances(recoveredRuntime.getTableIdentifier()).size()); } finally { @@ -218,6 +240,16 @@ public class TestProcessServiceRetryScheduling extends AMSManagerTestBase { } } + private TableProcessMeta getProcessMeta(long processId) { + return new TableProcessMetaReader().getProcessMeta(processId); + } + + private static class TableProcessMetaReader extends PersistentBase { + private TableProcessMeta getProcessMeta(long processId) { + return getAs(TableProcessMapper.class, mapper -> mapper.getProcessMeta(processId)); + } + } + private static class BlockingAwareFailingExecuteEngine implements ExecuteEngine { private final CountDownLatch submissionLatch; private final Map<String, AtomicInteger> statusChecks = new ConcurrentHashMap<>();
