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();
+ }
}