This is an automated email from the ASF dual-hosted git repository.
leonbao pushed a commit to branch 1.3.6-prepare
in repository https://gitbox.apache.org/repos/asf/incubator-dolphinscheduler.git
The following commit(s) were added to refs/heads/1.3.6-prepare by this push:
new e5e111d [1.3.6-prepare][Fix-4862][Server] Kill yarn application
command won't be execute whe… #4936 (#5072)
e5e111d is described below
commit e5e111dcf299b92f220b8e0096f3e1437c74e70d
Author: Kirs <[email protected]>
AuthorDate: Thu Mar 18 14:08:56 2021 +0800
[1.3.6-prepare][Fix-4862][Server] Kill yarn application command won't be
execute whe… #4936 (#5072)
pr #4936
issue #4862
---
.../server/utils/ProcessUtils.java | 5 ++--
.../server/worker/processor/TaskKillProcessor.java | 32 +++++++++++-----------
2 files changed, 18 insertions(+), 19 deletions(-)
diff --git
a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/utils/ProcessUtils.java
b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/utils/ProcessUtils.java
index 270f43e..f09e317 100644
---
a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/utils/ProcessUtils.java
+++
b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/utils/ProcessUtils.java
@@ -374,12 +374,11 @@ public class ProcessUtils {
OSUtils.exeCmd(cmd);
- // find log and kill yarn job
- killYarnJob(taskExecutionContext);
-
} catch (Exception e) {
logger.error("kill task failed", e);
}
+ // find log and kill yarn job
+ killYarnJob(taskExecutionContext);
}
/**
diff --git
a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillProcessor.java
b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillProcessor.java
index c9f2cd4..c5cfc01 100644
---
a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillProcessor.java
+++
b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillProcessor.java
@@ -102,11 +102,11 @@ public class TaskKillProcessor implements
NettyRequestProcessor {
* @return kill result
*/
private Pair<Boolean, List<String>> doKill(TaskKillRequestCommand
killCommand){
+ boolean processFlag = true;
List<String> appIds = Collections.emptyList();
+ int taskInstanceId = killCommand.getTaskInstanceId();
+ TaskExecutionContext taskExecutionContext =
taskExecutionContextCacheManager.getByTaskInstanceId(taskInstanceId);
try {
- int taskInstanceId = killCommand.getTaskInstanceId();
- TaskExecutionContext taskExecutionContext =
taskExecutionContextCacheManager.getByTaskInstanceId(taskInstanceId);
-
Integer processId = taskExecutionContext.getProcessId();
if (processId.equals(0)) {
@@ -121,17 +121,16 @@ public class TaskKillProcessor implements
NettyRequestProcessor {
OSUtils.exeCmd(cmd);
- // find log and kill yarn job
- appIds =
killYarnJob(Host.of(taskExecutionContext.getHost()).getIp(),
- taskExecutionContext.getLogPath(),
- taskExecutionContext.getExecutePath(),
- taskExecutionContext.getTenantCode());
-
- return Pair.of(true, appIds);
} catch (Exception e) {
+ processFlag = false;
logger.error("kill task error", e);
}
- return Pair.of(false, appIds);
+ // find log and kill yarn job
+ Pair<Boolean, List<String>> yarnResult =
killYarnJob(Host.of(taskExecutionContext.getHost()).getIp(),
+ taskExecutionContext.getLogPath(),
+ taskExecutionContext.getExecutePath(),
+ taskExecutionContext.getTenantCode());
+ return Pair.of(processFlag && yarnResult.getLeft(),
yarnResult.getRight());
}
/**
@@ -162,24 +161,25 @@ public class TaskKillProcessor implements
NettyRequestProcessor {
* @param logPath logPath
* @param executePath executePath
* @param tenantCode tenantCode
- * @return List<String> appIds
+ * @return Pair<Boolean, List<String>> yarn kill result
*/
- private List<String> killYarnJob(String host, String logPath, String
executePath, String tenantCode) {
+ private Pair<Boolean, List<String>> killYarnJob(String host, String
logPath, String executePath, String tenantCode) {
LogClientService logClient = null;
try {
logClient = new LogClientService();
logger.info("view log host : {},logPath : {}", host,logPath);
String log = logClient.viewLog(host, Constants.RPC_PORT, logPath);
+ List<String> appIds = Collections.emptyList();
if (StringUtils.isNotEmpty(log)) {
- List<String> appIds = LoggerUtils.getAppIds(log, logger);
+ appIds = LoggerUtils.getAppIds(log, logger);
if (StringUtils.isEmpty(executePath)) {
logger.error("task instance execute path is empty");
throw new RuntimeException("task instance execute path is
empty");
}
if (appIds.size() > 0) {
ProcessUtils.cancelApplication(appIds, logger, tenantCode,
executePath);
- return appIds;
+ return Pair.of(true, appIds);
}
}
} catch (Exception e) {
@@ -189,7 +189,7 @@ public class TaskKillProcessor implements
NettyRequestProcessor {
logClient.close();
}
}
- return Collections.EMPTY_LIST;
+ return Pair.of(false, Collections.emptyList());
}
}