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