fix the #14729 problem (#14902)

Co-authored-by: 刘阳 <liuy2590@chinaunicom.cn>
Co-authored-by: xiangzihao <460888207@qq.com>
(cherry picked from commit e748c2eb9a)
This commit is contained in:
LiuCanWu 2023-09-21 18:46:14 +08:00 committed by Jay Chung
parent 23f6494ad6
commit 06a555e494
1 changed files with 6 additions and 2 deletions

View File

@ -273,6 +273,8 @@ public class FlinkArgsUtils {
args.add(others);
}
// determine yarn queue
determinedYarnQueue(args, flinkParameters, deployMode, flinkVersion);
ProgramType programType = flinkParameters.getProgramType();
String mainClass = flinkParameters.getMainClass();
if (programType != null && programType != ProgramType.PYTHON && StringUtils.isNotEmpty(mainClass)) {
@ -295,8 +297,6 @@ public class FlinkArgsUtils {
args.add(ParameterUtils.convertParameterPlaceholders(mainArgs, ParameterUtils.convert(paramsMap)));
}
// determine yarn queue
determinedYarnQueue(args, flinkParameters, deployMode, flinkVersion);
return args;
}
@ -310,8 +310,10 @@ public class FlinkArgsUtils {
} else {
doAddQueue(args, flinkParameters, FlinkConstants.FLINK_YARN_QUEUE_FOR_MODE);
}
break;
case APPLICATION:
doAddQueue(args, flinkParameters, FlinkConstants.FLINK_YARN_QUEUE_FOR_TARGETS);
break;
}
}
@ -323,9 +325,11 @@ public class FlinkArgsUtils {
switch (option) {
case FlinkConstants.FLINK_YARN_QUEUE_FOR_TARGETS:
args.add(String.format(FlinkConstants.FLINK_YARN_QUEUE_FOR_TARGETS + "=%s", yarnQueue));
break;
case FlinkConstants.FLINK_YARN_QUEUE_FOR_MODE:
args.add(FlinkConstants.FLINK_YARN_QUEUE_FOR_MODE);
args.add(yarnQueue);
break;
}
}
}