Skip to content

Commit 8e35436

Browse files
author
whiletrue
authored
Merge pull request #330 from liubin2048/1.10_release
local模式taskmanager.numberOfTaskSlots指定失效问题修复
2 parents 5c6a062 + b2c9d7c commit 8e35436

File tree

1 file changed

+1
-1
lines changed

1 file changed

+1
-1
lines changed

core/src/main/java/com/dtstack/flink/sql/environment/MyLocalStreamEnvironment.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -105,13 +105,13 @@ public JobExecutionResult execute(StreamGraph streamGraph) throws Exception {
105105
configuration.addAll(jobGraph.getJobConfiguration());
106106

107107
configuration.setString(TaskManagerOptions.MANAGED_MEMORY_SIZE.key(), "512M");
108-
configuration.setInteger(TaskManagerOptions.NUM_TASK_SLOTS.key(), jobGraph.getMaximumParallelism());
109108

110109
// add (and override) the settings with what the user defined
111110
configuration.addAll(this.conf);
112111

113112
MiniClusterConfiguration.Builder configBuilder = new MiniClusterConfiguration.Builder();
114113
configBuilder.setConfiguration(configuration);
114+
configBuilder.setNumSlotsPerTaskManager(jobGraph.getMaximumParallelism());
115115

116116
if (LOG.isInfoEnabled()) {
117117
LOG.info("Running job on local embedded Flink mini cluster");

0 commit comments

Comments
 (0)