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 de6205c0e [AMORO-4264][AMS] Use a shared pool for Iceberg maintenance
planning (#4311)
de6205c0e is described below
commit de6205c0e92bbf3139e25067d735f61ba95b7f41
Author: Power John <[email protected]>
AuthorDate: Tue Aug 18 15:35:53 2026 +0800
[AMORO-4264][AMS] Use a shared pool for Iceberg maintenance planning (#4311)
* [AMORO-4264][AMS] Isolate snapshot expiration planning
* [AMORO-4264][AMS] Use a shared pool for Iceberg maintenance planning
---------
Co-authored-by: ZhouJinsong <[email protected]>
---
.../apache/amoro/server/AmoroManagementConf.java | 7 +++
.../apache/amoro/server/AmoroServiceContainer.java | 7 +--
.../amoro/server/TestAmoroManagementConf.java | 8 +++
.../iceberg/maintainer/IcebergTableMaintainer.java | 2 +
.../org/apache/amoro/utils/IcebergThreadPools.java | 30 +++++++++++
.../maintainer/TestIcebergTableMaintainer.java | 62 ++++++++++++++++++++++
.../apache/amoro/utils/TestIcebergThreadPools.java | 14 +++++
charts/amoro/templates/amoro-configmap.yaml | 1 +
dist/src/main/amoro-bin/conf/config.yaml | 1 +
docs/configuration/ams-config.md | 1 +
10 files changed, 130 insertions(+), 3 deletions(-)
diff --git
a/amoro-ams/src/main/java/org/apache/amoro/server/AmoroManagementConf.java
b/amoro-ams/src/main/java/org/apache/amoro/server/AmoroManagementConf.java
index f38f6f5ac..824ee080d 100644
--- a/amoro-ams/src/main/java/org/apache/amoro/server/AmoroManagementConf.java
+++ b/amoro-ams/src/main/java/org/apache/amoro/server/AmoroManagementConf.java
@@ -121,6 +121,13 @@ public class AmoroManagementConf {
"Sets the size of the worker pool. The worker pool limits the
number of tasks concurrently processing "
+ "manifests in the base table implementation across all
concurrent planning or commit operations.");
+ public static final ConfigOption<Integer>
TABLE_MANIFEST_IO_MAINTENANCE_THREAD_COUNT =
+ ConfigOptions.key("table-manifest-io.maintenance-thread-count")
+ .intType()
+ .defaultValue(10)
+ .withDescription(
+ "Sets the size of the worker pool used for manifest I/O across
best-effort table maintenance operations.");
+
public static final ConfigOption<Integer>
TABLE_MANIFEST_IO_PLANNING_THREAD_COUNT =
ConfigOptions.key("self-optimizing.plan-manifest-io-thread-count")
.intType()
diff --git
a/amoro-ams/src/main/java/org/apache/amoro/server/AmoroServiceContainer.java
b/amoro-ams/src/main/java/org/apache/amoro/server/AmoroServiceContainer.java
index e48961d9e..2e72cef20 100644
--- a/amoro-ams/src/main/java/org/apache/amoro/server/AmoroServiceContainer.java
+++ b/amoro-ams/src/main/java/org/apache/amoro/server/AmoroServiceContainer.java
@@ -610,9 +610,7 @@ public class AmoroServiceContainer {
return ConfigHelpers.convertConfigurationKeys(prefix, System.getenv());
}
- /**
- * Configures Iceberg's global worker pool and initializes self-optimizing
Iceberg I/O pools.
- */
+ /** Configures Iceberg's global worker pool and initializes Amoro Iceberg
I/O pools. */
private void initIcebergThreadPools() {
int workerThreadPoolSize =
Math.max(
@@ -630,7 +628,10 @@ public class AmoroServiceContainer {
Math.max(
Runtime.getRuntime().availableProcessors() / 2,
serviceConfig.getInteger(AmoroManagementConf.TABLE_MANIFEST_IO_COMMIT_THREAD_COUNT));
+ int maintenanceThreadPoolSize =
+
serviceConfig.getInteger(AmoroManagementConf.TABLE_MANIFEST_IO_MAINTENANCE_THREAD_COUNT);
IcebergThreadPools.init(planningThreadPoolSize, commitThreadPoolSize);
+ IcebergThreadPools.initMaintenanceThreadPool(maintenanceThreadPoolSize);
}
private void initContainerConfig() {
diff --git
a/amoro-ams/src/test/java/org/apache/amoro/server/TestAmoroManagementConf.java
b/amoro-ams/src/test/java/org/apache/amoro/server/TestAmoroManagementConf.java
index 414f760c2..aff2a1708 100644
---
a/amoro-ams/src/test/java/org/apache/amoro/server/TestAmoroManagementConf.java
+++
b/amoro-ams/src/test/java/org/apache/amoro/server/TestAmoroManagementConf.java
@@ -88,6 +88,14 @@ public class TestAmoroManagementConf {
Duration.ofDays(7),
serviceConfig.get(AmoroManagementConf.PROCESS_HISTORY_DATA_KEEP_TIME));
}
+ @Test
+ void testMaintenanceManifestIoThreadCountDefault() {
+ Configurations serviceConfig = new Configurations();
+ Assertions.assertEquals(
+ 10,
+
serviceConfig.getInteger(AmoroManagementConf.TABLE_MANIFEST_IO_MAINTENANCE_THREAD_COUNT));
+ }
+
@Test
void testDeprecatedIntegerConfigDefaults() {
Configurations serviceConfig = new Configurations();
diff --git
a/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/IcebergTableMaintainer.java
b/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/IcebergTableMaintainer.java
index 68d90e3a1..5333856ce 100644
---
a/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/IcebergTableMaintainer.java
+++
b/amoro-format-iceberg/src/main/java/org/apache/amoro/formats/iceberg/maintainer/IcebergTableMaintainer.java
@@ -41,6 +41,7 @@ import
org.apache.amoro.shade.guava32.com.google.common.collect.Iterables;
import org.apache.amoro.shade.guava32.com.google.common.collect.Maps;
import org.apache.amoro.shade.guava32.com.google.common.collect.Sets;
import org.apache.amoro.table.TableIdentifier;
+import org.apache.amoro.utils.IcebergThreadPools;
import org.apache.amoro.utils.TableFileUtil;
import org.apache.iceberg.ContentFile;
import org.apache.iceberg.ContentScanTask;
@@ -223,6 +224,7 @@ public class IcebergTableMaintainer implements
TableMaintainer {
.retainLast(Math.max(minCount, 1))
.expireOlderThan(olderThan)
.deleteWith(expiredFileCleaner::addFile)
+ .planWith(IcebergThreadPools.getMaintenanceExecutor())
.cleanExpiredFiles(
true) /* enable clean only for collecting the expired files,
will delete them later */;
// iceberg auto-selects IncrementalFileCleanup for single-ref tables. That
strategy walks the
diff --git
a/amoro-format-iceberg/src/main/java/org/apache/amoro/utils/IcebergThreadPools.java
b/amoro-format-iceberg/src/main/java/org/apache/amoro/utils/IcebergThreadPools.java
index 96c0f8fe5..1a93173b3 100644
---
a/amoro-format-iceberg/src/main/java/org/apache/amoro/utils/IcebergThreadPools.java
+++
b/amoro-format-iceberg/src/main/java/org/apache/amoro/utils/IcebergThreadPools.java
@@ -33,6 +33,7 @@ public class IcebergThreadPools {
private static final String PLANNING_POOL_NAME_PREFIX =
"iceberg-planning-pool";
private static final String COMMIT_POOL_NAME_PREFIX = "iceberg-commit-pool";
+ private static final String MAINTENANCE_POOL_NAME_PREFIX =
"iceberg-maintenance-pool";
private static final Map<String, Integer> POOL_SIZES = new
ConcurrentHashMap<>();
private static final Map<String, ExecutorService> POOLS = new
ConcurrentHashMap<>();
@@ -90,6 +91,35 @@ public class IcebergThreadPools {
return getThreadPool(COMMIT_POOL_NAME_PREFIX);
}
+ /**
+ * Initializes the process-wide Iceberg pool for best-effort table
maintenance.
+ *
+ * <p>Repeated initialization is ignored because existing pools cannot be
resized.
+ *
+ * @param poolSize number of worker threads
+ */
+ public static void initMaintenanceThreadPool(int poolSize) {
+ newThreadPool(MAINTENANCE_POOL_NAME_PREFIX, poolSize);
+ }
+
+ /**
+ * Return an {@link ExecutorService} for best-effort Iceberg table
maintenance.
+ *
+ * <p>Snapshot expiration planning is the first consumer of this shared
maintenance pool. Other
+ * maintenance operations can migrate to it as their Iceberg APIs expose
executor hooks.
+ *
+ * <p>The size of this pool is controlled by the AMS configuration {@code
+ * table-manifest-io.maintenance-thread-count}.
+ *
+ * <p>Before the dedicated pool is initialized, this returns Iceberg's
global worker pool.
+ *
+ * @return the maintenance pool, or Iceberg's global worker pool if the
dedicated pool has not
+ * been initialized
+ */
+ public static ExecutorService getMaintenanceExecutor() {
+ return getThreadPool(MAINTENANCE_POOL_NAME_PREFIX);
+ }
+
/**
* Returns the registered Iceberg thread pool for the given name prefix.
*
diff --git
a/amoro-format-iceberg/src/test/java/org/apache/amoro/formats/iceberg/maintainer/TestIcebergTableMaintainer.java
b/amoro-format-iceberg/src/test/java/org/apache/amoro/formats/iceberg/maintainer/TestIcebergTableMaintainer.java
new file mode 100644
index 000000000..da5639af0
--- /dev/null
+++
b/amoro-format-iceberg/src/test/java/org/apache/amoro/formats/iceberg/maintainer/TestIcebergTableMaintainer.java
@@ -0,0 +1,62 @@
+/*
+ * 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.formats.iceberg.maintainer;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+import org.apache.amoro.io.AuthenticatedFileIO;
+import org.apache.amoro.maintainer.TableMaintainerContext;
+import org.apache.amoro.table.TableIdentifier;
+import org.apache.amoro.utils.IcebergThreadPools;
+import org.apache.iceberg.ExpireSnapshots;
+import org.apache.iceberg.Table;
+import org.junit.jupiter.api.Test;
+
+class TestIcebergTableMaintainer {
+
+ @Test
+ void testExpireSnapshotsUsesMaintenancePool() {
+ IcebergThreadPools.initMaintenanceThreadPool(1);
+
+ Table table = mock(Table.class);
+ ExpireSnapshots expireSnapshots = mock(ExpireSnapshots.class);
+ when(table.name()).thenReturn("test_table");
+ when(table.io()).thenReturn(mock(AuthenticatedFileIO.class));
+ when(table.expireSnapshots()).thenReturn(expireSnapshots);
+ when(expireSnapshots.retainLast(1)).thenReturn(expireSnapshots);
+ when(expireSnapshots.expireOlderThan(100L)).thenReturn(expireSnapshots);
+ when(expireSnapshots.deleteWith(any())).thenReturn(expireSnapshots);
+ when(expireSnapshots.planWith(any())).thenReturn(expireSnapshots);
+ when(expireSnapshots.cleanExpiredFiles(true)).thenReturn(expireSnapshots);
+
+ IcebergTableMaintainer tableMaintainer =
+ new IcebergTableMaintainer(
+ table,
+ TableIdentifier.of("test_catalog", "test_database", "test_table"),
+ mock(TableMaintainerContext.class));
+
+ tableMaintainer.expireSnapshots(100L, 1);
+
+
verify(expireSnapshots).planWith(IcebergThreadPools.getMaintenanceExecutor());
+ verify(expireSnapshots).commit();
+ }
+}
diff --git
a/amoro-format-iceberg/src/test/java/org/apache/amoro/utils/TestIcebergThreadPools.java
b/amoro-format-iceberg/src/test/java/org/apache/amoro/utils/TestIcebergThreadPools.java
index a5bfb75c1..f1f3eafbf 100644
---
a/amoro-format-iceberg/src/test/java/org/apache/amoro/utils/TestIcebergThreadPools.java
+++
b/amoro-format-iceberg/src/test/java/org/apache/amoro/utils/TestIcebergThreadPools.java
@@ -41,12 +41,17 @@ public class TestIcebergThreadPools {
Assert.assertSame(registeredPool,
IcebergThreadPools.getThreadPool("registered-test-pool"));
IcebergThreadPools.init(1, 1);
+ IcebergThreadPools.initMaintenanceThreadPool(1);
ExecutorService planningPool = IcebergThreadPools.getPlanningExecutor();
ExecutorService commitPool = IcebergThreadPools.getCommitExecutor();
+ ExecutorService maintenancePool =
IcebergThreadPools.getMaintenanceExecutor();
Assert.assertNotSame(workerPool, planningPool);
Assert.assertNotSame(workerPool, commitPool);
+ Assert.assertNotSame(workerPool, maintenancePool);
Assert.assertNotSame(planningPool, commitPool);
+ Assert.assertNotSame(planningPool, maintenancePool);
+ Assert.assertNotSame(commitPool, maintenancePool);
Assert.assertTrue(
planningPool
.submit(() -> Thread.currentThread().getName())
@@ -57,9 +62,18 @@ public class TestIcebergThreadPools {
.submit(() -> Thread.currentThread().getName())
.get(10, TimeUnit.SECONDS)
.startsWith("iceberg-commit-pool-"));
+ Assert.assertTrue(
+ maintenancePool
+ .submit(() -> Thread.currentThread().getName())
+ .get(10, TimeUnit.SECONDS)
+ .startsWith("iceberg-maintenance-pool-"));
IcebergThreadPools.init(2, 2);
+ IcebergThreadPools.initMaintenanceThreadPool(2);
Assert.assertSame(planningPool, IcebergThreadPools.getPlanningExecutor());
Assert.assertSame(commitPool, IcebergThreadPools.getCommitExecutor());
+ Assert.assertSame(maintenancePool,
IcebergThreadPools.getMaintenanceExecutor());
+ Assert.assertThrows(
+ IllegalArgumentException.class, () ->
IcebergThreadPools.initMaintenanceThreadPool(0));
}
}
diff --git a/charts/amoro/templates/amoro-configmap.yaml
b/charts/amoro/templates/amoro-configmap.yaml
index 7a9f14ba5..26e690f0f 100644
--- a/charts/amoro/templates/amoro-configmap.yaml
+++ b/charts/amoro/templates/amoro-configmap.yaml
@@ -101,6 +101,7 @@ data:
# optional features
table-manifest-io:
thread-count: 20
+ maintenance-thread-count: 10
database:
type: {{ .Values.amoroConf.database.type }}
diff --git a/dist/src/main/amoro-bin/conf/config.yaml
b/dist/src/main/amoro-bin/conf/config.yaml
index f4cdaaac8..d4e60010a 100644
--- a/dist/src/main/amoro-bin/conf/config.yaml
+++ b/dist/src/main/amoro-bin/conf/config.yaml
@@ -110,6 +110,7 @@ ams:
# optional features
table-manifest-io:
thread-count: 20
+ maintenance-thread-count: 10
catalog-meta-cache:
expiration-interval: 60s
diff --git a/docs/configuration/ams-config.md b/docs/configuration/ams-config.md
index 2a6c6b6e8..e2746d79b 100644
--- a/docs/configuration/ams-config.md
+++ b/docs/configuration/ams-config.md
@@ -116,6 +116,7 @@ table td:last-child, table th:last-child { width: 40%;
word-break: break-all; }
| self-optimizing.runtime-data-keep-time | 30 d | Duration that
self-optimizing runtime data is retained. |
| server-bind-host | 0.0.0.0 | The host bound to the server. |
| server-expose-host | | The exposed host of the server. |
+| table-manifest-io.maintenance-thread-count | 10 | Sets the size of the
worker pool used for manifest I/O across best-effort table maintenance
operations. |
| table-manifest-io.thread-count | 20 | Sets the size of the worker pool. The
worker pool limits the number of tasks concurrently processing manifests in the
base table implementation across all concurrent planning or commit operations. |
| terminal.backend | local | Terminal backend implementation. local, kyuubi
and custom are valid values. |
| terminal.factory | <undefined> | Session factory implement of
terminal, `terminal.backend` must be `custom` if this is set. |