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.