|
|
@ -19,25 +19,18 @@ package org.apache.dolphinscheduler.service.task; |
|
|
|
|
|
|
|
|
|
|
|
import static java.lang.String.format; |
|
|
|
import static java.lang.String.format; |
|
|
|
|
|
|
|
|
|
|
|
import org.apache.dolphinscheduler.common.enums.PluginType; |
|
|
|
|
|
|
|
import org.apache.dolphinscheduler.dao.PluginDao; |
|
|
|
|
|
|
|
import org.apache.dolphinscheduler.dao.entity.PluginDefine; |
|
|
|
|
|
|
|
import org.apache.dolphinscheduler.plugin.task.api.TaskChannel; |
|
|
|
import org.apache.dolphinscheduler.plugin.task.api.TaskChannel; |
|
|
|
import org.apache.dolphinscheduler.plugin.task.api.TaskChannelFactory; |
|
|
|
import org.apache.dolphinscheduler.plugin.task.api.TaskChannelFactory; |
|
|
|
import org.apache.dolphinscheduler.plugin.task.api.TaskPluginException; |
|
|
|
import org.apache.dolphinscheduler.plugin.task.api.TaskPluginException; |
|
|
|
import org.apache.dolphinscheduler.plugin.task.api.parameters.AbstractParameters; |
|
|
|
import org.apache.dolphinscheduler.plugin.task.api.parameters.AbstractParameters; |
|
|
|
import org.apache.dolphinscheduler.plugin.task.api.parameters.ParametersNode; |
|
|
|
import org.apache.dolphinscheduler.plugin.task.api.parameters.ParametersNode; |
|
|
|
import org.apache.dolphinscheduler.spi.params.PluginParamsTransfer; |
|
|
|
|
|
|
|
import org.apache.dolphinscheduler.spi.params.base.PluginParams; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
import java.util.Collections; |
|
|
|
import java.util.Collections; |
|
|
|
import java.util.HashSet; |
|
|
|
import java.util.HashMap; |
|
|
|
import java.util.List; |
|
|
|
|
|
|
|
import java.util.Map; |
|
|
|
import java.util.Map; |
|
|
|
import java.util.Objects; |
|
|
|
import java.util.Objects; |
|
|
|
import java.util.ServiceLoader; |
|
|
|
import java.util.ServiceLoader; |
|
|
|
import java.util.Set; |
|
|
|
import java.util.concurrent.atomic.AtomicBoolean; |
|
|
|
import java.util.concurrent.ConcurrentHashMap; |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
import org.slf4j.Logger; |
|
|
|
import org.slf4j.Logger; |
|
|
|
import org.slf4j.LoggerFactory; |
|
|
|
import org.slf4j.LoggerFactory; |
|
|
@ -47,23 +40,42 @@ import org.springframework.stereotype.Component; |
|
|
|
public class TaskPluginManager { |
|
|
|
public class TaskPluginManager { |
|
|
|
private static final Logger logger = LoggerFactory.getLogger(TaskPluginManager.class); |
|
|
|
private static final Logger logger = LoggerFactory.getLogger(TaskPluginManager.class); |
|
|
|
|
|
|
|
|
|
|
|
private final Map<String, TaskChannel> taskChannelMap = new ConcurrentHashMap<>(); |
|
|
|
private final Map<String, TaskChannelFactory> taskChannelFactoryMap = new HashMap<>(); |
|
|
|
|
|
|
|
private final Map<String, TaskChannel> taskChannelMap = new HashMap<>(); |
|
|
|
|
|
|
|
|
|
|
|
private final PluginDao pluginDao; |
|
|
|
private final AtomicBoolean loadedFlag = new AtomicBoolean(false); |
|
|
|
|
|
|
|
|
|
|
|
public TaskPluginManager(PluginDao pluginDao) { |
|
|
|
/** |
|
|
|
this.pluginDao = pluginDao; |
|
|
|
* Load task plugins from classpath. |
|
|
|
|
|
|
|
*/ |
|
|
|
|
|
|
|
public void loadPlugin() { |
|
|
|
|
|
|
|
if (!loadedFlag.compareAndSet(false, true)) { |
|
|
|
|
|
|
|
logger.warn("The task plugin has already been loaded"); |
|
|
|
|
|
|
|
return; |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
ServiceLoader.load(TaskChannelFactory.class).forEach(factory -> { |
|
|
|
|
|
|
|
final String name = factory.getName(); |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
logger.info("Registering task plugin: {}", name); |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
if (taskChannelFactoryMap.containsKey(name)) { |
|
|
|
|
|
|
|
throw new TaskPluginException(format("Duplicate task plugins named '%s'", name)); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
taskChannelFactoryMap.put(name, factory); |
|
|
|
|
|
|
|
taskChannelMap.put(name, factory.create()); |
|
|
|
|
|
|
|
|
|
|
|
private void loadTaskChannel(TaskChannelFactory taskChannelFactory) { |
|
|
|
logger.info("Registered task plugin: {}", name); |
|
|
|
TaskChannel taskChannel = taskChannelFactory.create(); |
|
|
|
}); |
|
|
|
taskChannelMap.put(taskChannelFactory.getName(), taskChannel); |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
public Map<String, TaskChannel> getTaskChannelMap() { |
|
|
|
public Map<String, TaskChannel> getTaskChannelMap() { |
|
|
|
return Collections.unmodifiableMap(taskChannelMap); |
|
|
|
return Collections.unmodifiableMap(taskChannelMap); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
public Map<String, TaskChannelFactory> getTaskChannelFactoryMap() { |
|
|
|
|
|
|
|
return Collections.unmodifiableMap(taskChannelFactoryMap); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
public TaskChannel getTaskChannel(String type) { |
|
|
|
public TaskChannel getTaskChannel(String type) { |
|
|
|
return this.getTaskChannelMap().get(type); |
|
|
|
return this.getTaskChannelMap().get(type); |
|
|
|
} |
|
|
|
} |
|
|
@ -85,30 +97,4 @@ public class TaskPluginManager { |
|
|
|
return taskChannel.parseParameters(parametersNode); |
|
|
|
return taskChannel.parseParameters(parametersNode); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
public void installPlugin() { |
|
|
|
|
|
|
|
final Set<String> names = new HashSet<>(); |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
ServiceLoader.load(TaskChannelFactory.class).forEach(factory -> { |
|
|
|
|
|
|
|
final String name = factory.getName(); |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
logger.info("Registering task plugin: {}", name); |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
if (!names.add(name)) { |
|
|
|
|
|
|
|
throw new TaskPluginException(format("Duplicate task plugins named '%s'", name)); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
loadTaskChannel(factory); |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
logger.info("Registered task plugin: {}", name); |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
List<PluginParams> params = factory.getParams(); |
|
|
|
|
|
|
|
String paramsJson = PluginParamsTransfer.transferParamsToJson(params); |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
PluginDefine pluginDefine = new PluginDefine(name, PluginType.TASK.getDesc(), paramsJson); |
|
|
|
|
|
|
|
int count = pluginDao.addOrUpdatePluginDefine(pluginDefine); |
|
|
|
|
|
|
|
if (count <= 0) { |
|
|
|
|
|
|
|
throw new TaskPluginException("Failed to update task plugin: " + name); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
}); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|