提交 e2486188 编写于 作者: M Maximilian Michels

[hotfix] respect default local number of task managers

上级 e20c6390
......@@ -1074,7 +1074,7 @@ public class ExecutionGraph implements Serializable {
LOG.warn("Received accumulator result for unknown execution {}.", execID);
}
} catch (Exception e) {
LOG.error("Cannot update accumulators for job " + jobID, e);
LOG.error("Cannot update accumulators for job {}.", jobID, e);
}
}
......
......@@ -158,7 +158,8 @@ class LocalFlinkMiniCluster(
ConfigConstants.DEFAULT_TASK_MANAGER_NETWORK_NUM_BUFFERS) * bufferSize.toLong
val numTaskManager = config.getInteger(
ConfigConstants.LOCAL_NUMBER_TASK_MANAGER, 1)
ConfigConstants.LOCAL_NUMBER_TASK_MANAGER,
ConfigConstants.DEFAULT_LOCAL_NUMBER_TASK_MANAGER)
val memoryFraction = config.getFloat(
ConfigConstants.TASK_MANAGER_MEMORY_FRACTION_KEY,
......@@ -184,7 +185,8 @@ class LocalFlinkMiniCluster(
config.setString(ConfigConstants.JOB_MANAGER_IPC_ADDRESS_KEY, hostname)
config.setInteger(ConfigConstants.LOCAL_NUMBER_TASK_MANAGER, 1)
config.setInteger(ConfigConstants.LOCAL_NUMBER_TASK_MANAGER,
ConfigConstants.DEFAULT_LOCAL_NUMBER_TASK_MANAGER)
// Reduce number of threads for local execution
config.setInteger(NettyConfig.NUM_THREADS_CLIENT, 1)
......
Markdown is supported
0% .
You are about to add 0 people to the discussion. Proceed with caution.
先完成此消息的编辑!
想要评论请 注册