From 31dc7be74a7bb38fc7601f990d9812940cd6c1d7 Mon Sep 17 00:00:00 2001 From: ConradJam Date: Sun, 16 Aug 2026 13:56:46 +0800 Subject: [PATCH] [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 --- .../server/DefaultOptimizingService.java | 5 ++- .../server/resource/OptimizerInstance.java | 4 ++- .../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 eb943e6057..3b59004dc8 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 @@ -1081,8 +1081,11 @@ protected void processTask(OptimizerGroupKeepingTask keepingTask) { .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 0130704d47..70dfcaccdf 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 6e0591d18f..1a820190b9 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 void testMinParallelismResetToZeroWhenNoResource() throws InterruptedExce + ":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 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.