diff --git a/dolphinscheduler-task-plugin/dolphinscheduler-task-procedure/src/main/java/org/apache/dolphinscheduler/plugin/task/procedure/ProcedureTask.java b/dolphinscheduler-task-plugin/dolphinscheduler-task-procedure/src/main/java/org/apache/dolphinscheduler/plugin/task/procedure/ProcedureTask.java index d0b42aeecdfe..bca2d4bbf70f 100644 --- a/dolphinscheduler-task-plugin/dolphinscheduler-task-procedure/src/main/java/org/apache/dolphinscheduler/plugin/task/procedure/ProcedureTask.java +++ b/dolphinscheduler-task-plugin/dolphinscheduler-task-procedure/src/main/java/org/apache/dolphinscheduler/plugin/task/procedure/ProcedureTask.java @@ -43,6 +43,7 @@ import java.sql.CallableStatement; import java.sql.Connection; import java.sql.SQLException; +import java.sql.Statement; import java.sql.Types; import java.util.HashMap; import java.util.Map; @@ -60,6 +61,8 @@ public class ProcedureTask extends AbstractTask { private final ProcedureTaskExecutionContext procedureTaskExecutionContext; + private volatile Statement sessionStatement; + /** * constructor * @@ -105,30 +108,48 @@ public void handle(TaskCallBack taskCallBack) throws TaskException { } String proceduerSql = formatSql(sqlParamsMap, paramsMap); // call method - try (CallableStatement stmt = connection.prepareCall(proceduerSql)) { + try (CallableStatement tmpStatement = connection.prepareCall(proceduerSql)) { + sessionStatement = tmpStatement; // set timeout - setTimeout(stmt); + setTimeout(tmpStatement); // outParameterMap - Map outParameterMap = getOutParameterMap(stmt, sqlParamsMap, paramsMap); + Map outParameterMap = getOutParameterMap(tmpStatement, sqlParamsMap, paramsMap); - stmt.executeUpdate(); + tmpStatement.executeUpdate(); // print the output parameters to the log - printOutParameter(stmt, outParameterMap); + printOutParameter(tmpStatement, outParameterMap); setExitStatusCode(EXIT_CODE_SUCCESS); } } catch (Exception e) { + if (exitStatusCode == TaskConstants.EXIT_CODE_KILL) { + log.info("This procedure task has been killed"); + return; + } setExitStatusCode(EXIT_CODE_FAILURE); - log.error("procedure task error", e); + log.error("Failed to execute this procedure task", e); throw new TaskException("Execute procedure task failed", e); } } @Override public void cancel() throws TaskException { - + if (sessionStatement != null) { + try { + log.info("Try to cancel this procedure task"); + sessionStatement.cancel(); + setExitStatusCode(TaskConstants.EXIT_CODE_KILL); + log.info("This procedure task was canceled"); + } catch (Exception ex) { + log.warn("Failed to cancel this procedure task", ex); + throw new TaskException("Cancel this procedure task failed", ex); + } + } else { + log.info( + "Attempted to cancel this procedure task, but no active statement exists. Possible reasons: task not started, already completed, or canceled."); + } } private String formatSql(Map sqlParamsMap, Map paramsMap) {