github-advanced-security[bot] commented on code in PR #18667:
URL:
https://github.com/apache/dolphinscheduler/pull/18667#discussion_r4171236632
##########
dolphinscheduler-task-plugin/dolphinscheduler-task-flink/src/main/java/org/apache/dolphinscheduler/plugin/task/flink/FlinkArgsUtils.java:
##########
@@ -63,29 +83,122 @@
/**
* build flink cancel command line
- * @param taskExecutionContext
- * @return
+ *
+ * @param jobId the Flink JobID printed by `flink run`, it is not the
YARN/K8s application id
+ * @return argument list
*/
- public static List<String> buildCancelCommandLine(TaskExecutionContext
taskExecutionContext) {
+ public static List<String> buildCancelCommandLine(String jobId) {
List<String> args = new ArrayList<>();
args.add(FlinkConstants.FLINK_COMMAND);
args.add(FlinkConstants.FLINK_CANCEL);
- args.add(taskExecutionContext.getAppIds());
+ args.add(jobId);
return args;
}
/**
* build flink savepoint command line, the savepoint folder should be set
in flink conf
- * @return
+ *
+ * @param jobId the Flink JobID printed by `flink run`, it is not the
YARN/K8s application id
+ * @return argument list
*/
- public static List<String> buildSavePointCommandLine(TaskExecutionContext
taskExecutionContext) {
+ public static List<String> buildSavePointCommandLine(String jobId) {
List<String> args = new ArrayList<>();
args.add(FlinkConstants.FLINK_COMMAND);
args.add(FlinkConstants.FLINK_SAVEPOINT);
- args.add(taskExecutionContext.getAppIds());
+ args.add(jobId);
return args;
}
+ /**
+ * Execute a flink command in the same environment as the task itself.
+ *
+ * <p>The task script is executed by a shell interceptor which sources
shell.env_source_list and the
+ * task's custom environment, resolves the task parameters and runs as the
task's tenant. The
+ * cancel / savepoint commands must reuse that environment, otherwise
placeholders like
+ * ${FLINK_HOME} cannot be resolved when they are only defined in the
selected task environment.
+ *
+ * @param taskExecutionContext task execution context
+ * @param args the command arguments
+ * @return true if the command finished successfully
+ */
+ public static boolean executeCommand(TaskExecutionContext
taskExecutionContext, List<String> args) {
+ return executeCommand(taskExecutionContext, args,
OSUtils.isSudoEnable());
+ }
+
+ /**
+ * Same as {@link #executeCommand(TaskExecutionContext, List)}, with the
sudo mode injected so that
+ * the command execution can be verified without requiring sudo
permissions.
+ */
+ static boolean executeCommand(TaskExecutionContext taskExecutionContext,
List<String> args, boolean sudoEnable) {
+ Process process = null;
+ try {
+ IShellInterceptorBuilder shellInterceptorBuilder =
ShellInterceptorBuilderFactory.newBuilder()
+ .shellDirectory(taskExecutionContext.getExecutePath())
+ .shellName(taskExecutionContext.getTaskAppId() +
FLINK_COMMAND_SHELL_NAME_SUFFIX
Review Comment:
## CodeQL / Deprecated method or constructor invocation
Invoking [TaskExecutionContext.getTaskAppId](1) should be avoided because it
has been deprecated.
[Show more
details](https://github.com/apache/dolphinscheduler/security/code-scanning/5865)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]