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());
     }
 
 }

Reply via email to