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 7333642cf [hotfix][AMS] Clean up orphan optimizer records at startup
(#4320)
7333642cf is described below
commit 7333642cfb92e391fc4280553f673ea9829c6644
Author: ConradJam <[email protected]>
AuthorDate: Tue Aug 25 11:23:26 2026 +0800
[hotfix][AMS] Clean up orphan optimizer records at startup (#4320)
* [hotfix][AMS] Clean up orphan optimizer records at startup instead of NPE
An optimizer row referencing a resource group that no longer exists
(e.g. group dropped while AMS was down, or the canDeleteResourceGroup
check-then-act race) made loadOptimizingQueues throw an NPE during
initHandler and prevented AMS from starting, with no self-healing path.
Extract registerOptimizers(): register when the group exists, otherwise
WARN and remove the orphan row via the existing unregisterOptimizer
cleanup, mirroring the tolerant handling already present in the
follower-sync path (registerOptimizerWithoutPersist).
Verified by new test testOrphanedOptimizerRecordMustNotBreakInitialization
(red: NPE reproduced on old behavior; green: 26/26 in
TestDefaultOptimizingService).
Fix record:
docs/fix-records/2026-08-15-fix-01-orphan-optimizer-startup-npe.md
* [hotfix][AMS] Prevent optimizer resource group races in HA mode
---
.../amoro/server/DefaultOptimizingService.java | 48 ++++-
.../server/persistence/mapper/OptimizerMapper.java | 7 +
.../server/persistence/mapper/ResourceMapper.java | 7 +
.../server/resource/DefaultOptimizerManager.java | 23 ++-
.../amoro/server/TestDefaultOptimizingService.java | 113 ++++++++++++
.../TestOptimizerResourceGroupConcurrency.java | 203 +++++++++++++++++++++
6 files changed, 396 insertions(+), 5 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 eb943e605..f2aee8e4a 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
@@ -212,7 +212,7 @@ public class DefaultOptimizingService extends
StatedPersistentBase
optimizerGroupKeeper.keepInTouch(groupName, 1);
optimizerScaleKeeper.watch(group);
});
- optimizers.forEach(optimizer -> registerOptimizer(optimizer, false));
+ registerOptimizers(optimizers);
// Avoid keeping the tables in processing/pending status forever in below
cases:
// 1) Resource group does not exist
// 2) The AMS restarts after the tables disable self-optimizing but before
the optimizing
@@ -239,7 +239,17 @@ public class DefaultOptimizingService extends
StatedPersistentBase
private void registerOptimizer(OptimizerInstance optimizer, boolean
needPersistent) {
if (needPersistent) {
- doAs(OptimizerMapper.class, mapper -> mapper.insertOptimizer(optimizer));
+ doAsTransaction(
+ () -> {
+ String groupName =
+ getAs(
+ ResourceMapper.class,
+ mapper ->
mapper.selectResourceGroupNameForUpdate(optimizer.getGroupName()));
+ if (groupName == null) {
+ throw new ObjectNotExistsException("Optimizer group " +
optimizer.getGroupName());
+ }
+ doAs(OptimizerMapper.class, mapper ->
mapper.insertOptimizer(optimizer));
+ });
}
OptimizingQueue optimizingQueue =
optimizingQueueByGroup.get(optimizer.getGroupName());
@@ -250,6 +260,40 @@ public class DefaultOptimizingService extends
StatedPersistentBase
optimizerScaleKeeper.onOptimizerRegistered(optimizer);
}
+ /**
+ * Registers optimizers recovered from persistence at startup. A missing
local queue may be a
+ * stale snapshot in an HA deployment, so a non-empty optimizer group is
removed only when its
+ * persisted resource group is also absent. Empty group names are always
treated as orphaned.
+ */
+ void registerOptimizers(List<OptimizerInstance> optimizers) {
+ for (OptimizerInstance optimizer : optimizers) {
+ String groupName = optimizer.getGroupName();
+ if (groupName != null
+ && !groupName.isEmpty()
+ && optimizingQueueByGroup.containsKey(groupName)) {
+ registerOptimizer(optimizer, false);
+ } else {
+ long deleted =
+ updateAs(
+ OptimizerMapper.class,
+ mapper ->
mapper.deleteOptimizerIfResourceGroupAbsent(optimizer.getToken()));
+ if (deleted == 1) {
+ LOG.warn(
+ "Remove orphan optimizer {}, its resource group {} does not
exist",
+ optimizer.getToken(),
+ groupName);
+ } else {
+ LOG.warn(
+ "Skip recovering optimizer {} due to a stale local resource
group snapshot:"
+ + " group {} is unavailable locally, but keep its shared
record because the"
+ + " resource group still exists",
+ optimizer.getToken(),
+ groupName);
+ }
+ }
+ }
+ }
+
private void unregisterOptimizer(String token) {
doAs(OptimizerMapper.class, mapper -> mapper.deleteOptimizer(token));
OptimizingQueue optimizingQueue = optimizingQueueByToken.remove(token);
diff --git
a/amoro-ams/src/main/java/org/apache/amoro/server/persistence/mapper/OptimizerMapper.java
b/amoro-ams/src/main/java/org/apache/amoro/server/persistence/mapper/OptimizerMapper.java
index 3a0c10d5d..abe45e067 100644
---
a/amoro-ams/src/main/java/org/apache/amoro/server/persistence/mapper/OptimizerMapper.java
+++
b/amoro-ams/src/main/java/org/apache/amoro/server/persistence/mapper/OptimizerMapper.java
@@ -52,6 +52,13 @@ public interface OptimizerMapper {
@Delete("DELETE FROM optimizer WHERE token = #{token}")
void deleteOptimizer(@Param("token") String token);
+ @Delete(
+ "DELETE FROM optimizer WHERE token = #{token}"
+ + " AND (optimizer.group_name IS NULL OR optimizer.group_name = ''"
+ + " OR NOT EXISTS (SELECT 1 FROM resource_group"
+ + " WHERE resource_group.group_name = optimizer.group_name))")
+ int deleteOptimizerIfResourceGroupAbsent(@Param("token") String token);
+
@Select(
"SELECT token, resource_id, group_name, container_name, start_time,
touch_time,"
+ "thread_count, total_memory, properties FROM optimizer")
diff --git
a/amoro-ams/src/main/java/org/apache/amoro/server/persistence/mapper/ResourceMapper.java
b/amoro-ams/src/main/java/org/apache/amoro/server/persistence/mapper/ResourceMapper.java
index 319b214f5..dbe49eae1 100644
---
a/amoro-ams/src/main/java/org/apache/amoro/server/persistence/mapper/ResourceMapper.java
+++
b/amoro-ams/src/main/java/org/apache/amoro/server/persistence/mapper/ResourceMapper.java
@@ -58,6 +58,13 @@ public interface ResourceMapper {
})
ResourceGroup selectResourceGroup(@Param("resourceGroup") String groupName);
+ @Select(
+ "<script>"
+ + "SELECT group_name FROM resource_group WHERE group_name =
#{resourceGroup} FOR UPDATE"
+ + "<if test=\"_databaseId == 'derby'\"> WITH RS</if>"
+ + "</script>")
+ String selectResourceGroupNameForUpdate(@Param("resourceGroup") String
groupName);
+
@Select(
"SELECT resource_id, group_name, container_name, thread_count,
total_memory, properties"
+ " FROM resource WHERE group_name = #{resourceGroup}")
diff --git
a/amoro-ams/src/main/java/org/apache/amoro/server/resource/DefaultOptimizerManager.java
b/amoro-ams/src/main/java/org/apache/amoro/server/resource/DefaultOptimizerManager.java
index b916e319c..d139e2604 100644
---
a/amoro-ams/src/main/java/org/apache/amoro/server/resource/DefaultOptimizerManager.java
+++
b/amoro-ams/src/main/java/org/apache/amoro/server/resource/DefaultOptimizerManager.java
@@ -33,6 +33,7 @@ import
org.apache.amoro.shade.guava32.com.google.common.base.Preconditions;
import org.apache.amoro.table.TableProperties;
import java.util.List;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.stream.Collectors;
/**
@@ -95,9 +96,25 @@ public class DefaultOptimizerManager extends PersistentBase
implements Optimizer
@Override
public void deleteResourceGroup(String groupName) {
- if (canDeleteResourceGroup(groupName)) {
- doAs(ResourceMapper.class, mapper ->
mapper.deleteResourceGroup(groupName));
- } else {
+ AtomicBoolean groupInUse = new AtomicBoolean(false);
+ doAsTransaction(
+ () -> {
+ String lockedGroupName =
+ getAs(
+ ResourceMapper.class,
+ mapper ->
mapper.selectResourceGroupNameForUpdate(groupName));
+ if (lockedGroupName == null) {
+ return;
+ }
+
+ if (!canDeleteResourceGroup(groupName)) {
+ groupInUse.set(true);
+ return;
+ }
+ doAs(ResourceMapper.class, mapper ->
mapper.deleteResourceGroup(groupName));
+ });
+ // Keep the existing exception type instead of letting the transaction
helper wrap it.
+ if (groupInUse.get()) {
throw new RuntimeException(
String.format(
"The resource group %s cannot be deleted because it is currently
in use.",
diff --git
a/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultOptimizingService.java
b/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultOptimizingService.java
index 5c9d52554..e146a57ee 100644
---
a/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultOptimizingService.java
+++
b/amoro-ams/src/test/java/org/apache/amoro/server/TestDefaultOptimizingService.java
@@ -37,6 +37,7 @@ import org.apache.amoro.catalog.BasicCatalogTestHelper;
import org.apache.amoro.catalog.CatalogTestHelper;
import org.apache.amoro.config.OptimizingConfig;
import org.apache.amoro.config.TableConfiguration;
+import org.apache.amoro.exception.ObjectNotExistsException;
import org.apache.amoro.exception.PluginRetryAuthException;
import org.apache.amoro.exception.TaskRuntimeException;
import org.apache.amoro.io.MixedDataTestHelpers;
@@ -53,6 +54,7 @@ import org.apache.amoro.server.optimizing.OptimizingStatus;
import org.apache.amoro.server.optimizing.TaskRuntime;
import org.apache.amoro.server.persistence.SqlSessionFactoryProvider;
import org.apache.amoro.server.persistence.TableRuntimeMeta;
+import org.apache.amoro.server.persistence.mapper.OptimizerMapper;
import org.apache.amoro.server.persistence.mapper.TableProcessMapper;
import org.apache.amoro.server.persistence.mapper.TableRuntimeMapper;
import org.apache.amoro.server.process.TableProcessMeta;
@@ -173,6 +175,100 @@ public class TestDefaultOptimizingService extends
AMSTableTestBase {
PluginRetryAuthException.class, () ->
optimizingService().pollTask("whatever", THREAD_ID));
}
+ @Test
+ public void testOrphanedOptimizerRecordMustNotBreakInitialization() {
+ // An optimizer row whose resource group no longer exists (e.g. the group
was dropped while
+ // AMS was down) must be cleaned up instead of failing initialization with
an NPE.
+ OptimizerRegisterInfo registerInfo = buildRegisterInfo();
+ registerInfo.setGroupName("group-dropped-while-ams-down");
+ OptimizerInstance orphan = new OptimizerInstance(registerInfo, "local");
+ OptimizerRegisterInfo emptyGroupRegisterInfo = buildRegisterInfo();
+ emptyGroupRegisterInfo.setGroupName("");
+ emptyGroupRegisterInfo.setResourceId("resource-with-empty-group");
+ OptimizerInstance emptyGroupOrphan = new
OptimizerInstance(emptyGroupRegisterInfo, "local");
+ optimizerManager().createResourceGroup(new ResourceGroup.Builder("",
"local").build());
+ try (SqlSession session =
SqlSessionFactoryProvider.getInstance().get().openSession(true)) {
+ session.getMapper(OptimizerMapper.class).insertOptimizer(orphan);
+
session.getMapper(OptimizerMapper.class).insertOptimizer(emptyGroupOrphan);
+ }
+
+ try {
+ // Exercise the production startup path rather than calling the recovery
helper directly.
+ Assertions.assertDoesNotThrow(this::reload);
+
+ List<OptimizerInstance> remaining;
+ try (SqlSession session =
SqlSessionFactoryProvider.getInstance().get().openSession(true)) {
+ remaining = session.getMapper(OptimizerMapper.class).selectAll();
+ }
+ Assertions.assertFalse(
+ remaining.stream().anyMatch(o ->
o.getToken().equals(orphan.getToken())),
+ "orphaned optimizer record should be removed during initialization");
+ Assertions.assertFalse(
+ remaining.stream().anyMatch(o ->
o.getToken().equals(emptyGroupOrphan.getToken())),
+ "an optimizer record with an empty group should always be removed");
+ } finally {
+ optimizingService().deleteOptimizer("",
emptyGroupOrphan.getResourceId());
+ optimizingService().deleteResourceGroup("");
+ optimizerManager().deleteResourceGroup("");
+ }
+ }
+
+ @Test
+ public void testValidOptimizerRecordMustSurviveStaleQueueSnapshot() {
+ String groupName = "group-created-by-another-ams";
+ ResourceGroup group = new ResourceGroup.Builder(groupName,
"local").build();
+ optimizerManager().createResourceGroup(group);
+ optimizingService().createResourceGroup(group);
+ // Simulate a startup snapshot taken before another AMS created the
persisted group.
+ optimizingService().deleteResourceGroup(groupName);
+
+ OptimizerRegisterInfo registerInfo = buildRegisterInfo();
+ registerInfo.setGroupName(groupName);
+ registerInfo.setResourceId("resource-created-by-another-ams");
+ OptimizerInstance optimizer = new OptimizerInstance(registerInfo, "local");
+ insertOptimizer(optimizer);
+
+ try {
+ optimizingService().registerOptimizers(Lists.newArrayList(optimizer));
+
+ Assertions.assertTrue(
+ optimizerExists(optimizer.getToken()),
+ "a stale local queue snapshot must not delete a valid shared
optimizer record");
+ } finally {
+ deleteOptimizerRecord(optimizer.getToken());
+ optimizerManager().deleteResourceGroup(groupName);
+ }
+ }
+
+ @Test
+ public void testAuthenticateMustRejectStaleLocalQueueAfterGroupDeletion() {
+ String groupName = "group-deleted-by-another-ams";
+ ResourceGroup group = new ResourceGroup.Builder(groupName,
"local").build();
+ optimizerManager().createResourceGroup(group);
+ optimizingService().createResourceGroup(group);
+ // Keep the local queue but remove the shared database row, as can happen
before watcher sync.
+ optimizerManager().deleteResourceGroup(groupName);
+
+ OptimizerRegisterInfo registerInfo = buildRegisterInfo();
+ registerInfo.setGroupName(groupName);
+ registerInfo.setResourceId("resource-for-deleted-group");
+
+ try {
+ Assertions.assertThrows(
+ ObjectNotExistsException.class, () ->
optimizingService().authenticate(registerInfo));
+ Assertions.assertFalse(
+ optimizerManager().listOptimizers().stream()
+ .anyMatch(optimizer ->
groupName.equals(optimizer.getGroupName())),
+ "authentication must not persist an optimizer for a deleted resource
group");
+ } finally {
+ optimizerManager().listOptimizers().stream()
+ .filter(optimizer -> groupName.equals(optimizer.getGroupName()))
+ .map(OptimizerInstance::getToken)
+ .forEach(this::deleteOptimizerRecord);
+ optimizingService().deleteResourceGroup(groupName);
+ }
+ }
+
@Test
public void testPollOnce() {
// 1.poll task
@@ -832,6 +928,23 @@ public class TestDefaultOptimizingService extends
AMSTableTestBase {
return registerInfo;
}
+ private void insertOptimizer(OptimizerInstance optimizer) {
+ try (SqlSession session =
SqlSessionFactoryProvider.getInstance().get().openSession(true)) {
+ session.getMapper(OptimizerMapper.class).insertOptimizer(optimizer);
+ }
+ }
+
+ private boolean optimizerExists(String optimizerToken) {
+ return optimizerManager().listOptimizers().stream()
+ .anyMatch(optimizer -> optimizerToken.equals(optimizer.getToken()));
+ }
+
+ private void deleteOptimizerRecord(String optimizerToken) {
+ try (SqlSession session =
SqlSessionFactoryProvider.getInstance().get().openSession(true)) {
+ session.getMapper(OptimizerMapper.class).deleteOptimizer(optimizerToken);
+ }
+ }
+
@SuppressWarnings("unchecked")
private Map<String, ?> optimizerState(String fieldName) throws Exception {
Field field = DefaultOptimizingService.class.getDeclaredField(fieldName);
diff --git
a/amoro-ams/src/test/java/org/apache/amoro/server/TestOptimizerResourceGroupConcurrency.java
b/amoro-ams/src/test/java/org/apache/amoro/server/TestOptimizerResourceGroupConcurrency.java
new file mode 100644
index 000000000..dabd829a3
--- /dev/null
+++
b/amoro-ams/src/test/java/org/apache/amoro/server/TestOptimizerResourceGroupConcurrency.java
@@ -0,0 +1,203 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.amoro.server;
+
+import org.apache.amoro.api.OptimizerRegisterInfo;
+import org.apache.amoro.exception.ObjectNotExistsException;
+import org.apache.amoro.resource.ResourceGroup;
+import org.apache.amoro.server.persistence.SqlSessionFactoryProvider;
+import org.apache.amoro.server.persistence.mapper.OptimizerMapper;
+import org.apache.amoro.server.persistence.mapper.ResourceMapper;
+import org.apache.amoro.server.resource.OptimizerInstance;
+import org.apache.ibatis.session.SqlSession;
+import org.junit.Test;
+import org.junit.jupiter.api.Assertions;
+
+import java.util.HashMap;
+import java.util.UUID;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+
+public class TestOptimizerResourceGroupConcurrency extends AMSServiceTestBase {
+
+ private static final long BLOCKED_ASSERTION_TIMEOUT_MS = 500;
+ private static final long COMPLETION_TIMEOUT_SECONDS = 5;
+
+ @Test
+ public void testDeletingMissingGroupRemainsIdempotent() {
+ String groupName = "missing-" + UUID.randomUUID();
+
+ Assertions.assertDoesNotThrow(() ->
optimizerManager().deleteResourceGroup(groupName));
+ Assertions.assertDoesNotThrow(() ->
optimizerManager().deleteResourceGroup(groupName));
+ Assertions.assertNull(optimizerManager().getResourceGroup(groupName));
+ }
+
+ @Test(timeout = 30_000)
+ public void testDeleteFirstMakesConcurrentRegistrationFail() throws
Exception {
+ String groupName = "d-first-" + UUID.randomUUID();
+ createGroup(groupName);
+ OptimizerRegisterInfo registerInfo = buildRegisterInfo(groupName);
+ CountDownLatch registrationStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newSingleThreadExecutor();
+ Future<String> registration = null;
+
+ try {
+ try (SqlSession deleteSession = openSession(false)) {
+ ResourceMapper deleteMapper =
deleteSession.getMapper(ResourceMapper.class);
+ Assertions.assertEquals(
+ groupName,
deleteMapper.selectResourceGroupNameForUpdate(groupName));
+
+ registration =
+ executor.submit(
+ () -> {
+ registrationStarted.countDown();
+ return optimizingService().authenticate(registerInfo);
+ });
+ Assertions.assertTrue(
+ registrationStarted.await(COMPLETION_TIMEOUT_SECONDS,
TimeUnit.SECONDS));
+ Future<String> blockedRegistration = registration;
+ Assertions.assertThrows(
+ TimeoutException.class,
+ () -> blockedRegistration.get(BLOCKED_ASSERTION_TIMEOUT_MS,
TimeUnit.MILLISECONDS));
+
+ deleteMapper.deleteResourceGroup(groupName);
+ deleteSession.commit();
+ }
+
+ Future<String> completedRegistration = registration;
+ ExecutionException failure =
+ Assertions.assertThrows(
+ ExecutionException.class,
+ () -> completedRegistration.get(COMPLETION_TIMEOUT_SECONDS,
TimeUnit.SECONDS));
+ Assertions.assertEquals(ObjectNotExistsException.class,
failure.getCause().getClass());
+ Assertions.assertNull(optimizerManager().getResourceGroup(groupName));
+
Assertions.assertTrue(optimizerManager().listOptimizers(groupName).isEmpty());
+ } finally {
+ if (registration != null) {
+ registration.cancel(true);
+ }
+ try {
+ shutdown(executor);
+ } finally {
+ cleanGroup(groupName);
+ }
+ }
+ }
+
+ @Test(timeout = 30_000)
+ public void testRegistrationFirstMakesConcurrentDeletionFail() throws
Exception {
+ String groupName = "r-first-" + UUID.randomUUID();
+ createGroup(groupName);
+ OptimizerInstance optimizer = new
OptimizerInstance(buildRegisterInfo(groupName), "local");
+ CountDownLatch deletionStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newSingleThreadExecutor();
+ Future<?> deletion = null;
+
+ try {
+ try (SqlSession registrationSession = openSession(false)) {
+ ResourceMapper resourceMapper =
registrationSession.getMapper(ResourceMapper.class);
+ Assertions.assertEquals(
+ groupName,
resourceMapper.selectResourceGroupNameForUpdate(groupName));
+
registrationSession.getMapper(OptimizerMapper.class).insertOptimizer(optimizer);
+
+ deletion =
+ executor.submit(
+ () -> {
+ deletionStarted.countDown();
+ optimizerManager().deleteResourceGroup(groupName);
+ });
+
Assertions.assertTrue(deletionStarted.await(COMPLETION_TIMEOUT_SECONDS,
TimeUnit.SECONDS));
+ Future<?> blockedDeletion = deletion;
+ Assertions.assertThrows(
+ TimeoutException.class,
+ () -> blockedDeletion.get(BLOCKED_ASSERTION_TIMEOUT_MS,
TimeUnit.MILLISECONDS));
+
+ registrationSession.commit();
+ }
+
+ Future<?> completedDeletion = deletion;
+ ExecutionException failure =
+ Assertions.assertThrows(
+ ExecutionException.class,
+ () -> completedDeletion.get(COMPLETION_TIMEOUT_SECONDS,
TimeUnit.SECONDS));
+ Assertions.assertEquals(RuntimeException.class,
failure.getCause().getClass());
+
Assertions.assertTrue(failure.getCause().getMessage().contains("currently in
use"));
+ Assertions.assertNotNull(optimizerManager().getResourceGroup(groupName));
+ Assertions.assertTrue(
+ optimizerManager().listOptimizers(groupName).stream()
+ .anyMatch(record ->
optimizer.getToken().equals(record.getToken())));
+ } finally {
+ if (deletion != null) {
+ deletion.cancel(true);
+ }
+ try {
+ shutdown(executor);
+ } finally {
+ cleanGroup(groupName);
+ }
+ }
+ }
+
+ private void createGroup(String groupName) {
+ ResourceGroup resourceGroup = new ResourceGroup.Builder(groupName,
"local").build();
+ optimizerManager().createResourceGroup(resourceGroup);
+ optimizingService().createResourceGroup(resourceGroup);
+ }
+
+ private OptimizerRegisterInfo buildRegisterInfo(String groupName) {
+ OptimizerRegisterInfo registerInfo = new OptimizerRegisterInfo();
+ registerInfo.setProperties(new HashMap<>());
+ registerInfo.setThreadCount(1);
+ registerInfo.setMemoryMb(1024);
+ registerInfo.setGroupName(groupName);
+ registerInfo.setResourceId("resource-" + UUID.randomUUID());
+ registerInfo.setStartTime(System.currentTimeMillis());
+ return registerInfo;
+ }
+
+ private SqlSession openSession(boolean autoCommit) {
+ return
SqlSessionFactoryProvider.getInstance().get().openSession(autoCommit);
+ }
+
+ private void cleanGroup(String groupName) {
+ try {
+ optimizingService().deleteResourceGroup(groupName);
+ } catch (RuntimeException ignored) {
+ // The queue may already have been removed by the group watcher.
+ }
+ try (SqlSession session = openSession(true)) {
+ OptimizerMapper optimizerMapper =
session.getMapper(OptimizerMapper.class);
+ optimizerMapper.selectAll().stream()
+ .filter(optimizer -> groupName.equals(optimizer.getGroupName()))
+ .map(OptimizerInstance::getToken)
+ .forEach(optimizerMapper::deleteOptimizer);
+ session.getMapper(ResourceMapper.class).deleteResourceGroup(groupName);
+ }
+ }
+
+ private void shutdown(ExecutorService executor) throws InterruptedException {
+ executor.shutdownNow();
+
Assertions.assertTrue(executor.awaitTermination(COMPLETION_TIMEOUT_SECONDS,
TimeUnit.SECONDS));
+ }
+}