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 eee2e3573 [AMORO-4126][ams] Default waitAppCompletion=false for Spark 
optimizer on Kubernetes (#4289)
eee2e3573 is described below

commit eee2e3573180ae943adfa867018bdf310027efdf
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]>
---
 .../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();
+  }
 }

Reply via email to