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 4f26c4574 [hotfix][Optimizer] Preserve keeper visibility on lookup
failure (#4325)
4f26c4574 is described below
commit 4f26c4574c5200d4904822408db7a96f4f59326c
Author: ConradJam <[email protected]>
AuthorDate: Wed Aug 26 10:53:22 2026 +0800
[hotfix][Optimizer] Preserve keeper visibility on lookup failure (#4325)
[hotfix][Optimizer] Volatile optimizer touch time; keep groups watched on
container lookup failure
touchTime is written by thrift heartbeat threads and read by the keeper
thread without a shared lock - a plain long field is a JMM data race
(stale reads can expire live optimizers; word tearing on 32-bit JVMs).
Declare it volatile.
OptimizerGroupKeeper.processTask called Containers.get outside the
try/finally whose finally-block re-queues the group via keepInTouch.
An unknown container name threw IllegalArgumentException straight into
AbstractKeeper.run's swallow-all catch, permanently removing the group
from scale-out monitoring until restart. Move the lookup and the cast
inside the try.
Regression test testUnknownContainerKeepsGroupWatchedAndResetsMinParallelism
(red before: min-parallelism stayed 2, group silently dropped).
Fix record:
docs/fix-records/2026-08-16-fix-15-keeper-visibility-and-container-lookup.md
---
.../amoro/server/DefaultOptimizingService.java | 5 +++-
.../amoro/server/resource/OptimizerInstance.java | 4 ++-
.../amoro/server/TestOptimizerGroupKeeper.java | 31 ++++++++++++++++++++++
3 files changed, 38 insertions(+), 2 deletions(-)
diff --git
a/amoro-ams/src/main/java/org/apache/amoro/server/DefaultOptimizingService.java
b/amoro-ams/src/main/java/org/apache/amoro/server/DefaultOptimizingService.java
index f2aee8e4a..afccf8a97 100644
---
a/amoro-ams/src/main/java/org/apache/amoro/server/DefaultOptimizingService.java
+++
b/amoro-ams/src/main/java/org/apache/amoro/server/DefaultOptimizingService.java
@@ -1125,8 +1125,11 @@ public class DefaultOptimizingService extends
StatedPersistentBase
.setProperties(resourceGroup.getProperties())
.setThreadCount(requiredCores)
.build();
- ResourceContainer rc = Containers.get(resource.getContainerName());
try {
+ // Containers.get throws for an unknown container name; it must stay
inside the try so
+ // the finally-block keepInTouch still re-queues the group - otherwise
a single lookup
+ // failure silently removes the group from scale-out monitoring until
restart.
+ ResourceContainer rc = Containers.get(resource.getContainerName());
((AbstractOptimizerContainer) rc).requestResource(resource);
optimizerManager.createResource(resource);
} finally {
diff --git
a/amoro-ams/src/main/java/org/apache/amoro/server/resource/OptimizerInstance.java
b/amoro-ams/src/main/java/org/apache/amoro/server/resource/OptimizerInstance.java
index 0130704d4..70dfcaccd 100644
---
a/amoro-ams/src/main/java/org/apache/amoro/server/resource/OptimizerInstance.java
+++
b/amoro-ams/src/main/java/org/apache/amoro/server/resource/OptimizerInstance.java
@@ -28,7 +28,9 @@ public class OptimizerInstance extends Resource {
private String token;
private long startTime;
- private long touchTime;
+ // Written by thrift heartbeat threads (touch) and read by the keeper thread
for expiry
+ // detection without a shared lock; volatile guarantees the keeper observes
fresh heartbeats.
+ private volatile long touchTime;
public OptimizerInstance() {}
diff --git
a/amoro-ams/src/test/java/org/apache/amoro/server/TestOptimizerGroupKeeper.java
b/amoro-ams/src/test/java/org/apache/amoro/server/TestOptimizerGroupKeeper.java
index 6e0591d18..1a820190b 100644
---
a/amoro-ams/src/test/java/org/apache/amoro/server/TestOptimizerGroupKeeper.java
+++
b/amoro-ams/src/test/java/org/apache/amoro/server/TestOptimizerGroupKeeper.java
@@ -276,6 +276,37 @@ public class TestOptimizerGroupKeeper extends
AMSTableTestBase {
+ ":min-parallelism should be reset to 0 when no resources
available and no optimizer exists");
}
+ @Test
+ public void testUnknownContainerKeepsGroupWatchedAndResetsMinParallelism()
+ throws InterruptedException {
+ // Containers.get throws for an unknown container name. The lookup used to
sit outside the
+ // try/finally, so the exception skipped keepInTouch and silently removed
the group from
+ // scale-out monitoring forever. The keeper must keep watching and
eventually reset
+ // min-parallelism like any other permanently-failing scale-out.
+ scaleOutCallCount.set(0);
+ String groupName = TEST_GROUP_NAME + "-7";
+ this.currentGroupName = groupName;
+ Map<String, String> properties = Maps.newHashMap();
+ properties.put(OptimizerProperties.OPTIMIZER_GROUP_MIN_PARALLELISM, "2");
+ properties.put("memory", "1024");
+ ResourceGroup resourceGroup =
+ new ResourceGroup.Builder(groupName, "unknown-container-x")
+ .addProperties(properties)
+ .build();
+
+ optimizerManager().createResourceGroup(resourceGroup);
+ optimizingService().createResourceGroup(resourceGroup);
+
+ Thread.sleep(300);
+
+ ResourceGroup updatedGroup =
optimizerManager().getResourceGroup(groupName);
+ Assertions.assertEquals(
+ "0",
+
updatedGroup.getProperties().get(OptimizerProperties.OPTIMIZER_GROUP_MIN_PARALLELISM),
+ groupName
+ + ":keeper must keep watching an unknown-container group and reset
min-parallelism");
+ }
+
/**
* Test scenario 4: When no resources but has optimizer, min-parallelism
will be reset to
* optimizer's executionParallel.