diff --git a/amoro-ams/src/main/java/org/apache/amoro/server/optimizing/OptimizingQueue.java b/amoro-ams/src/main/java/org/apache/amoro/server/optimizing/OptimizingQueue.java index 178bebaa39..2627078b59 100644 --- a/amoro-ams/src/main/java/org/apache/amoro/server/optimizing/OptimizingQueue.java +++ b/amoro-ams/src/main/java/org/apache/amoro/server/optimizing/OptimizingQueue.java @@ -775,7 +775,10 @@ private void resetTask(TaskRuntime taskRuntime) { @Override public boolean isClosed() { - return status == ProcessStatus.KILLED; + // close() sets CLOSED (KILLED is reserved for kill flows); checking only KILLED made + // this predicate permanently false, so the acceptResult guard against late results on a + // closed process could never fire. + return status == ProcessStatus.CLOSED || status == ProcessStatus.KILLED; } @Override diff --git a/amoro-ams/src/test/java/org/apache/amoro/server/optimizing/TestOptimizingQueue.java b/amoro-ams/src/test/java/org/apache/amoro/server/optimizing/TestOptimizingQueue.java index 3955e335e3..37e2c63000 100644 --- a/amoro-ams/src/test/java/org/apache/amoro/server/optimizing/TestOptimizingQueue.java +++ b/amoro-ams/src/test/java/org/apache/amoro/server/optimizing/TestOptimizingQueue.java @@ -858,6 +858,8 @@ public void testProcessCloseKeepsLastOptimizedSnapshotId() { // Close process without success (simulates group change / forced termination) process.close(false); + Assert.assertTrue(process.isClosed()); + // lastOptimizedSnapshotId and lastOptimizedChangeSnapshotId should NOT be updated Assert.assertEquals(snapshotIdBeforePlanning, tableRuntime.getLastOptimizedSnapshotId()); Assert.assertEquals(