Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
B
byit-myth-job
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
liyuan
byit-myth-job
Commits
8bbeeed6
Commit
8bbeeed6
authored
Jan 03, 2020
by
huangfusuper
Browse files
Options
Browse Files
Download
Plain Diff
Merge remote-tracking branch 'origin/developer' into developer
parents
2cdc2029
c602e122
Hide whitespace changes
Inline
Side-by-side
Showing
21 changed files
with
564 additions
and
288 deletions
+564
-288
ApiFlowController.java
...h-admin/src/main/java/com/byit/api/ApiFlowController.java
+2
-2
ApiFlowService.java
...-admin/src/main/java/com/byit/service/ApiFlowService.java
+2
-2
ApiFlowServiceImpl.java
...c/main/java/com/byit/service/impl/ApiFlowServiceImpl.java
+288
-24
FlowVersionMapper.java
...core/src/main/java/com/byit/mapper/FlowVersionMapper.java
+7
-0
JobTaskMapper.java
...min-core/src/main/java/com/byit/mapper/JobTaskMapper.java
+7
-0
NodeDependencyMapper.java
...e/src/main/java/com/byit/mapper/NodeDependencyMapper.java
+4
-1
NodeMapper.java
...-admin-core/src/main/java/com/byit/mapper/NodeMapper.java
+10
-0
NodeVersionDependencyMapper.java
...ain/java/com/byit/mapper/NodeVersionDependencyMapper.java
+13
-0
NodeVersionMapper.java
...core/src/main/java/com/byit/mapper/NodeVersionMapper.java
+7
-0
RunRecordingMapper.java
...ore/src/main/java/com/byit/mapper/RunRecordingMapper.java
+19
-0
FlowVersion.java
...-admin-core/src/main/java/com/byit/model/FlowVersion.java
+4
-3
NodeVersionDependencyKey.java
...rc/main/java/com/byit/model/NodeVersionDependencyKey.java
+31
-0
FlowServiceImpl.java
.../src/main/java/com/byit/service/impl/FlowServiceImpl.java
+1
-1
Generator-config.xml
...e/myth-admin-core/src/main/resources/Generator-config.xml
+1
-13
FlowVersionMapper.xml
...dmin-core/src/main/resources/mapper/FlowVersionMapper.xml
+14
-42
JobTaskMapper.xml
...th-admin-core/src/main/resources/mapper/JobTaskMapper.xml
+6
-0
NodeDependencyMapper.xml
...n-core/src/main/resources/mapper/NodeDependencyMapper.xml
+20
-1
NodeMapper.xml
.../myth-admin-core/src/main/resources/mapper/NodeMapper.xml
+51
-97
NodeVersionDependencyMapper.xml
...src/main/resources/mapper/NodeVersionDependencyMapper.xml
+41
-0
NodeVersionMapper.xml
...dmin-core/src/main/resources/mapper/NodeVersionMapper.xml
+13
-102
RunRecordingMapper.xml
...min-core/src/main/resources/mapper/RunRecordingMapper.xml
+23
-0
No files found.
byit-myth-admin/src/main/java/com/byit/api/ApiFlowController.java
View file @
8bbeeed6
...
...
@@ -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"
;
}
...
...
@@ -34,7 +34,7 @@ public class ApiFlowController {
@PostMapping
(
"delete"
)
@ApiOperation
(
"删除工作流,若当前工作流被依赖则删除失败,只允许删除不被依赖的工作流,当前工作流依赖其他工作流不影响"
)
public
String
deleteFlow
(
String
flowName
,
String
workspaceName
){
apiFlowService
.
deleteFlow
(
flowName
);
apiFlowService
.
deleteFlow
(
flowName
,
workspaceName
);
return
"SUCCESS"
;
}
...
...
byit-myth-admin/src/main/java/com/byit/service/ApiFlowService.java
View file @
8bbeeed6
...
...
@@ -10,7 +10,7 @@ import com.byit.job.dto.plugin.PluginPackage;
*/
public
interface
ApiFlowService
{
void
deleteFlow
(
String
flowName
);
void
deleteFlow
(
String
flowName
,
String
workspaceName
);
void
start
(
String
flowName
,
String
workspaceName
);
...
...
@@ -22,7 +22,7 @@ public interface ApiFlowService {
void
reStartSchedule
(
String
runId
);
void
publishFlow
(
PluginPackage
pluginPackage
);
void
publishFlow
(
PluginPackage
pluginPackage
)
throws
Exception
;
/**
* 判断工作流石佛存在
...
...
byit-myth-admin/src/main/java/com/byit/service/impl/ApiFlowServiceImpl.java
View file @
8bbeeed6
...
...
@@ -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,16 +37,42 @@ 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
;
@Resource
private
RunRecordingMapper
runRecordingMapper
;
@Resource
private
JobTaskMapper
jobTaskMapper
;
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
void
publishFlow
(
PluginPackage
pluginPackage
)
{
public
void
publishFlow
(
PluginPackage
pluginPackage
)
throws
Exception
{
//校验参数
log
.
info
(
"校验参数是否符合规范"
);
Workspace
workspace
=
validate
(
pluginPackage
);
saveFlow
(
pluginPackage
.
getFlow
(),
workspace
.
getWorkspaceId
(),
false
);
//保存工作流
log
.
info
(
"保存工作流和节点信息"
);
Flow
flow
=
saveFlow
(
pluginPackage
.
getFlow
(),
workspace
.
getWorkspaceId
(),
false
);
//生成工作流版本
log
.
info
(
"生成版本"
);
FlowVersion
flowVersion
=
saveFlowVersion
(
flow
);
log
.
info
(
"插件端通过API保存成功!"
);
}
/**
...
...
@@ -56,7 +81,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 +91,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
());
flow
.
setRepeatCount
(
pluginFlow
.
getConfig
().
getRepeatCount
());
flow
.
setRemainingCount
(
pluginFlow
.
getConfig
().
getRepeatCount
());
//设置可执行次数
//若没设置或者设置为-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 +119,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 +325,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 +360,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,9 +380,87 @@ public class ApiFlowServiceImpl implements ApiFlowService {
}
}
@Override
public
void
deleteFlow
(
String
flowName
)
{
/**
* 工作流生成新版本
* @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
,
String
workspaceName
)
{
log
.
info
(
"删除工作空间【{}】---工作流【{}】调度"
,
workspaceName
,
flowName
);
Workspace
workspace
=
workspaceMapper
.
getByName
(
workspaceName
);
ValidationUtil
.
dataNotNull
(
workspace
,
workspaceName
+
"工作空间不存在"
);
Flow
flow
=
flowMapper
.
getByWorkSpaceAndName
(
workspace
.
getWorkspaceId
(),
flowName
);
ValidationUtil
.
dataNotNull
(
flow
,
flowName
+
"工作流不存在"
);
ValidationUtil
.
isTrueValidation
(
FlowPropertyEnum
.
IS_INNER
.
getCode
().
equals
(
flow
.
getIsInner
()),
"内嵌工作流不允许删除!"
);
//判断是否在调度中
List
<
RunRecording
>
recordingList
=
runRecordingMapper
.
findOnScheduleByFlowId
(
flow
.
getFlowId
());
ValidationUtil
.
isTrueValidation
(
null
!=
recordingList
&&
recordingList
.
size
()
>
0
,
"工作流已在调度中不允许撤销调度!"
);
List
<
RunRecording
>
unStartRecordingList
=
runRecordingMapper
.
findUnStartByFlowId
(
flow
.
getFlowId
());
//删除对应的task记录
unStartRecordingList
.
forEach
(
runRecording
->
{
ValidationUtil
.
isTrueValidation
((
runRecording
.
getTriggerTime
()
-
System
.
currentTimeMillis
())
>
10000
,
"工作流已在调度中不允许删除"
);
runRecordingMapper
.
deleteByRunId
(
runRecording
.
getRunId
());
jobTaskMapper
.
deleteByRunId
(
runRecording
.
getRunId
());
});
//删除工作流及下属节点
}
@Override
...
...
@@ -225,7 +470,26 @@ public class ApiFlowServiceImpl implements ApiFlowService {
@Override
public
void
repealSchedule
(
String
flowName
,
String
workspaceName
)
{
log
.
info
(
"撤销工作空间【{}】---工作流【{}】调度"
,
workspaceName
,
flowName
);
Workspace
workspace
=
workspaceMapper
.
getByName
(
workspaceName
);
ValidationUtil
.
dataNotNull
(
workspace
,
workspaceName
+
"工作空间不存在"
);
Flow
flow
=
flowMapper
.
getByWorkSpaceAndName
(
workspace
.
getWorkspaceId
(),
flowName
);
ValidationUtil
.
dataNotNull
(
flow
,
flowName
+
"工作流不存在"
);
ValidationUtil
.
isTrueValidation
(
FlowPropertyEnum
.
IS_INNER
.
getCode
().
equals
(
flow
.
getIsInner
()),
"内嵌工作流不允许撤销调度!"
);
//判断是否在调度中
List
<
RunRecording
>
recordingList
=
runRecordingMapper
.
findOnScheduleByFlowId
(
flow
.
getFlowId
());
ValidationUtil
.
isTrueValidation
(
null
!=
recordingList
&&
recordingList
.
size
()
>
0
,
"工作流已在调度中不允许撤销调度!"
);
List
<
RunRecording
>
unStartRecordingList
=
runRecordingMapper
.
findUnStartByFlowId
(
flow
.
getFlowId
());
//删除对应的task记录
unStartRecordingList
.
forEach
(
runRecording
->
{
ValidationUtil
.
isTrueValidation
((
runRecording
.
getTriggerTime
()
-
System
.
currentTimeMillis
())
>
10000
,
"工作流已在调度中不允许撤销调度"
);
runRecordingMapper
.
deleteByRunId
(
runRecording
.
getRunId
());
jobTaskMapper
.
deleteByRunId
(
runRecording
.
getRunId
());
});
flow
.
setStartUp
(
FlowPropertyEnum
.
NO_START
.
getCode
());
flowMapper
.
updateByIdSelective
(
flow
);
log
.
info
(
"工作流【{}】调度撤销成功!"
,
flowName
);
}
@Override
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/FlowVersionMapper.java
View file @
8bbeeed6
...
...
@@ -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
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/JobTaskMapper.java
View file @
8bbeeed6
...
...
@@ -51,4 +51,10 @@ public interface JobTaskMapper {
* @param jobTasks
*/
void
deleteInId
(
@Param
(
"jobTasks"
)
List
<
JobTask
>
jobTasks
);
/**
* 根据runid删除运行信息
* @param runId
*/
int
deleteByRunId
(
String
runId
);
}
\ No newline at end of file
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/NodeDependencyMapper.java
View file @
8bbeeed6
...
...
@@ -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
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/NodeMapper.java
View file @
8bbeeed6
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
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/NodeVersionDependencyMapper.java
0 → 100644
View file @
8bbeeed6
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
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/NodeVersionMapper.java
View file @
8bbeeed6
...
...
@@ -11,4 +11,10 @@ public interface NodeVersionMapper {
int
updateByIdSelective
(
NodeVersion
record
);
/**
* 将当前工作流下版本表中所有节点都设置为不是当前版本
* @param flowId
* @return
*/
int
updateNotCurrentVersion
(
Integer
flowId
);
}
\ No newline at end of file
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/RunRecordingMapper.java
View file @
8bbeeed6
...
...
@@ -57,4 +57,22 @@ public interface RunRecordingMapper {
*/
int
deleteById
(
Integer
recordingId
);
/**
* 根据工作流id查询正在调度中的工作流
* @param flowId
*/
List
<
RunRecording
>
findOnScheduleByFlowId
(
Integer
flowId
);
/**
* 根据工作流id查询未开始的工作流
* @param flowId
* @return
*/
List
<
RunRecording
>
findUnStartByFlowId
(
Integer
flowId
);
/**
* 根据runid删除工作流的运行记录
* @param runId
*/
int
deleteByRunId
(
String
runId
);
}
\ No newline at end of file
byit-myth-core/myth-admin-core/src/main/java/com/byit/model/FlowVersion.java
View file @
8bbeeed6
...
...
@@ -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
;
/**
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/model/NodeVersionDependencyKey.java
0 → 100644
View file @
8bbeeed6
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
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/impl/FlowServiceImpl.java
View file @
8bbeeed6
...
...
@@ -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
.
find
DependId
ByNodeId
(
node
.
getNodeId
());
nodeVoList
.
add
(
nodeVo
);
});
flowVo
.
setNodeVoList
(
nodeVoList
);
...
...
byit-myth-core/myth-admin-core/src/main/resources/Generator-config.xml
View file @
8bbeeed6
...
...
@@ -55,19 +55,7 @@
type=
"XMLMAPPER"
>
<property
name=
"enableSubPackages"
value=
"false"
/>
</javaClientGenerator>
<table
tableName=
"email_alarm"
domainObjectName=
"EmailAlarm"
/>
<table
tableName=
"flow"
domainObjectName=
"Flow"
/>
<table
tableName=
"flow_version"
domainObjectName=
"FlowVersion"
/>
<table
tableName=
"job_task"
domainObjectName=
"JobTask"
/>
<table
tableName=
"job_task_run_log"
domainObjectName=
"JobTaskRunLog"
/>
<table
tableName=
"job_task_schedule"
domainObjectName=
"JobTaskSchedule"
/>
<table
tableName=
"node"
domainObjectName=
"Node"
/>
<table
tableName=
"node_version"
domainObjectName=
"NodeVersion"
/>
<table
tableName=
"run_recording"
domainObjectName=
"RunRecording"
/>
<table
tableName=
"source_history"
domainObjectName=
"SourceHistory"
/>
<table
tableName=
"flow_dependent"
domainObjectName=
"FlowDependent"
/>
<table
tableName=
"node_dependency"
domainObjectName=
"NodeDependency"
/>
<table
tableName=
"workspace"
domainObjectName=
"Workspace"
/>
<table
tableName=
"node_version_dependency"
domainObjectName=
"NodeVersionDependency"
/>
</context>
...
...
byit-myth-core/myth-admin-core/src/main/resources/mapper/FlowVersionMapper.xml
View file @
8bbeeed6
...
...
@@ -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
byit-myth-core/myth-admin-core/src/main/resources/mapper/JobTaskMapper.xml
View file @
8bbeeed6
...
...
@@ -344,4 +344,9 @@
delete from job_task
where id = #{id,jdbcType=INTEGER}
</delete>
<delete
id=
"deleteByRunId"
parameterType=
"java.lang.Integer"
>
delete from job_task
where run_id = #{runId,jdbcType=INTEGER}
</delete>
</mapper>
\ No newline at end of file
byit-myth-core/myth-admin-core/src/main/resources/mapper/NodeDependencyMapper.xml
View file @
8bbeeed6
...
...
@@ -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})
...
...
byit-myth-core/myth-admin-core/src/main/resources/mapper/NodeMapper.xml
View file @
8bbeeed6
...
...
@@ -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
byit-myth-core/myth-admin-core/src/main/resources/mapper/NodeVersionDependencyMapper.xml
0 → 100644
View file @
8bbeeed6
<?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
byit-myth-core/myth-admin-core/src/main/resources/mapper/NodeVersionMapper.xml
View file @
8bbeeed6
...
...
@@ -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
byit-myth-core/myth-admin-core/src/main/resources/mapper/RunRecordingMapper.xml
View file @
8bbeeed6
...
...
@@ -49,11 +49,34 @@
from run_recording
where recording_id = #{recordingId,jdbcType=INTEGER}
</select>
<select
id=
"findOnScheduleByFlowId"
parameterType=
"java.lang.Integer"
resultMap=
"BaseResultMap"
>
select
<include
refid=
"Base_Column_List"
/>
from run_recording
where flow_id = #{flowId,jdbcType=INTEGER}
and flow_status not in ('1','4')
</select>
<select
id=
"findUnStartByFlowId"
parameterType=
"java.lang.Integer"
resultMap=
"BaseResultMap"
>
select
<include
refid=
"Base_Column_List"
/>
from run_recording
where flow_id = #{flowId,jdbcType=INTEGER}
and flow_status = '1'
</select>
<delete
id=
"deleteById"
parameterType=
"java.lang.Integer"
>
<!-- generated @mbg.generated date: 2019-12-25 -->
delete from run_recording
where recording_id = #{recordingId,jdbcType=INTEGER}
</delete>
<delete
id=
"deleteByRunId"
parameterType=
"java.lang.Integer"
>
delete from run_recording
where run_id = #{runId,jdbcType=INTEGER}
</delete>
<insert
id=
"saveRunRecording"
parameterType=
"com.byit.model.RunRecording"
>
<!-- generated @mbg.generated date: 2019-12-25 -->
insert into run_recording
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment