Commit ebbc1f71 by guominglei

工作流插件端保存

parent 8bb343fa
......@@ -26,7 +26,7 @@ public class ApiFlowController {
@PostMapping("publish")
@ApiOperation("发布工作流,并开始调度")
public String publishFlow(@RequestBody PluginPackage pluginPackage){
public String publishFlow(@RequestBody PluginPackage pluginPackage) throws Exception {
apiFlowService.publishFlow(pluginPackage);
return "SUCCESS";
}
......
......@@ -22,7 +22,7 @@ public interface ApiFlowService {
void reStartSchedule(String runId);
void publishFlow(PluginPackage pluginPackage);
void publishFlow(PluginPackage pluginPackage) throws Exception;
/**
* 判断工作流石佛存在
......
......@@ -2,27 +2,26 @@ package com.byit.service.impl;
import com.byit.enums.DagCheckEnum;
import com.byit.enums.FlowPropertyEnum;
import com.byit.enums.NodePropertyEnum;
import com.byit.job.dto.plugin.PluginBaseNode;
import com.byit.job.dto.plugin.PluginFlow;
import com.byit.job.dto.plugin.PluginNode;
import com.byit.job.dto.plugin.PluginPackage;
import com.byit.job.enums.plugin.PluginNodeTypeEnum;
import com.byit.job.utils.CronExpression;
import com.byit.mapper.FlowMapper;
import com.byit.mapper.WorkspaceMapper;
import com.byit.model.Flow;
import com.byit.model.Workspace;
import com.byit.mapper.*;
import com.byit.model.*;
import com.byit.service.ApiFlowService;
import com.byit.util.ApiFlowDagCheck;
import com.byit.utils.ValidationUtil;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.BeanUtils;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import javax.annotation.Resource;
import java.util.Date;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.*;
/**
* @description: 工作流的Api请求业务处理实现类
......@@ -38,15 +37,32 @@ public class ApiFlowServiceImpl implements ApiFlowService {
private FlowMapper flowMapper;
@Resource
private FlowVersionMapper flowVersionMapper;
@Resource
private WorkspaceMapper workspaceMapper;
@Resource
private NodeMapper nodeMapper;
@Resource
private NodeVersionMapper nodeVersionMapper;
@Resource
private NodeDependencyMapper nodeDependencyMapper;
@Resource
private NodeVersionDependencyMapper nodeVersionDependencyMapper;
@Override
@Transactional(rollbackFor = Exception.class)
public void publishFlow(PluginPackage pluginPackage) {
public void publishFlow(PluginPackage pluginPackage) throws Exception {
//校验参数
Workspace workspace = validate(pluginPackage);
saveFlow(pluginPackage.getFlow(), workspace.getWorkspaceId(), false);
//保存工作流
Flow flow = saveFlow(pluginPackage.getFlow(), workspace.getWorkspaceId(), false);
//生成工作流版本
FlowVersion flowVersion = saveFlowVersion(flow);
}
......@@ -56,7 +72,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
* @param isInnner
* @param workspaceId
*/
private void saveFlow(PluginFlow pluginFlow, Integer workspaceId, boolean isInnner) {
private Flow saveFlow(PluginFlow pluginFlow, Integer workspaceId, boolean isInnner) throws Exception {
Flow flow = new Flow();
flow.setFlowName(pluginFlow.getName());
......@@ -66,15 +82,23 @@ public class ApiFlowServiceImpl implements ApiFlowService {
flow.setAuthor(pluginFlow.getAuthor());
flow.setFlowDesc(pluginFlow.getDesc());
flow.setPrincipal(pluginFlow.getPrincipal());
flow.setWorkspaceId(workspaceId);
flow.setStartUp(FlowPropertyEnum.IS_START.getCode());
flow.setExecType(pluginFlow.getConfig().getExecType());
flow.setPriority(pluginFlow.getConfig().getPriority());
flow.setAlarmlAction(pluginFlow.getConfig().getAlarmlAction());
flow.setAlarmEmail(pluginFlow.getConfig().getAlarmEmail());
//设置可执行次数
//若没设置或者设置为-1,则响应设置剩余的执行次数
if (null != pluginFlow.getConfig().getRepeatCount() && !"-1".equals(pluginFlow.getConfig().getRepeatCount())){
flow.setRepeatCount(pluginFlow.getConfig().getRepeatCount());
flow.setRemainingCount(pluginFlow.getConfig().getRepeatCount());
}else {
flow.setRepeatCount(-1);
}
flow.setScheduleFollow(pluginFlow.getConfig().getScheduleFollow());
flow.setFlowCron(pluginFlow.getConfig().getFlowCron());
flow.setTriggerNextTime(StringUtils.isEmpty(pluginFlow.getConfig().getFlowCron()) ? null : new CronExpression(pluginFlow.getConfig().getFlowCron()).getNextValidTimeAfter(new Date()).getTime());
//设置超时时间,未设置默认30分钟
flow.setFlowTimeout(null == pluginFlow.getConfig().getFlowTimeout() ? 1000 * 60 * 30 : pluginFlow.getConfig().getFlowTimeout());
......@@ -86,28 +110,151 @@ public class ApiFlowServiceImpl implements ApiFlowService {
flow.setFlowId(oldFlow.getFlowId());
flow.setVersionName("V." + (versionTag + 1));
flowMapper.updateByIdSelective(flow);
deleteNodeByFlow(flow.getFlowId());
saveNode(pluginFlow.getNodeList(), flow.getFlowId());
updateNodeByFlow(pluginFlow.getNodeList(), flow);
}else {
flow.setVersionName("V.1");
Integer flowId = flowMapper.insertSelective(flow);
saveNode(pluginFlow.getNodeList(), flow.getFlowId());
saveNode(pluginFlow.getNodeList(), flow);
}
return flow;
}
/**
* 删除原有工作流下的节点
* @param flowId
* @param pluginNodeList
* @param flow
*/
private void deleteNodeByFlow(Integer flowId) {
private void updateNodeByFlow(List<PluginBaseNode> pluginNodeList, Flow flow) throws Exception{
//将工作流下的所有节点设为不在工作流调度上
log.info("将工作流【{}】所有的节点移下调度,并删除相关的依赖关系", flow.getFlowName());
//查询当前工作流下所有在调度上的节点
List<Node> nodeList = nodeMapper.findOnforkByFlowId(flow.getFlowId());
//删除在调度上的节点的相关依赖
nodeList.forEach(node -> nodeDependencyMapper.deleteByNodeId(node.getNodeId()));
nodeMapper.updateOnforkDown(flow.getFlowId());
//查找并删除工作流下的虚节点
List<Node> virtualNodeList = nodeMapper.findVirtualByFlowId(flow.getFlowId());
if (null != virtualNodeList && virtualNodeList.size() > 0){
for (Node node : virtualNodeList) {
Flow virtualFlow = new Flow();
virtualFlow.setFlowId(node.getMapFlowId());
//将内嵌的工作流设置为不是内嵌
virtualFlow.setIsInner(FlowPropertyEnum.ISNOT_INNER.getCode());
virtualFlow.setStartUp(FlowPropertyEnum.NO_START.getCode());
flowMapper.updateByIdSelective(virtualFlow);
}
}
//删除工作流下虚拟的节点
nodeMapper.deleteVirtualNode(flow.getFlowId());
//重新组建节点和依赖关系
HashMap<String, Integer> nameIdRel = new HashMap<>();
//保存节点信息
for(PluginBaseNode pluginNode : pluginNodeList){
Node node = nodeMapper.getByNameAndFlow(pluginNode.getName(), flow.getFlowId());
if (null == node){
node = new Node();
buildNode(pluginNode, flow, node);
Integer nodeId = nodeMapper.insertSelective(node);
}else {
buildNode(pluginNode, flow, node);
nodeMapper.updateByIdSelective(node);
}
nameIdRel.put(node.getNodeName(), node.getNodeId());
}
buildDepend(pluginNodeList, nameIdRel);
}
/**
* 保存工作流
* @param nodeList
* @param flowId
* @param flow
*/
private void saveNode(List<PluginBaseNode> nodeList, Flow flow) throws Exception{
HashMap<String, Integer> nameIdRel = new HashMap<>();
//保存节点信息
for(PluginBaseNode pluginNode : nodeList){
Node node = new Node();
buildNode(pluginNode, flow, node);
Integer nodeId = nodeMapper.insertSelective(node);
nameIdRel.put(node.getNodeName(), node.getNodeId());
}
buildDepend(nodeList, nameIdRel);
}
/**
* 设置节点信息
* @param pluginNode
* @param flow
* @param node
* @throws Exception
*/
private void buildNode(PluginBaseNode pluginNode, Flow flow, Node node) throws Exception{
node.setVersionName(flow.getVersionName());
node.setFlowId(flow.getFlowId());
node.setOnFork(NodePropertyEnum.ON_FORK.getCode());
node.setAddTime(new Date());
node.setAuthor(pluginNode.getAuthor());
node.setNodeName(pluginNode.getName());
if (PluginNodeTypeEnum.FLOW.getCode().equals(pluginNode.getType())){
//如果是内嵌工作流先保存工作流信息
Flow innerFlow = saveFlow((PluginFlow) pluginNode, flow.getWorkspaceId(), true);
node.setIsVirtual(NodePropertyEnum.IS_VIRTUAL.getCode());
node.setMapFlowId(innerFlow.getFlowId());
}else {
node.setHandlerName(((PluginNode)pluginNode).getHandlerName());
node.setJobType(((PluginNode)pluginNode).getJobType());
node.setRunSource(((PluginNode)pluginNode).getRunSource());
node.setRunSourceDesc(((PluginNode)pluginNode).getRunSourceDesc());
node.setRunCommand(((PluginNode)pluginNode).getRunCommand());
node.setRunParam(((PluginNode)pluginNode).getRunParam());
node.setSourcePrincipal(((PluginNode)pluginNode).getSourcePrincipal());
node.setBlockStrategy(StringUtils.isEmpty(((PluginNode)pluginNode).getConfig().getBlockStrategy()) ? "0" : ((PluginNode)pluginNode).getConfig().getBlockStrategy());
node.setGatewayToken(((PluginNode)pluginNode).getConfig().getGatewayToken());
node.setPluginToken(((PluginNode)pluginNode).getConfig().getPluginToken());
node.setNodeCron(((PluginNode)pluginNode).getConfig().getNodeCron());
node.setTriggerNextTime(StringUtils.isEmpty(((PluginNode)pluginNode).getConfig().getNodeCron()) ? null : new CronExpression(((PluginNode)pluginNode).getConfig().getNodeCron()).getNextValidTimeAfter(new Date()).getTime());
node.setPluginUrls(((PluginNode)pluginNode).getConfig().getPluginUrls());
node.setRoutingStrategy(StringUtils.isEmpty(((PluginNode)pluginNode).getConfig().getRoutingStrategy()) ? "RANDOM" : ((PluginNode)pluginNode).getConfig().getRoutingStrategy());
node.setPriority(StringUtils.isEmpty(((PluginNode)pluginNode).getConfig().getPriority()) ? "1" : ((PluginNode)pluginNode).getConfig().getPriority());
//设置失败重试
if (null != ((PluginNode)pluginNode).getConfig().getFailedRetryCount()){
node.setFailedRetryCount(((PluginNode)pluginNode).getConfig().getFailedRetryCount());
node.setFailedRetryInterval(((PluginNode)pluginNode).getConfig().getFailedRetryInterval());
}
//设置剩余执行次数
if (null != ((PluginNode)pluginNode).getConfig().getRepeatCount() && !"-1".equals(((PluginNode)pluginNode).getConfig().getRepeatCount())){
node.setRemainingCount(((PluginNode)pluginNode).getConfig().getRepeatCount());
node.setRepeatCount(((PluginNode)pluginNode).getConfig().getRepeatCount());
}else {
node.setRemainingCount(-1);
}
node.setIsVirtual(NodePropertyEnum.ISNOT_VIRTUAL.getCode());
}
}
/**
* 重新组织节点依赖关系
* @param nodeList
* @param nameIdRel
*/
private void saveNode(List<PluginBaseNode> nodeList, Integer flowId) {
public void buildDepend(List<PluginBaseNode> nodeList, HashMap<String, Integer> nameIdRel){
for (PluginBaseNode pluginNode : nodeList){
pluginNode.getDependNodeNameList().forEach(dependNodeName -> {
Integer nodeId = nameIdRel.get(pluginNode.getName());
Integer dependNodeId = nameIdRel.get(dependNodeName);
ValidationUtil.isTrueValidation(null == nodeId , pluginNode.getName() + "节点不存在!");
ValidationUtil.isTrueValidation(null == dependNodeId , "依赖的节点" + dependNodeName + "不存在!");
nodeDependencyMapper.insert(nodeId, dependNodeId);
});
}
}
/**
......@@ -169,9 +316,19 @@ public class ApiFlowServiceImpl implements ApiFlowService {
//校验父工作流和内嵌工作流的执行类型是否一致
ValidationUtil.isTrueValidation(!(flow.getConfig().getExecType().equals(((PluginFlow) node).getConfig().getExecType())), "内嵌工作流的调度配置必须和父工作流配置保持一致!");
validateFlow(workspace.getWorkspaceId(), (PluginFlow) node);
if(FlowPropertyEnum.NO_SCHEDULE.getCode().equals(flow.getConfig().getScheduleFollow())){
ValidationUtil.dataNotBank( ((PluginFlow) node).getConfig().getFlowCron(), "工作流设置为不跟随调度时节点必须设置调度时间!");
}
toFlowSet.add(node.getName());
lineMap.put(flow.getName() + "_" + node.getName(), 1);
checkFlow(workspace, (PluginFlow) node, flowNameSet, fromFlowSet, toFlowSet, lineMap);
}else {
ValidationUtil.dataNotBank(((PluginNode) node).getJobType(), "节点类型不允许为空!");
if(FlowPropertyEnum.NO_SCHEDULE.getCode().equals(flow.getConfig().getScheduleFollow())) {
ValidationUtil.dataNotNull(((PluginNode) node).getConfig(), "节点配置信息不允许为空!");
ValidationUtil.dataNotBank(((PluginNode) node).getConfig().getNodeCron(), "工作流设置为不跟随调度时节点必须设置调度时间!");
}
}
});
//校验同一个工作流下的任务节点是否存在环路
......@@ -194,6 +351,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
@Override
public void validateFlow(Integer workspaceId, PluginFlow pluginFlow) {
ValidationUtil.dataNotBank(pluginFlow.getName(), "工作流名称不允许为空!");
ValidationUtil.dataNotNull(pluginFlow.getNodeList(), "工作流下属节点不允许为空!");
ValidationUtil.dataNotNull(pluginFlow.getConfig(), "工作流配置不允许为空!");
ValidationUtil.dataNotBank(pluginFlow.getConfig().getExecType(), "工作流的调度类型不允许为空!");
if (FlowPropertyEnum.SCHEDULE_MODE.getCode().equals(pluginFlow.getConfig().getExecType())){
......@@ -213,6 +371,68 @@ public class ApiFlowServiceImpl implements ApiFlowService {
}
}
/**
* 工作流生成新版本
* @param flow
* @return
*/
private FlowVersion saveFlowVersion(Flow flow) {
Date currentDate = new Date();
HashMap<Integer, Integer> idVersionIdRel = new HashMap<>();
List<NodeDependencyKey> dependencyKeyList = new ArrayList<>();
FlowVersion flowVersion = new FlowVersion();
BeanUtils.copyProperties(flow, flowVersion);
if (!"V.1".equals(flow.getVersionName())){
//查找当前版本的工作流,如果有,则设置为不是,若没有跳过
FlowVersion oldFlowVersion = flowVersionMapper.getByFlowIdAndThis(flow.getFlowId());
if (null != oldFlowVersion){
oldFlowVersion.setVersionMark(FlowPropertyEnum.ISNOT_CURRENTVERSION.getCode());
flowVersionMapper.updateByIdSelective(oldFlowVersion);
}
}
//设置为是当前版本的工作流
flowVersion.setVersionMark(FlowPropertyEnum.IS_CURRENTVERSION.getCode());
flowVersion.setAddTime(currentDate);
flowVersionMapper.insertSelective(flowVersion);
//查询在调度上的节点
List<Node> nodeList = nodeMapper.findOnforkByFlowId(flow.getFlowId());
//将节点版本表所有的节点设置为不是当前版本
nodeVersionMapper.updateNotCurrentVersion(flow.getFlowId());
nodeList.forEach(node -> {
//查询本节点的依赖关系,若不为空则添加到总的依赖集合中
List<NodeDependencyKey> nodeDependList = nodeDependencyMapper.findByNodeId(node.getNodeId());
if (null != nodeDependList && nodeDependList.size() > 0){
dependencyKeyList.addAll(nodeDependList);
}
NodeVersion nodeVersion = new NodeVersion();
BeanUtils.copyProperties(node, nodeVersion);
nodeVersion.setAddTime(currentDate);
nodeVersion.setVersionMark(FlowPropertyEnum.IS_CURRENTVERSION.getCode());
if (null != node.getIsVirtual() && NodePropertyEnum.IS_VIRTUAL.equals(node.getIsVirtual())){
Flow innerFlow = flowMapper.getById(node.getMapFlowId());
FlowVersion innerFlowVersion = saveFlowVersion(innerFlow);
node.setMapFlowId(innerFlowVersion.getFlowVersionId());
}
nodeVersionMapper.insertSelective(nodeVersion);
idVersionIdRel.put(node.getNodeId(), nodeVersion.getNodeVersionId());
});
//添加版本的节点的依赖关系
dependencyKeyList.forEach(nodeDependencyKey -> {
Integer nodeVersionId = idVersionIdRel.get(nodeDependencyKey.getNodeId());
Integer dependVersionId = idVersionIdRel.get(nodeDependencyKey.getDependencyId());
ValidationUtil.dataNotNull(nodeVersionId , "没有找到对应的节点");
ValidationUtil.dataNotNull(dependVersionId , "没有找到对应的依赖的节点");
nodeVersionDependencyMapper.insert(nodeVersionId, dependVersionId);
});
return flowVersion;
}
@Override
public void deleteFlow(String flowName) {
......
......@@ -11,4 +11,10 @@ public interface FlowVersionMapper {
int updateByIdSelective(FlowVersion record);
/**
* 根据当前工作流id查找当前应用的版本
* @param flowId
* @return
*/
FlowVersion getByFlowIdAndThis(Integer flowId);
}
\ No newline at end of file
......@@ -24,5 +24,7 @@ public interface NodeDependencyMapper {
* @param nodeId
* @return
*/
List<Integer> findByNodeId(Integer nodeId);
List<Integer> findDependIdByNodeId(Integer nodeId);
List<NodeDependencyKey> findByNodeId(Integer nodeId);
}
\ No newline at end of file
package com.byit.mapper;
import com.byit.model.Node;
import org.apache.ibatis.annotations.Param;
import java.util.List;
......@@ -46,4 +47,12 @@ public interface NodeMapper {
* @param flowId
*/
void deleteVirtualNode(Integer flowId);
/**
* 根据工作流id和节点名称查询节点信息
* @param nodeName
* @param flowId
* @return
*/
Node getByNameAndFlow(@Param("nodeName") String nodeName, @Param("flowId") Integer flowId);
}
\ No newline at end of file
package com.byit.mapper;
import com.byit.model.NodeVersionDependencyKey;
import org.apache.ibatis.annotations.Param;
public interface NodeVersionDependencyMapper {
int deleteById(NodeVersionDependencyKey key);
int insert(@Param("nodeVersionId") Integer nodeVersionId, @Param("dependencyId") Integer dependencyId);
int insertSelective(NodeVersionDependencyKey record);
}
\ No newline at end of file
......@@ -11,4 +11,10 @@ public interface NodeVersionMapper {
int updateByIdSelective(NodeVersion record);
/**
* 将当前工作流下版本表中所有节点都设置为不是当前版本
* @param flowId
* @return
*/
int updateNotCurrentVersion(Integer flowId);
}
\ No newline at end of file
......@@ -2,9 +2,10 @@ package com.byit.model;
import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty;
import lombok.Data;
import java.io.Serializable;
import java.util.Date;
import lombok.Data;
/**
*
......@@ -97,9 +98,9 @@ public class FlowVersion implements Serializable {
private Integer repeatCount;
/**
* 当前版本的标志 this
* 当前版本的标志 0 是, 1 不是
*/
@ApiModelProperty("当前版本的标志 this")
@ApiModelProperty("当前版本的标志 0 是, 1 不是")
private String versionMark;
/**
......
package com.byit.model;
import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty;
import lombok.Data;
import java.io.Serializable;
/**
*
*/
@ApiModel
@Data
public class NodeVersionDependencyKey implements Serializable {
/**
* 节点版本id
*/
@ApiModelProperty("节点版本id")
private Integer nodeVersionId;
/**
* 依赖的id
*/
@ApiModelProperty("依赖的id")
private Integer dependencyId;
/**
*/
private static final long serialVersionUID = 1L;
}
\ No newline at end of file
......@@ -215,7 +215,7 @@ public class FlowServiceImpl implements FlowService {
nodeList.forEach(node -> {
NodeVo nodeVo = new NodeVo();
BeanUtils.copyProperties(node, nodeVo);
List<Integer> dependNodeIdList = nodeDependencyMapper.findByNodeId(node.getNodeId());
List<Integer> dependNodeIdList = nodeDependencyMapper.findDependIdByNodeId(node.getNodeId());
nodeVoList.add(nodeVo);
});
flowVo.setNodeVoList(nodeVoList);
......
......@@ -30,6 +30,7 @@
flow_name, flow_node_count, alarml_action, schedule_follow, priority, remove_mark,
repeat_count, version_mark, workspace_id, author, principal, version_name, is_inner
</sql>
<select id="getById" parameterType="java.lang.Integer" resultMap="BaseResultMap">
<!-- generated @mbg.generated date: 2019-12-31 -->
select
......@@ -37,29 +38,22 @@
from flow_version
where flow_version_id = #{flowVersionId,jdbcType=INTEGER}
</select>
<select id="getByFlowIdAndThis" parameterType="java.lang.Integer" resultMap="BaseResultMap">
select
<include refid="Base_Column_List" />
from flow_version
where flow_id = #{flowId,jdbcType=INTEGER}
and version_mark = '0'
</select>
<delete id="deleteById" parameterType="java.lang.Integer">
<!-- generated @mbg.generated date: 2019-12-31 -->
delete from flow_version
where flow_version_id = #{flowVersionId,jdbcType=INTEGER}
</delete>
<insert id="insert" parameterType="com.byit.model.FlowVersion">
<!-- generated @mbg.generated date: 2019-12-31 -->
insert into flow_version (flow_version_id, add_time, alarm_email,
exec_type, flow_cron, flow_desc,
flow_id, flow_name, flow_node_count,
alarml_action, schedule_follow, priority,
remove_mark, repeat_count, version_mark,
workspace_id, author, principal,
version_name, is_inner)
values (#{flowVersionId,jdbcType=INTEGER}, #{addTime,jdbcType=TIMESTAMP}, #{alarmEmail,jdbcType=VARCHAR},
#{execType,jdbcType=CHAR}, #{flowCron,jdbcType=VARCHAR}, #{flowDesc,jdbcType=VARCHAR},
#{flowId,jdbcType=INTEGER}, #{flowName,jdbcType=VARCHAR}, #{flowNodeCount,jdbcType=INTEGER},
#{alarmlAction,jdbcType=CHAR}, #{scheduleFollow,jdbcType=CHAR}, #{priority,jdbcType=CHAR},
#{removeMark,jdbcType=CHAR}, #{repeatCount,jdbcType=INTEGER}, #{versionMark,jdbcType=CHAR},
#{workspaceId,jdbcType=INTEGER}, #{author,jdbcType=VARCHAR}, #{principal,jdbcType=VARCHAR},
#{versionName,jdbcType=VARCHAR}, #{isInner,jdbcType=CHAR})
</insert>
<insert id="insertSelective" parameterType="com.byit.model.FlowVersion">
<insert id="insertSelective" useGeneratedKeys="true" keyProperty="flowVersionId" parameterType="com.byit.model.FlowVersion">
<!-- generated @mbg.generated date: 2019-12-31 -->
insert into flow_version
<trim prefix="(" suffix=")" suffixOverrides=",">
......@@ -251,28 +245,5 @@
</set>
where flow_version_id = #{flowVersionId,jdbcType=INTEGER}
</update>
<update id="updateById" parameterType="com.byit.model.FlowVersion">
<!-- generated @mbg.generated date: 2019-12-31 -->
update flow_version
set add_time = #{addTime,jdbcType=TIMESTAMP},
alarm_email = #{alarmEmail,jdbcType=VARCHAR},
exec_type = #{execType,jdbcType=CHAR},
flow_cron = #{flowCron,jdbcType=VARCHAR},
flow_desc = #{flowDesc,jdbcType=VARCHAR},
flow_id = #{flowId,jdbcType=INTEGER},
flow_name = #{flowName,jdbcType=VARCHAR},
flow_node_count = #{flowNodeCount,jdbcType=INTEGER},
alarml_action = #{alarmlAction,jdbcType=CHAR},
schedule_follow = #{scheduleFollow,jdbcType=CHAR},
priority = #{priority,jdbcType=CHAR},
remove_mark = #{removeMark,jdbcType=CHAR},
repeat_count = #{repeatCount,jdbcType=INTEGER},
version_mark = #{versionMark,jdbcType=CHAR},
workspace_id = #{workspaceId,jdbcType=INTEGER},
author = #{author,jdbcType=VARCHAR},
principal = #{principal,jdbcType=VARCHAR},
version_name = #{versionName,jdbcType=VARCHAR},
is_inner = #{isInner,jdbcType=CHAR}
where flow_version_id = #{flowVersionId,jdbcType=INTEGER}
</update>
</mapper>
\ No newline at end of file
......@@ -6,13 +6,32 @@
<id column="node_id" jdbcType="INTEGER" property="nodeId" />
<id column="dependency_id" jdbcType="INTEGER" property="dependencyId" />
</resultMap>
<select id="findDependIdByNodeId" parameterType="integer" resultType="integer">
select dependency_id
from node_dependency
where node_id = #{nodeId,jdbcType=INTEGER}
</select>
<select id="findByNodeId" parameterType="integer" resultMap="BaseResultMap">
select node_id, dependency_id
from node_dependency
where node_id = #{nodeId,jdbcType=INTEGER}
</select>
<delete id="deleteById" parameterType="com.byit.model.NodeDependencyKey">
<!-- generated @mbg.generated date: 2019-12-31 -->
delete from node_dependency
where node_id = #{nodeId,jdbcType=INTEGER}
and dependency_id = #{dependencyId,jdbcType=INTEGER}
</delete>
<insert id="insert" parameterType="com.byit.model.NodeDependencyKey">
<delete id="deleteByNodeId" parameterType="integer">
delete from node_dependency
where node_id = #{nodeId,jdbcType=INTEGER}
</delete>
<insert id="insert">
<!-- generated @mbg.generated date: 2019-12-31 -->
insert into node_dependency (node_id, dependency_id)
values (#{nodeId,jdbcType=INTEGER}, #{dependencyId,jdbcType=INTEGER})
......
......@@ -60,37 +60,50 @@
from node
where node_id = #{nodeId,jdbcType=INTEGER}
</select>
<select id="findOnforkByFlowId" parameterType="integer" resultMap="BaseResultMap">
select
<include refid="Base_Column_List" />
,
<include refid="Blob_Column_List" />
from node
where flow_id = #{flowId,jdbcType=INTEGER}
and on_fork = '0';
</select>
<select id="findVirtualByFlowId" parameterType="integer" resultMap="BaseResultMap">
select
<include refid="Base_Column_List" />
,
<include refid="Blob_Column_List" />
from node
where flow_id = #{flowId,jdbcType=INTEGER}
and is_virtual = '0';
</select>
<select id="getByNameAndFlow" resultMap="BaseResultMap">
select
<include refid="Base_Column_List" />
,
<include refid="Blob_Column_List" />
from node
where flow_id = #{flowId,jdbcType=INTEGER}
and node_name = #{nodeName,jdbcType=VARCHAR}
</select>
<delete id="deleteVirtualNode">
delete from node
where flow_id = #{flowId,jdbcType=INTEGER}
and is_virtual = '0';
</delete>
<delete id="deleteById" parameterType="java.lang.Integer">
<!-- generated @mbg.generated date: 2019-12-31 -->
delete from node
where node_id = #{nodeId,jdbcType=INTEGER}
</delete>
<insert id="insert" parameterType="com.byit.model.Node">
<!-- generated @mbg.generated date: 2019-12-31 -->
insert into node (node_id, block_strategy, plugin_token,
failed_retry_count, flow_id, gateway_token,
job_type, handler_name, node_cron,
node_desc, node_name, map_flow_id,
node_timeout, is_virtual, plugin_urls,
priority, remaining_count, repeat_count,
failed_retry_interval, routing_strategy, run_param,
run_source_desc, script_urls, source_principal,
source_update_time, trigger_next_time, author,
add_time, version_name, on_fork,
run_command, run_source)
values (#{nodeId,jdbcType=INTEGER}, #{blockStrategy,jdbcType=VARCHAR}, #{pluginToken,jdbcType=VARCHAR},
#{failedRetryCount,jdbcType=INTEGER}, #{flowId,jdbcType=INTEGER}, #{gatewayToken,jdbcType=VARCHAR},
#{jobType,jdbcType=VARCHAR}, #{handlerName,jdbcType=VARCHAR}, #{nodeCron,jdbcType=VARCHAR},
#{nodeDesc,jdbcType=VARCHAR}, #{nodeName,jdbcType=VARCHAR}, #{mapFlowId,jdbcType=INTEGER},
#{nodeTimeout,jdbcType=BIGINT}, #{isVirtual,jdbcType=CHAR}, #{pluginUrls,jdbcType=VARCHAR},
#{priority,jdbcType=CHAR}, #{remainingCount,jdbcType=INTEGER}, #{repeatCount,jdbcType=INTEGER},
#{failedRetryInterval,jdbcType=BIGINT}, #{routingStrategy,jdbcType=VARCHAR}, #{runParam,jdbcType=VARCHAR},
#{runSourceDesc,jdbcType=VARCHAR}, #{scriptUrls,jdbcType=VARCHAR}, #{sourcePrincipal,jdbcType=VARCHAR},
#{sourceUpdateTime,jdbcType=DATE}, #{triggerNextTime,jdbcType=BIGINT}, #{author,jdbcType=VARCHAR},
#{addTime,jdbcType=TIMESTAMP}, #{versionName,jdbcType=VARCHAR}, #{onFork,jdbcType=CHAR},
#{runCommand,jdbcType=VARCHAR}, #{runSource,jdbcType=LONGVARCHAR})
</insert>
<insert id="insertSelective" parameterType="com.byit.model.Node">
<insert id="insertSelective" useGeneratedKeys="true" keyProperty="nodeId" parameterType="com.byit.model.Node">
<!-- generated @mbg.generated date: 2019-12-31 -->
insert into node
<trim prefix="(" suffix=")" suffixOverrides=",">
......@@ -290,6 +303,16 @@
</if>
</trim>
</insert>
<!--将节点更改为在调度上-->
<update id="updateOnforkUp" parameterType="integer" >
update node set on_fork = '0' where node_id = #{nodeId, jdbcType=INTEGER}
</update>
<update id="updateOnforkDown" parameterType="integer" >
update node set on_fork = '1' where flow_id = #{flowId, jdbcType=INTEGER}
</update>
<update id="updateByIdSelective" parameterType="com.byit.model.Node">
<!-- generated @mbg.generated date: 2019-12-31 -->
update node
......@@ -390,75 +413,5 @@
</set>
where node_id = #{nodeId,jdbcType=INTEGER}
</update>
<update id="updateByPrimaryKeyWithBLOBs" parameterType="com.byit.model.Node">
<!-- generated @mbg.generated date: 2019-12-31 -->
update node
set block_strategy = #{blockStrategy,jdbcType=VARCHAR},
plugin_token = #{pluginToken,jdbcType=VARCHAR},
failed_retry_count = #{failedRetryCount,jdbcType=INTEGER},
flow_id = #{flowId,jdbcType=INTEGER},
gateway_token = #{gatewayToken,jdbcType=VARCHAR},
job_type = #{jobType,jdbcType=VARCHAR},
handler_name = #{handlerName,jdbcType=VARCHAR},
node_cron = #{nodeCron,jdbcType=VARCHAR},
node_desc = #{nodeDesc,jdbcType=VARCHAR},
node_name = #{nodeName,jdbcType=VARCHAR},
map_flow_id = #{mapFlowId,jdbcType=INTEGER},
node_timeout = #{nodeTimeout,jdbcType=BIGINT},
is_virtual = #{isVirtual,jdbcType=CHAR},
plugin_urls = #{pluginUrls,jdbcType=VARCHAR},
priority = #{priority,jdbcType=CHAR},
remaining_count = #{remainingCount,jdbcType=INTEGER},
repeat_count = #{repeatCount,jdbcType=INTEGER},
failed_retry_interval = #{failedRetryInterval,jdbcType=BIGINT},
routing_strategy = #{routingStrategy,jdbcType=VARCHAR},
run_param = #{runParam,jdbcType=VARCHAR},
run_source_desc = #{runSourceDesc,jdbcType=VARCHAR},
script_urls = #{scriptUrls,jdbcType=VARCHAR},
source_principal = #{sourcePrincipal,jdbcType=VARCHAR},
source_update_time = #{sourceUpdateTime,jdbcType=DATE},
trigger_next_time = #{triggerNextTime,jdbcType=BIGINT},
author = #{author,jdbcType=VARCHAR},
add_time = #{addTime,jdbcType=TIMESTAMP},
version_name = #{versionName,jdbcType=VARCHAR},
on_fork = #{onFork,jdbcType=CHAR},
run_command = #{runCommand,jdbcType=VARCHAR},
run_source = #{runSource,jdbcType=LONGVARCHAR}
where node_id = #{nodeId,jdbcType=INTEGER}
</update>
<update id="updateById" parameterType="com.byit.model.Node">
<!-- generated @mbg.generated date: 2019-12-31 -->
update node
set block_strategy = #{blockStrategy,jdbcType=VARCHAR},
plugin_token = #{pluginToken,jdbcType=VARCHAR},
failed_retry_count = #{failedRetryCount,jdbcType=INTEGER},
flow_id = #{flowId,jdbcType=INTEGER},
gateway_token = #{gatewayToken,jdbcType=VARCHAR},
job_type = #{jobType,jdbcType=VARCHAR},
handler_name = #{handlerName,jdbcType=VARCHAR},
node_cron = #{nodeCron,jdbcType=VARCHAR},
node_desc = #{nodeDesc,jdbcType=VARCHAR},
node_name = #{nodeName,jdbcType=VARCHAR},
map_flow_id = #{mapFlowId,jdbcType=INTEGER},
node_timeout = #{nodeTimeout,jdbcType=BIGINT},
is_virtual = #{isVirtual,jdbcType=CHAR},
plugin_urls = #{pluginUrls,jdbcType=VARCHAR},
priority = #{priority,jdbcType=CHAR},
remaining_count = #{remainingCount,jdbcType=INTEGER},
repeat_count = #{repeatCount,jdbcType=INTEGER},
failed_retry_interval = #{failedRetryInterval,jdbcType=BIGINT},
routing_strategy = #{routingStrategy,jdbcType=VARCHAR},
run_param = #{runParam,jdbcType=VARCHAR},
run_source_desc = #{runSourceDesc,jdbcType=VARCHAR},
script_urls = #{scriptUrls,jdbcType=VARCHAR},
source_principal = #{sourcePrincipal,jdbcType=VARCHAR},
source_update_time = #{sourceUpdateTime,jdbcType=DATE},
trigger_next_time = #{triggerNextTime,jdbcType=BIGINT},
author = #{author,jdbcType=VARCHAR},
add_time = #{addTime,jdbcType=TIMESTAMP},
version_name = #{versionName,jdbcType=VARCHAR},
on_fork = #{onFork,jdbcType=CHAR},
run_command = #{runCommand,jdbcType=VARCHAR}
where node_id = #{nodeId,jdbcType=INTEGER}
</update>
</mapper>
\ No newline at end of file
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.byit.mapper.NodeVersionDependencyMapper">
<resultMap id="BaseResultMap" type="com.byit.model.NodeVersionDependencyKey">
<!-- generated @mbg.generated date: 2020-01-03 -->
<id column="node_version_id" jdbcType="INTEGER" property="nodeVersionId" />
<id column="dependency_id" jdbcType="INTEGER" property="dependencyId" />
</resultMap>
<delete id="deleteById" parameterType="com.byit.model.NodeVersionDependencyKey">
<!-- generated @mbg.generated date: 2020-01-03 -->
delete from node_version_dependency
where node_version_id = #{nodeVersionId,jdbcType=INTEGER}
and dependency_id = #{dependencyId,jdbcType=INTEGER}
</delete>
<insert id="insert">
<!-- generated @mbg.generated date: 2020-01-03 -->
insert into node_version_dependency (node_version_id, dependency_id)
values (#{nodeVersionId,jdbcType=INTEGER}, #{dependencyId,jdbcType=INTEGER})
</insert>
<insert id="insertSelective" parameterType="com.byit.model.NodeVersionDependencyKey">
<!-- generated @mbg.generated date: 2020-01-03 -->
insert into node_version_dependency
<trim prefix="(" suffix=")" suffixOverrides=",">
<if test="nodeVersionId != null">
node_version_id,
</if>
<if test="dependencyId != null">
dependency_id,
</if>
</trim>
<trim prefix="values (" suffix=")" suffixOverrides=",">
<if test="nodeVersionId != null">
#{nodeVersionId,jdbcType=INTEGER},
</if>
<if test="dependencyId != null">
#{dependencyId,jdbcType=INTEGER},
</if>
</trim>
</insert>
</mapper>
\ No newline at end of file
......@@ -54,6 +54,7 @@
<!-- generated @mbg.generated date: 2019-12-31 -->
run_source
</sql>
<select id="getById" parameterType="java.lang.Integer" resultMap="ResultMapWithBLOBs">
<!-- generated @mbg.generated date: 2019-12-31 -->
select
......@@ -63,38 +64,14 @@
from node_version
where node_version_id = #{nodeVersionId,jdbcType=INTEGER}
</select>
<delete id="deleteById" parameterType="java.lang.Integer">
<!-- generated @mbg.generated date: 2019-12-31 -->
delete from node_version
where node_version_id = #{nodeVersionId,jdbcType=INTEGER}
</delete>
<insert id="insert" parameterType="com.byit.model.NodeVersion">
<!-- generated @mbg.generated date: 2019-12-31 -->
insert into node_version (node_version_id, node_id, block_strategy,
pugin_token, failed_retry_count, flow_id,
flow_version_id, gateway_token, job_type,
handler_name, node_cron, node_desc,
node_name, map_flow_id, node_timeout,
is_virtual, plugin_urls, priority,
remove_mark, repeat_count, failed_retry_interval,
routing_strategy, run_param, run_source_desc,
script_urls, source_principal, source_update_time,
version_mark, author, add_time,
version_name, on_fork, run_command,
run_source)
values (#{nodeVersionId,jdbcType=INTEGER}, #{nodeId,jdbcType=INTEGER}, #{blockStrategy,jdbcType=VARCHAR},
#{puginToken,jdbcType=VARCHAR}, #{failedRetryCount,jdbcType=INTEGER}, #{flowId,jdbcType=INTEGER},
#{flowVersionId,jdbcType=INTEGER}, #{gatewayToken,jdbcType=VARCHAR}, #{jobType,jdbcType=VARCHAR},
#{handlerName,jdbcType=VARCHAR}, #{nodeCron,jdbcType=VARCHAR}, #{nodeDesc,jdbcType=VARCHAR},
#{nodeName,jdbcType=VARCHAR}, #{mapFlowId,jdbcType=INTEGER}, #{nodeTimeout,jdbcType=BIGINT},
#{isVirtual,jdbcType=VARCHAR}, #{pluginUrls,jdbcType=VARCHAR}, #{priority,jdbcType=CHAR},
#{removeMark,jdbcType=CHAR}, #{repeatCount,jdbcType=INTEGER}, #{failedRetryInterval,jdbcType=BIGINT},
#{routingStrategy,jdbcType=VARCHAR}, #{runParam,jdbcType=VARCHAR}, #{runSourceDesc,jdbcType=VARCHAR},
#{scriptUrls,jdbcType=VARCHAR}, #{sourcePrincipal,jdbcType=VARCHAR}, #{sourceUpdateTime,jdbcType=TIMESTAMP},
#{versionMark,jdbcType=VARCHAR}, #{author,jdbcType=VARCHAR}, #{addTime,jdbcType=TIMESTAMP},
#{versionName,jdbcType=VARCHAR}, #{onFork,jdbcType=CHAR}, #{runCommand,jdbcType=VARCHAR},
#{runSource,jdbcType=LONGVARCHAR})
</insert>
<insert id="insertSelective" parameterType="com.byit.model.NodeVersion">
<!-- generated @mbg.generated date: 2019-12-31 -->
insert into node_version
......@@ -307,6 +284,13 @@
</if>
</trim>
</insert>
<update id="updateNotCurrentVersion" parameterType="integer">
update node_version
set version_mark = '1'
where flow_id = #{flowId}
</update>
<update id="updateByIdSelective" parameterType="com.byit.model.NodeVersion">
<!-- generated @mbg.generated date: 2019-12-31 -->
update node_version
......@@ -413,79 +397,5 @@
</set>
where node_version_id = #{nodeVersionId,jdbcType=INTEGER}
</update>
<update id="updateByPrimaryKeyWithBLOBs" parameterType="com.byit.model.NodeVersion">
<!-- generated @mbg.generated date: 2019-12-31 -->
update node_version
set node_id = #{nodeId,jdbcType=INTEGER},
block_strategy = #{blockStrategy,jdbcType=VARCHAR},
pugin_token = #{puginToken,jdbcType=VARCHAR},
failed_retry_count = #{failedRetryCount,jdbcType=INTEGER},
flow_id = #{flowId,jdbcType=INTEGER},
flow_version_id = #{flowVersionId,jdbcType=INTEGER},
gateway_token = #{gatewayToken,jdbcType=VARCHAR},
job_type = #{jobType,jdbcType=VARCHAR},
handler_name = #{handlerName,jdbcType=VARCHAR},
node_cron = #{nodeCron,jdbcType=VARCHAR},
node_desc = #{nodeDesc,jdbcType=VARCHAR},
node_name = #{nodeName,jdbcType=VARCHAR},
map_flow_id = #{mapFlowId,jdbcType=INTEGER},
node_timeout = #{nodeTimeout,jdbcType=BIGINT},
is_virtual = #{isVirtual,jdbcType=VARCHAR},
plugin_urls = #{pluginUrls,jdbcType=VARCHAR},
priority = #{priority,jdbcType=CHAR},
remove_mark = #{removeMark,jdbcType=CHAR},
repeat_count = #{repeatCount,jdbcType=INTEGER},
failed_retry_interval = #{failedRetryInterval,jdbcType=BIGINT},
routing_strategy = #{routingStrategy,jdbcType=VARCHAR},
run_param = #{runParam,jdbcType=VARCHAR},
run_source_desc = #{runSourceDesc,jdbcType=VARCHAR},
script_urls = #{scriptUrls,jdbcType=VARCHAR},
source_principal = #{sourcePrincipal,jdbcType=VARCHAR},
source_update_time = #{sourceUpdateTime,jdbcType=TIMESTAMP},
version_mark = #{versionMark,jdbcType=VARCHAR},
author = #{author,jdbcType=VARCHAR},
add_time = #{addTime,jdbcType=TIMESTAMP},
version_name = #{versionName,jdbcType=VARCHAR},
on_fork = #{onFork,jdbcType=CHAR},
run_command = #{runCommand,jdbcType=VARCHAR},
run_source = #{runSource,jdbcType=LONGVARCHAR}
where node_version_id = #{nodeVersionId,jdbcType=INTEGER}
</update>
<update id="updateById" parameterType="com.byit.model.NodeVersion">
<!-- generated @mbg.generated date: 2019-12-31 -->
update node_version
set node_id = #{nodeId,jdbcType=INTEGER},
block_strategy = #{blockStrategy,jdbcType=VARCHAR},
pugin_token = #{puginToken,jdbcType=VARCHAR},
failed_retry_count = #{failedRetryCount,jdbcType=INTEGER},
flow_id = #{flowId,jdbcType=INTEGER},
flow_version_id = #{flowVersionId,jdbcType=INTEGER},
gateway_token = #{gatewayToken,jdbcType=VARCHAR},
job_type = #{jobType,jdbcType=VARCHAR},
handler_name = #{handlerName,jdbcType=VARCHAR},
node_cron = #{nodeCron,jdbcType=VARCHAR},
node_desc = #{nodeDesc,jdbcType=VARCHAR},
node_name = #{nodeName,jdbcType=VARCHAR},
map_flow_id = #{mapFlowId,jdbcType=INTEGER},
node_timeout = #{nodeTimeout,jdbcType=BIGINT},
is_virtual = #{isVirtual,jdbcType=VARCHAR},
plugin_urls = #{pluginUrls,jdbcType=VARCHAR},
priority = #{priority,jdbcType=CHAR},
remove_mark = #{removeMark,jdbcType=CHAR},
repeat_count = #{repeatCount,jdbcType=INTEGER},
failed_retry_interval = #{failedRetryInterval,jdbcType=BIGINT},
routing_strategy = #{routingStrategy,jdbcType=VARCHAR},
run_param = #{runParam,jdbcType=VARCHAR},
run_source_desc = #{runSourceDesc,jdbcType=VARCHAR},
script_urls = #{scriptUrls,jdbcType=VARCHAR},
source_principal = #{sourcePrincipal,jdbcType=VARCHAR},
source_update_time = #{sourceUpdateTime,jdbcType=TIMESTAMP},
version_mark = #{versionMark,jdbcType=VARCHAR},
author = #{author,jdbcType=VARCHAR},
add_time = #{addTime,jdbcType=TIMESTAMP},
version_name = #{versionName,jdbcType=VARCHAR},
on_fork = #{onFork,jdbcType=CHAR},
run_command = #{runCommand,jdbcType=VARCHAR}
where node_version_id = #{nodeVersionId,jdbcType=INTEGER}
</update>
</mapper>
\ No newline at end of file
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment