This is an automated email from the ASF dual-hosted git repository.

zhoujinsong 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 72a772986 [hotfix][Optimizer] Remove zero-valued in-flight task 
counters (#4328)
72a772986 is described below

commit 72a7729860e284c88206479ffec1ff751751d463
Author: ConradJam <[email protected]>
AuthorDate: Tue Aug 25 14:21:12 2026 +0800

    [hotfix][Optimizer] Remove zero-valued in-flight task counters (#4328)
    
    [hotfix][Optimizer] Remove per-table in-flight task counters at zero
    
    optimizingTasksMap kept a zero-count AtomicInteger entry forever after
    a table's last task completed: poll created entries via
    computeIfAbsent, acceptResult decremented but always returned the
    holder, and neither releaseTable nor dispose cleaned up. Every table
    that ever had a polled task (including deleted tables and tables that
    moved groups) permanently occupied one map slot.
    
    Switch acceptResult to an atomic compute that removes the entry at zero
    (keeping the >0 decrement guard; late accepts on a removed entry are
    no-ops instead of resurrecting it), and clear the entry in
    releaseTable/dispose.
    
    Regression test testTasksMapEntryRemovedAfterTaskCompletes (red before:
    map size stayed 1). Full TestOptimizingQueue 47/47 green.
    Fix record: docs/fix-records/2026-08-16-fix-19-optimizing-tasks-map-leak.md
---
 .../amoro/server/optimizing/OptimizingQueue.java   | 13 ++++++++++--
 .../server/optimizing/TestOptimizingQueue.java     | 24 ++++++++++++++++++++++
 2 files changed, 35 insertions(+), 2 deletions(-)

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 2627078b5..9d579ef24 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
@@ -243,6 +243,9 @@ public class OptimizingQueue extends PersistentBase {
       process.close(false);
       clearProcess(process);
     }
+    // Drop the per-table in-flight counter: the table left this queue, and 
keeping the entry
+    // leaks one map slot per table that ever had a polled task.
+    optimizingTasksMap.remove(tableRuntime.getTableIdentifier());
     LOG.info(
         "Release queue {} with table {}",
         optimizerGroup.getName(),
@@ -511,6 +514,7 @@ public class OptimizingQueue extends PersistentBase {
 
   public void dispose() {
     this.metrics.unregister();
+    this.optimizingTasksMap.clear();
   }
 
   private TableOptimizingProcess findProcess(OptimizingTaskId taskId) {
@@ -697,13 +701,18 @@ public class OptimizingQueue extends PersistentBase {
     private void acceptResult(TaskRuntime<?> taskRuntime) {
       lock.lock();
       try {
-        optimizingTasksMap.computeIfPresent(
+        optimizingTasksMap.compute(
             tableRuntime.getTableIdentifier(),
             (k, v) -> {
+              if (v == null) {
+                return null;
+              }
               if (v.get() > 0) {
                 v.decrementAndGet();
               }
-              return v;
+              // Remove the entry at zero instead of leaving an empty counter 
behind:
+              // every table with a polled task would otherwise occupy a slot 
forever.
+              return v.get() > 0 ? v : null;
             });
         try {
           tableRuntime.addTaskQuota(taskRuntime.getCurrentQuota());
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 37e2c6300..b4d25425f 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
@@ -76,6 +76,7 @@ import org.junit.Test;
 import org.junit.runner.RunWith;
 import org.junit.runners.Parameterized;
 
+import java.lang.reflect.Field;
 import java.util.ArrayList;
 import java.util.Collections;
 import java.util.List;
@@ -956,6 +957,23 @@ public class TestOptimizingQueue extends AMSTableTestBase {
     Assert.assertEquals(OptimizingStatus.IDLE, 
tableRuntime.getOptimizingStatus());
   }
 
+  @Test
+  public void testInFlightCounterEntryRemovedAfterTaskCompletes() throws 
Exception {
+    DefaultTableRuntime tableRuntime = initTableWithFiles();
+    OptimizingQueue queue = buildOptimizingGroupService(tableRuntime);
+    TaskRuntime<?> task = queue.pollTask(optimizerThread, MAX_POLLING_TIME);
+    Assert.assertNotNull(task);
+    Assert.assertEquals(1, readOptimizingTasksMapSize(queue));
+
+    queue.ackTask(task.getTaskId(), optimizerThread);
+    queue.completeTask(
+        optimizerThread,
+        buildOptimizingTaskResult(task.getTaskId(), 
optimizerThread.getThreadId()));
+
+    Assert.assertEquals(0, readOptimizingTasksMapSize(queue));
+    queue.dispose();
+  }
+
   protected DefaultTableRuntime initTableWithFiles() {
     MixedTable mixedTable =
         (MixedTable) 
tableService().loadTable(serverTableIdentifier()).originalTable();
@@ -1065,6 +1083,12 @@ public class TestOptimizingQueue extends 
AMSTableTestBase {
     return optimizingTaskResult;
   }
 
+  private int readOptimizingTasksMapSize(OptimizingQueue queue) throws 
Exception {
+    Field field = OptimizingQueue.class.getDeclaredField("optimizingTasksMap");
+    field.setAccessible(true);
+    return ((Map<?, ?>) field.get(queue)).size();
+  }
+
   /**
    * Simulate the loadOptimizingQueues logic: tables whose persisted optimizer 
group no longer
    * exists (e.g., table config changed from an old deleted group to 
"default", but AMS restarted

Reply via email to