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 4e09b201c47c35f1b2c57408d37ec7642229c0fa Author: csurong <[email protected]> AuthorDate: Fri Jul 24 11:49:22 2026 +0800 [AMORO-4126][ams] Default waitAppCompletion=false for Spark optimizer on Kubernetes (#4289) [AMORO-4126][ams] Default waitAppCompletion to false on Kubernetes Co-authored-by: caosurong <[email protected]> (cherry picked from commit eee2e3573180ae943adfa867018bdf310027efdf) --- .../server/manager/SparkOptimizerContainer.java | 8 +++ .../manager/TestSparkOptimizerContainer.java | 73 ++++++++++++++++++++++ 2 files changed, 81 insertions(+) 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 8609cf352..d453798fe 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 @@ -202,6 +202,12 @@ public class SparkOptimizerContainer extends AbstractOptimizerContainer { private void addKubernetesProperties( Resource resource, SparkOptimizerContainer.SparkConf sparkConf) { + sparkConf.putToOptions( + SparkConfKeys.KUBERNETES_SUBMISSION_WAIT_APP_COMPLETION, + StringUtils.defaultIfEmpty( + sparkConf.configValue(SparkConfKeys.KUBERNETES_SUBMISSION_WAIT_APP_COMPLETION), + "false")); + String driverName = kubernetesDriverName(resource); sparkConf.putToOptions( SparkOptimizerContainer.SparkConfKeys.KUBERNETES_DRIVER_NAME, driverName); @@ -332,6 +338,8 @@ public class SparkOptimizerContainer extends AbstractOptimizerContainer { public static final String KUBERNETES_IMAGE_REF = "spark.kubernetes.container.image"; public static final String KUBERNETES_DRIVER_NAME = "spark.kubernetes.driver.pod.name"; 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_DRIVER_LABEL_PREFIX = "spark.kubernetes.driver.label."; public static final String KUBERNETES_DRA_ENABLED = "spark.dynamicAllocation.enabled"; 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 2936c83eb..5a432feb8 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 @@ -19,6 +19,8 @@ package org.apache.amoro.server.manager; import org.apache.amoro.OptimizerProperties; +import org.apache.amoro.resource.Resource; +import org.apache.amoro.resource.ResourceType; import org.apache.amoro.shade.guava32.com.google.common.collect.Maps; import org.junit.Assert; import org.junit.Test; @@ -58,4 +60,75 @@ public class TestSparkOptimizerContainer { Assert.assertTrue(sparkOptions.contains("--conf key2=value4")); Assert.assertTrue(sparkOptions.contains("--conf key5=value5")); } + + @Test + public void testKubernetesWaitAppCompletionDefaultsToFalse() { + SparkOptimizerContainer container = createContainer("k8s://https://127.0.0.1:6443"); + + String startupArgs = + container.buildOptimizerStartupArgsString(createResource(Maps.newHashMap())); + + Assert.assertTrue( + startupArgs.contains( + "--conf " + + SparkOptimizerContainer.SparkConfKeys.KUBERNETES_SUBMISSION_WAIT_APP_COMPLETION + + "=false")); + } + + @Test + public void testKubernetesWaitAppCompletionCanBeOverridden() { + SparkOptimizerContainer container = createContainer("k8s://https://127.0.0.1:6443"); + Map<String, String> resourceProperties = Maps.newHashMap(); + resourceProperties.put( + SparkOptimizerContainer.SparkConf.SPARK_PARAMETER_PREFIX + + SparkOptimizerContainer.SparkConfKeys.KUBERNETES_SUBMISSION_WAIT_APP_COMPLETION, + "true"); + + String startupArgs = + container.buildOptimizerStartupArgsString(createResource(resourceProperties)); + + Assert.assertTrue( + startupArgs.contains( + "--conf " + + SparkOptimizerContainer.SparkConfKeys.KUBERNETES_SUBMISSION_WAIT_APP_COMPLETION + + "=true")); + Assert.assertFalse( + startupArgs.contains( + "--conf " + + SparkOptimizerContainer.SparkConfKeys.KUBERNETES_SUBMISSION_WAIT_APP_COMPLETION + + "=false")); + } + + @Test + public void testYarnDoesNotSetKubernetesWaitAppCompletion() { + SparkOptimizerContainer container = createContainer("yarn"); + + String startupArgs = + container.buildOptimizerStartupArgsString(createResource(Maps.newHashMap())); + + Assert.assertFalse( + startupArgs.contains( + SparkOptimizerContainer.SparkConfKeys.KUBERNETES_SUBMISSION_WAIT_APP_COMPLETION)); + } + + private SparkOptimizerContainer createContainer(String master) { + Map<String, String> properties = Maps.newHashMap(containerProperties); + properties.put(SparkOptimizerContainer.SPARK_MASTER, master); + if (master.startsWith("k8s://")) { + properties.put( + SparkOptimizerContainer.SparkConf.SPARK_PARAMETER_PREFIX + + SparkOptimizerContainer.SparkConfKeys.KUBERNETES_IMAGE_REF, + "apache/amoro-spark-optimizer:test"); + } + SparkOptimizerContainer container = new SparkOptimizerContainer(); + container.init("spark", properties); + return container; + } + + private Resource createResource(Map<String, String> properties) { + return new Resource.Builder("spark", "test-group", ResourceType.OPTIMIZER) + .setThreadCount(1) + .setProperties(properties) + .build(); + } }
