diff --git a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/enums/TaskType.java b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/enums/TaskType.java index a6c539d72d..bbf00f6895 100644 --- a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/enums/TaskType.java +++ b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/enums/TaskType.java @@ -37,7 +37,7 @@ public enum TaskType { * 10 DATAX * 11 CONDITIONS * 12 SQOOP - * 13 WATERDROP + * 13 SEATUNNEL * 15 PIGEON */ SHELL(0, "SHELL"), @@ -53,7 +53,7 @@ public enum TaskType { DATAX(10, "DATAX"), CONDITIONS(11, "CONDITIONS"), SQOOP(12, "SQOOP"), - WATERDROP(13, "WATERDROP"), + SEATUNNEL(13, "SEATUNNEL"), SWITCH(14, "SWITCH"), PIGEON(15, "PIGEON"); diff --git a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/TaskParametersUtils.java b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/TaskParametersUtils.java index 781b83b828..792f6a5577 100644 --- a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/TaskParametersUtils.java +++ b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/TaskParametersUtils.java @@ -60,7 +60,7 @@ public class TaskParametersUtils { case "SUB_PROCESS": return JSONUtils.parseObject(parameter, SubProcessParameters.class); case "SHELL": - case "WATERDROP": + case "SEATUNNEL": return JSONUtils.parseObject(parameter, ShellParameters.class); case "PROCEDURE": return JSONUtils.parseObject(parameter, ProcedureParameters.class); diff --git a/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/pom.xml b/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/pom.xml new file mode 100644 index 0000000000..0bc7c46bed --- /dev/null +++ b/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/pom.xml @@ -0,0 +1,48 @@ + + + + + dolphinscheduler-task-plugin + org.apache.dolphinscheduler + 2.0.0-SNAPSHOT + + 4.0.0 + + dolphinscheduler-task-seatunnel + jar + + + + org.apache.dolphinscheduler + dolphinscheduler-spi + provided + + + org.apache.dolphinscheduler + dolphinscheduler-task-api + ${project.version} + + + + org.apache.commons + commons-collections4 + + + diff --git a/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelParameters.java b/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelParameters.java new file mode 100644 index 0000000000..5c02efee2c --- /dev/null +++ b/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelParameters.java @@ -0,0 +1,62 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.dolphinscheduler.plugin.task.seatunnel; + +import org.apache.dolphinscheduler.spi.task.AbstractParameters; +import org.apache.dolphinscheduler.spi.task.ResourceInfo; + +import java.util.List; + +public class SeatunnelParameters extends AbstractParameters { + + /** + * shell script + */ + private String rawScript; + + /** + * resource list + */ + private List resourceList; + + public String getRawScript() { + return rawScript; + } + + public void setRawScript(String rawScript) { + this.rawScript = rawScript; + } + + public List getResourceList() { + return resourceList; + } + + public void setResourceList(List resourceList) { + this.resourceList = resourceList; + } + + @Override + public boolean checkParameters() { + return rawScript != null && !rawScript.isEmpty(); + } + + @Override + public List getResourceFilesList() { + return resourceList; + } +} diff --git a/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelTask.java b/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelTask.java new file mode 100644 index 0000000000..58dc4bc4b6 --- /dev/null +++ b/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelTask.java @@ -0,0 +1,170 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.dolphinscheduler.plugin.task.seatunnel; + +import static org.apache.dolphinscheduler.spi.task.TaskConstants.EXIT_CODE_FAILURE; +import static org.apache.dolphinscheduler.spi.task.TaskConstants.RWXR_XR_X; + +import org.apache.dolphinscheduler.plugin.task.api.AbstractTaskExecutor; +import org.apache.dolphinscheduler.plugin.task.api.ShellCommandExecutor; +import org.apache.dolphinscheduler.plugin.task.api.TaskResponse; +import org.apache.dolphinscheduler.plugin.task.util.OSUtils; +import org.apache.dolphinscheduler.spi.task.AbstractParameters; +import org.apache.dolphinscheduler.spi.task.Property; +import org.apache.dolphinscheduler.spi.task.paramparser.ParamUtils; +import org.apache.dolphinscheduler.spi.task.paramparser.ParameterUtils; +import org.apache.dolphinscheduler.spi.task.request.TaskRequest; +import org.apache.dolphinscheduler.spi.utils.JSONUtils; + +import org.apache.commons.collections4.MapUtils; + +import java.io.File; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.StandardOpenOption; +import java.nio.file.attribute.FileAttribute; +import java.nio.file.attribute.PosixFilePermission; +import java.nio.file.attribute.PosixFilePermissions; +import java.util.HashMap; +import java.util.Map; +import java.util.Set; + +/** + * seatunnel task + */ +public class SeatunnelTask extends AbstractTaskExecutor { + + /** + * seatunnel parameters + */ + private SeatunnelParameters seatunnelParameters; + + /** + * shell command executor + */ + private ShellCommandExecutor shellCommandExecutor; + + /** + * taskExecutionContext + */ + private TaskRequest taskExecutionContext; + + /** + * constructor + * + * @param taskExecutionContext taskExecutionContext + */ + public SeatunnelTask(TaskRequest taskExecutionContext) { + super(taskExecutionContext); + + this.taskExecutionContext = taskExecutionContext; + this.shellCommandExecutor = new ShellCommandExecutor(this::logHandle, + taskExecutionContext, + logger); + } + + @Override + public void init() { + logger.info("seatunnel task params {}", taskExecutionContext.getTaskParams()); + + seatunnelParameters = JSONUtils.parseObject(taskExecutionContext.getTaskParams(), SeatunnelParameters.class); + + if (!seatunnelParameters.checkParameters()) { + throw new RuntimeException("seatunnel task params is not valid"); + } + } + + @Override + public void handle() throws Exception { + try { + // construct process + String command = buildCommand(); + TaskResponse commandExecuteResult = shellCommandExecutor.run(command); + setExitStatusCode(commandExecuteResult.getExitStatusCode()); + setAppIds(commandExecuteResult.getAppIds()); + setProcessId(commandExecuteResult.getProcessId()); + seatunnelParameters.dealOutParam(shellCommandExecutor.getVarPool()); + } catch (Exception e) { + logger.error("seatunnel task error", e); + setExitStatusCode(EXIT_CODE_FAILURE); + throw e; + } + } + + @Override + public void cancelApplication(boolean cancelApplication) throws Exception { + // cancel process + shellCommandExecutor.cancelApplication(); + } + + /** + * create command + * + * @return file name + * @throws Exception exception + */ + private String buildCommand() throws Exception { + // generate scripts + String fileName = String.format("%s/%s_node.%s", + taskExecutionContext.getExecutePath(), + taskExecutionContext.getTaskAppId(), OSUtils.isWindows() ? "bat" : "sh"); + + Path path = new File(fileName).toPath(); + + if (Files.exists(path)) { + return fileName; + } + + String script = seatunnelParameters.getRawScript().replaceAll("\\r\\n", "\n"); + script = parseScript(script); + seatunnelParameters.setRawScript(script); + + logger.info("raw script : {}", seatunnelParameters.getRawScript()); + logger.info("task execute path : {}", taskExecutionContext.getExecutePath()); + + Set perms = PosixFilePermissions.fromString(RWXR_XR_X); + FileAttribute> attr = PosixFilePermissions.asFileAttribute(perms); + + if (OSUtils.isWindows()) { + Files.createFile(path); + } else { + Files.createFile(path, attr); + } + + Files.write(path, seatunnelParameters.getRawScript().getBytes(), StandardOpenOption.APPEND); + + return fileName; + } + + @Override + public AbstractParameters getParameters() { + return seatunnelParameters; + } + + private String parseScript(String script) { + // combining local and global parameters + Map paramsMap = ParamUtils.convert(taskExecutionContext,getParameters()); + if (MapUtils.isEmpty(paramsMap)) { + paramsMap = new HashMap<>(); + } + if (MapUtils.isNotEmpty(taskExecutionContext.getParamsMap())) { + paramsMap.putAll(taskExecutionContext.getParamsMap()); + } + return ParameterUtils.convertParameterPlaceholders(script, ParamUtils.convert(paramsMap)); + } +} diff --git a/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelTaskChannel.java b/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelTaskChannel.java new file mode 100644 index 0000000000..54799aee05 --- /dev/null +++ b/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelTaskChannel.java @@ -0,0 +1,35 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.dolphinscheduler.plugin.task.seatunnel; + +import org.apache.dolphinscheduler.spi.task.TaskChannel; +import org.apache.dolphinscheduler.spi.task.request.TaskRequest; + +public class SeatunnelTaskChannel implements TaskChannel { + + @Override + public void cancelApplication(boolean status) { + + } + + @Override + public SeatunnelTask createTask(TaskRequest taskRequest) { + return new SeatunnelTask(taskRequest); + } + +} diff --git a/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelTaskChannelFactory.java b/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelTaskChannelFactory.java new file mode 100644 index 0000000000..15581dac75 --- /dev/null +++ b/dolphinscheduler-task-plugin/dolphinscheduler-task-seatunnel/src/main/java/org/apache/dolphinscheduler/plugin/task/seatunnel/SeatunnelTaskChannelFactory.java @@ -0,0 +1,45 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.dolphinscheduler.plugin.task.seatunnel; + +import org.apache.dolphinscheduler.spi.params.base.PluginParams; +import org.apache.dolphinscheduler.spi.task.TaskChannel; +import org.apache.dolphinscheduler.spi.task.TaskChannelFactory; + +import java.util.List; + +import com.google.auto.service.AutoService; + +@AutoService(TaskChannelFactory.class) +public class SeatunnelTaskChannelFactory implements TaskChannelFactory { + + @Override + public TaskChannel create() { + return new SeatunnelTaskChannel(); + } + + @Override + public String getName() { + return "SEATUNNEL"; + } + + @Override + public List getParams() { + return null; + } +} diff --git a/dolphinscheduler-task-plugin/pom.xml b/dolphinscheduler-task-plugin/pom.xml index 24c1a31c39..3295667d75 100644 --- a/dolphinscheduler-task-plugin/pom.xml +++ b/dolphinscheduler-task-plugin/pom.xml @@ -41,5 +41,6 @@ dolphinscheduler-task-sqoop dolphinscheduler-task-procedure dolphinscheduler-task-pigeon + dolphinscheduler-task-seatunnel diff --git a/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/canvas/taskbar.scss b/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/canvas/taskbar.scss index aec283f77f..a982dccbeb 100644 --- a/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/canvas/taskbar.scss +++ b/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/canvas/taskbar.scss @@ -102,8 +102,8 @@ &.icos-conditions { background-image: url("../images/task-icos/conditions.png"); } - &.icos-waterdrop { - background-image: url("../images/task-icos/waterdrop.png"); + &.icos-seatunnel { + background-image: url("../images/task-icos/seatunnel.png"); } &.icos-spark { background-image: url("../images/task-icos/spark.png"); @@ -162,8 +162,8 @@ &.icos-conditions { background-image: url("../images/task-icos/conditions_hover.png"); } - &.icos-waterdrop { - background-image: url("../images/task-icos/waterdrop_hover.png"); + &.icos-seatunnel { + background-image: url("../images/task-icos/seatunnel_hover.png"); } &.icos-spark { background-image: url("../images/task-icos/spark_hover.png"); diff --git a/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/config.js b/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/config.js index af81261117..4026e139ca 100755 --- a/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/config.js +++ b/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/config.js @@ -333,7 +333,7 @@ const tasksType = { desc: 'SWITCH', color: '#E46F13' }, - WATERDROP: { + SEATUNNEL: { desc: 'WATERDROP', color: '#646465', helperLinkDisable: true diff --git a/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/formModel/formModel.vue b/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/formModel/formModel.vue index 743f04b731..4b7ea189bb 100644 --- a/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/formModel/formModel.vue +++ b/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/formModel/formModel.vue @@ -464,15 +464,15 @@ :nodeData="nodeData" :postTasks="postTasks" > - - + - + @@ -507,7 +507,7 @@ import { findLocale } from '@/module/i18n/config' import mListBox from './tasks/_source/listBox' import mShell from './tasks/shell' - import mWaterdrop from './tasks/waterdrop' + import mSeatunnel from './tasks/seatunnel' import mSpark from './tasks/spark' import mFlink from './tasks/flink' import mPython from './tasks/python' @@ -1160,7 +1160,7 @@ mListBox, mMr, mShell, - mWaterdrop, + mSeatunnel, mSubProcess, mProcedure, mSql, diff --git a/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/formModel/tasks/waterdrop.vue b/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/formModel/tasks/seatunnel.vue similarity index 97% rename from dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/formModel/tasks/waterdrop.vue rename to dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/formModel/tasks/seatunnel.vue index a5ebc3443f..31536185cb 100644 --- a/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/formModel/tasks/waterdrop.vue +++ b/dolphinscheduler-ui/src/js/conf/home/pages/dag/_source/formModel/tasks/seatunnel.vue @@ -15,7 +15,7 @@ * limitations under the License. */