From 56419924e31277c6ede51ae6f46bfdf16279379f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BC=A0=E4=B8=87=E4=B9=89?= Date: Wed, 19 Aug 2026 07:31:37 +0800 Subject: [PATCH] [Fix] Fix memory sizes config not taking effect in Yarn mode --- .../flink/client/trait/FlinkClientTrait.java | 16 ++++ .../flink/client/test/SubmitRequestTest.java | 78 +++++++++++++++++++ 2 files changed, 94 insertions(+) diff --git a/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/java/org/apache/streampark/flink/client/trait/FlinkClientTrait.java b/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/java/org/apache/streampark/flink/client/trait/FlinkClientTrait.java index ca5f66a3ce..061977e8f5 100644 --- a/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/java/org/apache/streampark/flink/client/trait/FlinkClientTrait.java +++ b/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/java/org/apache/streampark/flink/client/trait/FlinkClientTrait.java @@ -329,6 +329,7 @@ private Configuration prepareConfig(SubmitRequest submitRequest) throws Exceptio applyCheckpointDefaults(submitRequest, flinkConfig); applySavepointConfig(submitRequest, flinkConfig); applyEnvProperties(submitRequest, flinkConfig); + applyAppProperties(submitRequest, flinkConfig); return flinkConfig; } @@ -455,6 +456,21 @@ private void applyEnvProperties(SubmitRequest submitRequest, Configuration flink } } + private void applyAppProperties(SubmitRequest submitRequest, Configuration flinkConfig) { + Map appProperties = submitRequest.appProperties(); + if (MapUtils.isEmpty(appProperties)) { + return; + } + for (Map.Entry entry : appProperties.entrySet()) { + logInfo( + "appProperties: " + + entry.getKey() + + " : " + + entry.getValue()); + flinkConfig.setString(entry.getKey(), entry.getValue()); + } + } + public abstract void setConfig(SubmitRequest submitRequest, Configuration flinkConf); public SavepointResponse triggerSavepoint(TriggerSavepointRequest savepointRequest) throws FlinkException { diff --git a/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/test/java/org/apache/streampark/flink/client/test/SubmitRequestTest.java b/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/test/java/org/apache/streampark/flink/client/test/SubmitRequestTest.java index 441b213f04..1a2bc0ca9d 100644 --- a/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/test/java/org/apache/streampark/flink/client/test/SubmitRequestTest.java +++ b/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/test/java/org/apache/streampark/flink/client/test/SubmitRequestTest.java @@ -23,6 +23,7 @@ import org.apache.streampark.common.enums.ApplicationType; import org.apache.streampark.common.enums.FlinkDeployMode; import org.apache.streampark.common.enums.FlinkJobType; +import org.apache.streampark.common.util.DeflaterUtils; import org.apache.streampark.flink.client.bean.SubmitApplicationSpec; import org.apache.streampark.flink.client.bean.SubmitRequest; @@ -134,6 +135,83 @@ void allowNonRestoredStateShouldDefaultToFalse() { assertThat(request.allowNonRestoredState()).isFalse(); } + @Test + void appPropertiesWithMemoryConfigFromYamlConf() { + String yamlContent = + "flink.property.jobmanager.memory.process.size: 2048m\n" + + "flink.property.taskmanager.memory.process.size: 4096m\n"; + String appConf = "yaml://" + DeflaterUtils.zipString(yamlContent); + SubmitApplicationSpec application = + SubmitApplicationSpec.builder() + .jobType(FlinkJobType.FLINK_JAR) + .appConf(appConf) + .build(); + SubmitRequest request = + new SubmitRequest( + FLINK_VERSION, + FlinkDeployMode.YARN_APPLICATION, + Collections.emptyMap(), + application, + null, + null, + null); + + assertThat(request.appProperties()) + .containsEntry("jobmanager.memory.process.size", "2048m") + .containsEntry("taskmanager.memory.process.size", "4096m"); + } + + @Test + void appPropertiesWithMemoryConfigFromJsonConf() { + String jmKey = ConfigKeys.KEY_FLINK_PROPERTY_PREFIX() + "jobmanager.memory.process.size"; + String tmKey = ConfigKeys.KEY_FLINK_PROPERTY_PREFIX() + "taskmanager.memory.process.size"; + String appConf = "json://{\"" + jmKey + "\":\"2g\",\"" + tmKey + "\":\"4g\"}"; + SubmitApplicationSpec application = + SubmitApplicationSpec.builder() + .jobType(FlinkJobType.FLINK_JAR) + .appConf(appConf) + .build(); + SubmitRequest request = + new SubmitRequest( + FLINK_VERSION, + FlinkDeployMode.YARN_APPLICATION, + Collections.emptyMap(), + application, + null, + null, + null); + + assertThat(request.appProperties()) + .containsEntry("jobmanager.memory.process.size", "2g") + .containsEntry("taskmanager.memory.process.size", "4g"); + } + + @Test + void appPropertiesWithMemoryConfigFromPropertiesConf() { + String propertiesContent = + "flink.property.jobmanager.memory.process.size=1024m\n" + + "flink.property.taskmanager.memory.process.size=2048m\n"; + String appConf = "prop://" + DeflaterUtils.zipString(propertiesContent); + SubmitApplicationSpec application = + SubmitApplicationSpec.builder() + .jobType(FlinkJobType.FLINK_JAR) + .appConf(appConf) + .build(); + SubmitRequest request = + new SubmitRequest( + FLINK_VERSION, + FlinkDeployMode.YARN_APPLICATION, + Collections.emptyMap(), + application, + null, + null, + null); + + assertThat(request.appProperties()) + .containsEntry("jobmanager.memory.process.size", "1024m") + .containsEntry("taskmanager.memory.process.size", "2048m"); + } + private static SubmitRequest createRequest(FlinkJobType jobType, String appConf) { SubmitApplicationSpec application = SubmitApplicationSpec.builder()