This is an automated email from the ASF dual-hosted git repository. xxubai pushed a commit to branch 0.9.x in repository https://gitbox.apache.org/repos/asf/amoro.git
commit f46f57d9423538520d7cafb05a398fa33acd5dbc Author: csurong <[email protected]> AuthorDate: Thu Jul 30 21:39:30 2026 +0800 [AMORO-4301][ams] Fix Kubernetes executor label prefix (#4302) Use Spark's executor label prefix for optimizer executor pods and add regression coverage for both driver and executor labels. Co-authored-by: caosurong <[email protected]> Co-authored-by: ZhouJinsong <[email protected]> (cherry picked from commit efbedc396fa178abe39ec24071f91f24163d8c9a) --- .../server/manager/SparkOptimizerContainer.java | 3 ++- .../manager/TestSparkOptimizerContainer.java | 22 ++++++++++++++++++++++ 2 files changed, 24 insertions(+), 1 deletion(-) diff --git a/amoro-ams/src/main/java/org/apache/amoro/server/manager/SparkOptimizerContainer.java b/amoro-ams/src/main/java/org/apache/amoro/server/manager/SparkOptimizerContainer.java index d453798fe..775ecfdbd 100644 --- a/amoro-ams/src/main/java/org/apache/amoro/server/manager/SparkOptimizerContainer.java +++ b/amoro-ams/src/main/java/org/apache/amoro/server/manager/SparkOptimizerContainer.java @@ -340,7 +340,8 @@ public class SparkOptimizerContainer extends AbstractOptimizerContainer { public static final String KUBERNETES_NAMESPACE = "spark.kubernetes.namespace"; public static final String KUBERNETES_SUBMISSION_WAIT_APP_COMPLETION = "spark.kubernetes.submission.waitAppCompletion"; - public static final String KUBERNETES_EXECUTOR_LABEL_PREFIX = "spark.kubernetes.driver.label."; + public static final String KUBERNETES_EXECUTOR_LABEL_PREFIX = + "spark.kubernetes.executor.label."; public static final String KUBERNETES_DRIVER_LABEL_PREFIX = "spark.kubernetes.driver.label."; public static final String KUBERNETES_DRA_ENABLED = "spark.dynamicAllocation.enabled"; public static final String KUBERNETES_DRA_MAX_EXECUTORS = diff --git a/amoro-ams/src/test/java/org/apache/amoro/server/manager/TestSparkOptimizerContainer.java b/amoro-ams/src/test/java/org/apache/amoro/server/manager/TestSparkOptimizerContainer.java index 5a432feb8..c109848aa 100644 --- a/amoro-ams/src/test/java/org/apache/amoro/server/manager/TestSparkOptimizerContainer.java +++ b/amoro-ams/src/test/java/org/apache/amoro/server/manager/TestSparkOptimizerContainer.java @@ -75,6 +75,17 @@ public class TestSparkOptimizerContainer { + "=false")); } + @Test + public void testKubernetesAddsLabelsToDriverAndExecutorPods() { + SparkOptimizerContainer container = createContainer("k8s://https://127.0.0.1:6443"); + Resource resource = createResource(Maps.newHashMap()); + + String startupArgs = container.buildOptimizerStartupArgsString(resource); + + assertPodLabels(startupArgs, "spark.kubernetes.driver.label.", resource); + assertPodLabels(startupArgs, "spark.kubernetes.executor.label.", resource); + } + @Test public void testKubernetesWaitAppCompletionCanBeOverridden() { SparkOptimizerContainer container = createContainer("k8s://https://127.0.0.1:6443"); @@ -131,4 +142,15 @@ public class TestSparkOptimizerContainer { .setProperties(properties) .build(); } + + private void assertPodLabels(String startupArgs, String labelPrefix, Resource resource) { + Assert.assertTrue( + startupArgs.contains( + "--conf " + labelPrefix + "optimizer-group=" + resource.getGroupName())); + Assert.assertTrue( + startupArgs.contains( + "--conf " + labelPrefix + "optimizer-implementation=spark-native-kubernetes")); + Assert.assertTrue( + startupArgs.contains("--conf " + labelPrefix + "optimizer-id=" + resource.getResourceId())); + } }
