Commit 136a6485 by guominglei

工作流的保存和创建

parent 31a9c9cf
package com.byit.enums;
/**
* @description: 节点类型
* @author: gml
* @create: 2019-12-25 12:01
*/
public enum NodeTypeEnum {
SHELL("SHELL"),
JAVA("JAVA"),
PYTHON("PYTHON"),
SQL("SQL"),
SCRIPT("SCRIPT"),
;
private String code;
private NodeTypeEnum(String code){
this.code = code;
}
public String getCode(){
return this.code;
}
public NodeTypeEnum getTypeByCode(String code){
for (NodeTypeEnum typeEnum : NodeTypeEnum.values()) {
if (typeEnum.getCode().equals(code)){
return typeEnum;
}
}
return null;
}
}
package com.byit.mapper;
import com.byit.model.FlowDependentKey;
import org.apache.ibatis.annotations.Param;
import java.util.List;
public interface FlowDependentMapper {
int deleteById(FlowDependentKey key);
int insert(FlowDependentKey record);
int insert(@Param("flowId")Integer flowId, @Param("dependFlowId") Integer dependFlowId);
int insertSelective(FlowDependentKey record);
int deleteByFlowId(Integer flowId);
/**
* 根据工作流Id查询本工作流的所有依赖关系
* @param flowId
* @return
*/
List<FlowDependentKey> findByFlowId(Integer flowId);
}
\ No newline at end of file
......@@ -11,4 +11,18 @@ public interface NodeMapper {
int updateByIdSelective(Node record);
/**
* 将节点更改为在调度上
* @param nodeId
* @return
*/
int updateOnforkUp(Integer nodeId);
/**
* 将工作流下的所有节点设为在调度下
* @param flowId
* @return
*/
int updateOnforkDown(Integer flowId);
}
\ No newline at end of file
......@@ -142,6 +142,12 @@ public class FlowVo implements Serializable {
private Integer remainingCount;
/**
* 工作流依赖的工作流id
*/
@ApiModelProperty("工作流依赖的工作流id")
private String flowDepend;
/**
* 在工作流调度中的节点
*/
@ApiModelProperty("在工作流调度中的节点")
......
package com.byit.model.vo;
import io.swagger.annotations.ApiModelProperty;
import lombok.Data;
import java.io.Serializable;
import java.util.Date;
......@@ -10,6 +11,7 @@ import java.util.Date;
* @author: gml
* @create: 2019-12-23 15:39
*/
@Data
public class NodeVo implements Serializable {
/**
* 当前版本节点主键
......
package com.byit.service.impl;
import com.byit.job.utils.CronExpression;
import com.byit.mapper.FlowDependentMapper;
import com.byit.mapper.FlowMapper;
import com.byit.mapper.NodeMapper;
import com.byit.model.Flow;
import com.byit.model.FlowDependentKey;
import com.byit.model.Node;
import com.byit.model.vo.FlowVo;
import com.byit.model.vo.NodeVo;
import com.byit.service.FlowService;
import com.byit.util.FlowInitUtil;
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.Arrays;
import java.util.Date;
import java.util.List;
/**
* @description: 工作流业务逻辑实现类
......@@ -18,33 +29,110 @@ import javax.annotation.Resource;
*/
@Service
@Slf4j
public class FlowServiceImpl implements FlowService {
@Resource
private FlowMapper flowMapper;
@Resource
private NodeMapper nodeMapper;
@Resource
private FlowDependentMapper flowDependentMapper;
@Override
public Flow saveJobFlow(FlowVo flowVo) {
return null;
//TODO 未做参数校验
Flow flow = new Flow();
BeanUtils.copyProperties(flowVo, flow);
//只要修改就设置为未启动
flow.setStartUp("0");
//设置在工作流上的节点相关信息
List<NodeVo> nodeVoList = flowVo.getNodeVoList();
if (null != nodeVoList && nodeVoList.size() > 0){
//设置在工作流上的节点数目
int nodeCount = nodeVoList.size();
log.info("将【{}】个节点在工作流【{}】调度流程中", nodeCount, flow.getFlowName());
flow.setFlowNodeCount(nodeCount);
//将工作流下的所有节点设为不在工作流调度上
log.info("将工作流【{}】所有的节点移下调度", flow.getFlowName());
nodeMapper.updateOnforkDown(flow.getFlowId());
//将对应的节点更改为在调度上
log.info("给工作流【{}】添加节点调度调度", flow.getFlowName());
nodeVoList.forEach(nodeVo -> nodeMapper.updateOnforkUp(nodeVo.getNodeId()));
}
//重组工作流依赖
reFlowDepend(flowVo);
//更新工作流的相关信息
flowMapper.updateByIdSelective(flow);
return flow;
}
/**
* 重新组织工作流的依赖关系
* @param flowVo
*/
private void reFlowDepend(FlowVo flowVo) {
//先去除所有的工作流依赖
List<FlowDependentKey> flowDependentKeyList = flowDependentMapper.findByFlowId(flowVo.getFlowId());
//解除跟上游工作流的绑定
log.info("将【{}】工作流上游工作流设置为没有下游工作流", flowVo.getFlowName());
for (FlowDependentKey flowDependentKey: flowDependentKeyList) {
Flow flow = flowMapper.getById(flowDependentKey.getDependFlowId());
flow.setIsHaveDepend("0");
flowMapper.updateByIdSelective(flow);
}
log.info("删除【{}】工作流的依赖关系", flowVo.getFlowName());
flowDependentMapper.deleteByFlowId(flowVo.getFlowId());
log.info("开始建立【{}】工作流的依赖关系,并将上游工作流设置为拥有下游工作流", flowVo.getFlowName());
//判断是否有工作流依赖
if (StringUtils.isNotEmpty(flowVo.getFlowDepend())){
String flowDepends = flowVo.getFlowDepend();
List<String> flowDependList = Arrays.asList(flowDepends);
for (String flowDependId : flowDependList){
Flow flow = flowMapper.getById(Integer.valueOf(flowDependId));
flow.setIsHaveDepend("1");
flowMapper.updateByIdSelective(flow);
flowDependentMapper.insert(flowVo.getFlowId(), Integer.valueOf(Integer.valueOf(flowDependId)));
}
}
}
@Override
@Transactional
public Flow createFlow(FlowVo flowVo) {
//TODO 未做参数校验
Flow flow = new Flow();
BeanUtils.copyProperties(flowVo, flow);
Date currentDate = new Date();
flow.setAddTime(currentDate);
flow.setAuthor("admin");
//新建的设置版本为V1
flow.setVersionName("V.1");
//新建的工作流设置为不启动
flow.setStartUp("0");
if(flow.getRemainingCount() != null && flow.getRemainingCount() > 0){
flow.setRepeatCount(flow.getRemainingCount());
}
//获取工作流id
int flowId = flowMapper.insertSelective(flow);
//初始化开始和结束节点
Node start = new Node();
Node end = new Node();
start.setFlowId(flowId);
end.setFlowId(flowId);
//如果工作流设置的节点不跟随工作流调度,设置开始和结束节点的调度为工作流的调度时间
if ("2".equals(flow.getScheduleFollow())){
start.setNodeCron(flow.getFlowCron());
end.setNodeCron(flow.getFlowCron());
}
return null;
//初始创建两个节点(开始和结束节点)
List<Node> nodeList = FlowInitUtil.InitNode(flow);
nodeList.forEach(node-> nodeMapper.insertSelective(node));
//将工作流的节点数改为2
flow.setFlowNodeCount(2);
flowMapper.updateByIdSelective(flow);
return flow;
}
@Override
......@@ -64,6 +152,19 @@ public class FlowServiceImpl implements FlowService {
@Override
public Boolean startFlow(FlowVo flowVo) {
//TODO 未做参数校验
Flow flow = new Flow();
BeanUtils.copyProperties(flowVo, flow);
//校验cron表达式
boolean isValid = CronExpression.isValidExpression(flow.getFlowCron());
if (!isValid){
log.error("工作流【{}】cron表达式不符合规范", flow.getFlowName());
return false;
}
if(flow.getRemainingCount() != null && flow.getRemainingCount() > 0){
log.info("设置工作流的剩余次数");
flow.setRepeatCount(flow.getRemainingCount());
}
return null;
}
}
//package com.byit.util;
//
//import lombok.Data;
//import org.apache.commons.lang3.StringUtils;
//
//import java.util.*;
//
//@Data
//public class DagCheck {
//
// //节点个数
// private int nodeNum;
// //连线个数
// private int lineNum;
// //节点入度
// private Map<Integer, Integer> importNum = new HashMap<>();
// //入口集合
// private Set<Integer> fromSet = new HashSet<>();
// //出口集合
// private Set<Integer> toSet = new HashSet<>();
// //节点集合
// private Set<Integer> nodeSet = new HashSet<>();
// //用队列保存拓扑序列
// private Queue<Integer> queue = new LinkedList<>();
// //存储连线信息
// private Map<String, Integer> graph = new HashMap<>();
// //工作流的key
// private String flowName;
// //依赖的节点
// private Set<String> dependencySet = new HashSet<>();
// //节点信息,节点Id和节点名称
// private Map<Integer, String> nodeInfo = new HashMap<>();
//
// public void initDag(List<FlowLine> flowLineList, List<JobFlowRelVo> jobFlowRelList, String flowName){
// //设置节点个数
// this.nodeNum = jobFlowRelList.size();
// //设置连线个数
// this.lineNum = flowLineList.size();
// //设置工作流key
// this.flowName = flowName;
//
// Map<Integer, String> dependMap = new HashMap<>();
//
// //初始化连线信息
// for (FlowLine flowLine : flowLineList){
// Integer source = flowLine.getFromJob();
// Integer target = flowLine.getToJob();
// fromSet.add(source);
// toSet.add(target);
// graph.put(flowLine.getFromJob() + "_" + flowLine.getToJob(), 1);
// String depend = dependMap.get(target);
// if(StringUtils.isEmpty(depend)){
// dependMap.put(target, String.valueOf(source));
// }else{
// dependMap.put(target, depend + "," + source);
// }
// }
//
// //设置节点信息
// for (JobFlowRelVo jobFlowRelVo: jobFlowRelList){
// nodeSet.add(jobFlowRelVo.getJobId());
// nodeInfo.put(jobFlowRelVo.getJobId(), jobFlowRelVo.getNodeName());
//
// String depends = dependMap.get(jobFlowRelVo.getJobId());
// //将没有依赖的节点放到队列中
// if(StringUtils.isEmpty(depends)){
// queue.offer(jobFlowRelVo.getJobId());
// continue;
// }
// //将依赖的节点放到集合中
// if(StringUtils.isNotEmpty(depends)){
// jobFlowRelVo.setJobDependencies(depends);
// List<String> dependenList = Arrays.asList(depends.split(","));
// dependencySet.addAll(dependenList);
// }
// }
//
// //初始化入度
// for (Integer fromJob : nodeSet){
// for (Integer toJob : nodeSet){
// Integer line = graph.get(fromJob + "_" + toJob);
// if (line != null && line == 1){
// Integer count = importNum.get(toJob) == null ? 1 : importNum.get(toJob) + 1;
// importNum.put(toJob, count);
// }
// }
// }
// }
//
// /**
// * 校验是否保存游离节点
// * @return
// */
// public DagCheckEnum checkFreeNode(){
// if(lineNum < nodeNum - 1){
// return DagCheckEnum.FREE;
// }
// for (Integer jobId : nodeSet){
// if (fromSet.contains(jobId)){
// continue;
// }
// if(toSet.contains(jobId)){
// continue;
// }
// return DagCheckEnum.FREE;
// }
//
// return DagCheckEnum.PASS;
// }
//
// /**
// * 校验是否存在环路
// * @return
// */
// public DagCheckEnum checkDag(){
// //入度为0的结点的个数,也就是入队个数
// int number = 0;
// //暂时存放拓扑序列
// Queue<Integer> temp = new LinkedList<Integer>();
// //删除这些被删除结点的出边(即对应结点入度减一)
// while(!queue.isEmpty()){
// int fromJob = queue.peek();
// temp.offer(queue.poll());
// number++;
// for(Integer toJob : nodeSet){
// String line = fromJob + "_" + toJob;
// if(graph.get(line) != null && graph.get(line) == 1){
// Integer num = importNum.get(toJob) - 1;
// importNum.put(toJob, num);
// //出现了新的入度为0的结点,删除
// if(num == 0){
// queue.offer(toJob);
// }
// }
// }
// }
// if(number != nodeNum){
// System.out.println("最后存在入度为1的结点,这个有向图是有回路的。");
// return DagCheckEnum.LOOP;
// }else{
// System.out.println("这个有向图不存在回路,拓扑序列为:" + temp.toString());
// return DagCheckEnum.PASS;
// }
// }
//
// /**
// * 校验是否单一起点
// * @return
// */
// public DagCheckEnum checkStart(){
// int startNum = queue.size();
// if(startNum == 0){
// return DagCheckEnum.LOOP;
// }else if(startNum == 1){
// if("start".equals(nodeInfo.get(queue.element()))){
// return DagCheckEnum.PASS;
// }else{
// return DagCheckEnum.STARTNAME_WRONG;
// }
// }
// return DagCheckEnum.STARTNAME_WRONG;
// }
//
// /**
// * 校验是否统一结尾
// * @return
// */
// public DagCheckEnum checkEnd(){
// List<Integer> tempList = new ArrayList<>();
// for (Integer jobId : nodeSet){
// if (dependencySet.contains(String.valueOf(jobId))){
// continue;
// }
// tempList.add(jobId);
// }
// int num = tempList.size();
// if(num == 0){
// return DagCheckEnum.LOOP;
// }else if(num == 1){
// Integer jobId = tempList.get(0);
// String nodeName = nodeInfo.get(jobId);
// if (flowName.equals(nodeName)){
// return DagCheckEnum.PASS;
// }else {
// return DagCheckEnum.ENDNAME_WRONG;
// }
// }
// return DagCheckEnum.ENDNAME_WRONG;
// }
//
//}
package com.byit.util;
import com.byit.enums.NodeTypeEnum;
import com.byit.model.Flow;
import com.byit.model.Node;
import java.util.ArrayList;
import java.util.List;
/**
* @description: 工作流初始化工具类
* @author: gml
* @create: 2019-12-25 14:01
*/
public class FlowInitUtil {
public static List<Node> InitNode(Flow flow){
List<Node> nodeList = new ArrayList<>();
//初始化开始和结束节点
Node start = new Node();
Node end = new Node();
//设置所属工作流id
start.setFlowId(flow.getFlowId());
end.setFlowId(flow.getFlowId());
//设置名称
start.setNodeName("start");
end.setNodeName("end");
//设置描述
start.setNodeDesc("开始节点");
end.setNodeDesc("结束节点");
//设置版本名称,跟随工作流的
start.setVersionName(flow.getVersionName());
end.setVersionName(flow.getVersionName());
//设置是否为虚节点
start.setIsVirtual("1");
end.setIsVirtual("1");
//如果工作流设置的节点不跟随工作流调度,设置开始和结束节点的调度为工作流的调度时间
if ("2".equals(flow.getScheduleFollow())){
start.setNodeCron(flow.getFlowCron());
end.setNodeCron(flow.getFlowCron());
}
//设置添加时间
start.setAddTime(flow.getAddTime());
end.setAddTime(flow.getAddTime());
//设置添加人
start.setAuthor(flow.getAuthor());
end.setAuthor(flow.getAuthor());
//设置类型为shell脚本
start.setJobType(NodeTypeEnum.SHELL.getCode());
end.setJobType(NodeTypeEnum.SHELL.getCode());
//设置运行内容
start.setRunSource("echo \"" + flow.getFlowName() +" run start \"");
end.setRunSource("echo \"" + flow.getFlowName() +" run end \"");
//设置负责人
start.setSourcePrincipal(flow.getPrincipal());
end.setSourcePrincipal(flow.getPrincipal());
//设置在工作流的调度上
start.setOnFork("0");
end.setOnFork("0");
nodeList.add(start);
nodeList.add(end);
return nodeList;
}
}
......@@ -6,17 +6,33 @@
<id column="flow_id" jdbcType="INTEGER" property="flowId" />
<id column="depend_flow_id" jdbcType="INTEGER" property="dependFlowId" />
</resultMap>
<select id="findByFlowId" parameterType="integer" resultMap="BaseResultMap">
select flow_id, depend_flow_id
from flow_dependent
where flow_id = #{flowId,jdbcType=INTEGER}
</select>
<delete id="deleteById" parameterType="com.byit.model.FlowDependentKey">
<!-- generated @mbg.generated date: 2019-12-25 -->
delete from flow_dependent
where flow_id = #{flowId,jdbcType=INTEGER}
and depend_flow_id = #{dependFlowId,jdbcType=INTEGER}
</delete>
<insert id="insert" parameterType="com.byit.model.FlowDependentKey">
<delete id="deleteByFlowId" parameterType="integer">
delete from flow_dependent
where flow_id = #{flowId,jdbcType=INTEGER}
</delete>
<insert id="insert" parameterType="integer">
<!-- generated @mbg.generated date: 2019-12-25 -->
insert into flow_dependent (flow_id, depend_flow_id)
values (#{flowId,jdbcType=INTEGER}, #{dependFlowId,jdbcType=INTEGER})
</insert>
<insert id="insertSelective" parameterType="com.byit.model.FlowDependentKey">
<!-- generated @mbg.generated date: 2019-12-25 -->
insert into flow_dependent
......
......@@ -47,7 +47,7 @@
where flow_id = #{flowId,jdbcType=INTEGER}
</delete>
<insert id="insertSelective" parameterType="com.byit.model.Flow">
<insert id="insertSelective" useGeneratedKeys="true" keyProperty="flowId" parameterType="com.byit.model.Flow">
<!-- generated @mbg.generated date: 2019-12-25 -->
insert into flow
<trim prefix="(" suffix=")" suffixOverrides=",">
......
......@@ -368,4 +368,13 @@
where node_id = #{nodeId,jdbcType=INTEGER}
</update>
<!--将节点更改为在调度上-->
<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>
</mapper>
\ No newline at end of file
......@@ -3,11 +3,25 @@ package com.byit.rpc.remoting.invoker.route;
import com.byit.rpc.remoting.invoker.route.impl.*;
public enum LoadBalance {
/**
* 随机
*/
RANDOM(new RpcLoadBalanceRandomStrategy()),
/**
* 轮询
*/
ROUND(new RpcLoadBalanceRoundStrategy()),
/**
* 最近最少使用
*/
LRU(new RpcLoadBalanceLRUStrategy()),
/**
* 最不常用
*/
LFU(new RpcLoadBalanceLFUStrategy()),
/**
* 哈希
*/
CONSISTENT_HASH(new RpcLoadBalanceConsistentHashStrategy());
......
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