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
9e41a972
Commit
9e41a972
authored
Feb 27, 2020
by
guominglei
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
重跑工作流和调用api工具开发
parent
8108428a
Hide whitespace changes
Inline
Side-by-side
Showing
6 changed files
with
274 additions
and
10 deletions
+274
-10
ApiFlowController.java
...h-admin/src/main/java/com/byit/api/ApiFlowController.java
+11
-4
ApiFlowService.java
...-admin/src/main/java/com/byit/service/ApiFlowService.java
+7
-1
ApiFlowServiceImpl.java
...c/main/java/com/byit/service/impl/ApiFlowServiceImpl.java
+54
-2
NodeMapper.java
...-admin-core/src/main/java/com/byit/mapper/NodeMapper.java
+8
-0
NodeMapper.xml
.../myth-admin-core/src/main/resources/mapper/NodeMapper.xml
+9
-0
JobUtils.java
...xecutor-plugin/src/main/java/com/byit/utils/JobUtils.java
+185
-3
No files found.
byit-myth-admin/src/main/java/com/byit/api/ApiFlowController.java
View file @
9e41a972
...
...
@@ -82,10 +82,17 @@ public class ApiFlowController {
return
"SUCCESS"
;
}
@PostMapping
(
"reRun"
)
@ApiOperation
(
"重跑"
)
public
String
reRun
(
String
param
){
apiFlowService
.
reRun
(
param
);
@PostMapping
(
"reRunJob"
)
@ApiOperation
(
"重跑节点"
)
public
String
reRunJob
(
String
param
){
apiFlowService
.
reRunJob
(
param
);
return
"SUCCESS"
;
}
@PostMapping
(
"reRunFlow"
)
@ApiOperation
(
"重跑工作流"
)
public
String
reRunFlow
(
String
param
){
apiFlowService
.
reRunFlow
(
param
);
return
"SUCCESS"
;
}
...
...
byit-myth-admin/src/main/java/com/byit/service/ApiFlowService.java
View file @
9e41a972
...
...
@@ -37,11 +37,17 @@ public interface ApiFlowService {
* 重跑任务
* @param param
*/
void
reRun
(
String
param
);
void
reRun
Job
(
String
param
);
/**
* 手动置为成功
* @param param
*/
void
madeSuccess
(
String
param
);
/**
* 重跑工作流
* @param param
*/
void
reRunFlow
(
String
param
);
}
byit-myth-admin/src/main/java/com/byit/service/impl/ApiFlowServiceImpl.java
View file @
9e41a972
...
...
@@ -416,7 +416,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
* @param param
*/
@Override
public
void
reRun
(
String
param
)
{
public
void
reRun
Job
(
String
param
)
{
ValidationUtil
.
dataNotBank
(
param
,
"请求参数不允许为空!"
);
JSONObject
jsonObject
=
JSON
.
parseObject
(
param
);
//获取工作空间名称
...
...
@@ -471,7 +471,7 @@ public class ApiFlowServiceImpl implements ApiFlowService {
//校验通过,开始设置重跑
//判断重跑机制(单节点重跑,节点及下游重跑
if
(!
"1"
.
equals
(
runState
))
{
//如果是只重跑当前节点
if
(!
"1"
.
equals
(
runState
))
{
//如果
不
是只重跑当前节点
//查询依赖本节点的节点,并添加到集合中
List
<
Integer
>
subNodeIdList
=
nodeDependencyMapper
.
findSubNodeList
(
node
.
getNodeId
());
if
(
subNodeIdList
!=
null
&&
subNodeIdList
.
size
()
>
0
){
...
...
@@ -483,6 +483,58 @@ public class ApiFlowServiceImpl implements ApiFlowService {
}
@Override
public
void
reRunFlow
(
String
param
){
ValidationUtil
.
dataNotBank
(
param
,
"请求参数不允许为空!"
);
JSONObject
jsonObject
=
JSON
.
parseObject
(
param
);
//获取工作空间名称
String
workspaceName
=
jsonObject
.
getString
(
"workspaceName"
);
ValidationUtil
.
dataNotBank
(
workspaceName
,
"工作空间名称不允许为空!"
);
//获取工作流名称
String
flowName
=
jsonObject
.
getString
(
"flowName"
);
ValidationUtil
.
dataNotBank
(
flowName
,
"工作流名称不允许为空!"
);
//获取要重跑的运行记录id
String
runId
=
jsonObject
.
getString
(
"runId"
);
ValidationUtil
.
dataNotBank
(
runId
,
"运行实例id不允许为空!"
);
//开始校验
Workspace
workspace
=
workspaceMapper
.
getByName
(
workspaceName
);
ValidationUtil
.
dataNotNull
(
workspace
,
workspaceName
+
"工作空间不存在"
);
Flow
flow
=
flowMapper
.
getByWorkSpaceAndName
(
workspace
.
getWorkspaceId
(),
flowName
);
ValidationUtil
.
dataNotNull
(
flow
,
flowName
+
"工作流不存在"
);
//获取运行日志实例
RunRecording
runRecording
=
runRecordingMapper
.
findRunRecordingByFlowIdAndRunId
(
flow
.
getFlowId
(),
runId
);
ValidationUtil
.
dataNotNull
(
runRecording
,
"没有找到对应的运行记录"
);
//判断是否是内嵌工作流
if
(
FlowPropertyEnum
.
IS_INNER
.
getCode
().
equals
(
flow
.
getIsInner
())){
//是内嵌工作流
Node
node
=
nodeMapper
.
getByMapFlowId
(
flow
.
getFlowId
());
Map
<
String
,
Object
>
requestMap
=
new
HashMap
<>(
10
);
requestMap
.
put
(
"workspaceName"
,
workspaceName
);
requestMap
.
put
(
"flowName"
,
flowName
);
requestMap
.
put
(
"runId"
,
runId
);
requestMap
.
put
(
"runState"
,
"1"
);
requestMap
.
put
(
"nodeName"
,
node
.
getNodeName
());
reRunJob
(
JSON
.
toJSONString
(
requestMap
));
}
else
{
//不是内嵌工作流
Long
triggerTime
=
System
.
currentTimeMillis
();
String
reRunId
=
UUID
.
randomUUID
().
toString
().
replace
(
"-"
,
""
);
List
<
JobTask
>
jobTaskList
=
new
ArrayList
<>();
List
<
Node
>
nodeList
=
nodeMapper
.
findOnforkByFlowId
(
flow
.
getFlowId
());
ValidationUtil
.
dataNotNull
(
nodeList
,
"该工作流没有在调度上的任务"
);
nodeList
.
forEach
(
node
->
{
JobTask
jobTask
=
new
JobTask
();
BeanUtils
.
copyProperties
(
node
,
jobTask
);
jobTask
.
setTriggerTime
(
triggerTime
);
jobTask
.
setTriggerStatus
(
"1"
);
jobTask
.
setRunId
(
reRunId
);
jobTask
.
setReRunId
(
runId
);
jobTaskList
.
add
(
jobTask
);
});
jobTaskMapper
.
saveJobTasks
(
jobTaskList
);
}
}
private
void
addDependNode
(
String
reRunId
,
String
runId
,
Long
triggerTime
,
List
<
JobTask
>
jobTaskList
,
List
<
Integer
>
subNodeIdList
)
{
subNodeIdList
.
forEach
(
childNodeId
->
{
Node
subNode
=
nodeMapper
.
getById
(
childNodeId
);
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/NodeMapper.java
View file @
9e41a972
...
...
@@ -78,4 +78,11 @@ public interface NodeMapper {
* @param flowId
*/
int
deleteByFlowId
(
Integer
flowId
);
/**
* 根据内嵌工作流的id来查找结点
* @param mapFlowId
* @return
*/
Node
getByMapFlowId
(
Integer
mapFlowId
);
}
\ No newline at end of file
byit-myth-core/myth-admin-core/src/main/resources/mapper/NodeMapper.xml
View file @
9e41a972
...
...
@@ -114,6 +114,15 @@
where flow_id = #{flowId,jdbcType=INTEGER}
</select>
<select
id=
"getByMapFlowId"
parameterType=
"integer"
resultMap=
"BaseResultMap"
>
select
<include
refid=
"Base_Column_List"
/>
,
<include
refid=
"Blob_Column_List"
/>
from node
where map_flow_id = #{mapFlowId,jdbcType=INTEGER}
</select>
<delete
id=
"deleteVirtualNode"
>
delete from node
where flow_id = #{flowId,jdbcType=INTEGER}
...
...
byit-myth-executor/myth-executor-plugin/src/main/java/com/byit/utils/JobUtils.java
View file @
9e41a972
...
...
@@ -37,10 +37,41 @@ public class JobUtils {
*/
private
static
final
String
REQUEST_WORKSPACE_ADD
=
"/api/workspace/add"
;
/**
* 开始工作流
*/
private
static
final
String
REQUEST_FLOW_START
=
"/api/flow/start"
;
/**
* 暂停工作流
*/
private
static
final
String
REQUEST_FLOW_STOP
=
"/api/flow/stopSchedule"
;
/**
* 删除工作流
*/
private
static
final
String
REQUEST_FLOW_DELETE
=
"/api/flow/deleteFlow"
;
/**
* 撤销工作流调度
*/
private
static
final
String
REQUEST_FLOW_REPEAL
=
"/api/flow/repealSchedule"
;
/**
* 重新开始撤销的工作流的调度
*/
private
static
final
String
REQUEST_FLOW_REREPEAL
=
"/api/flow/reStartSchedule"
;
/**
* 杀死任务
*/
private
static
final
String
REQUEST_FLOW_KILL_JOB
=
"/api/flow/killJob"
;
/**
* 杀死任务
*/
private
static
final
String
REQUEST_FLOW_KILL_FLOW
=
"/api/flow/killFlow"
;
/**
* 重跑节点
*/
private
static
final
String
REQUEST_FLOW_RERUNJOB
=
"/api/flow/reRunJob"
;
/**
* 手动置为成功
*/
private
static
final
String
REQUEST_FLOW_MAKESUCCESS
=
"/api/flow/madeSuccess"
;
private
static
final
String
SERVER_PORT
=
"8998"
;
/**
* 当前项目运行环境 jar file
...
...
@@ -96,6 +127,24 @@ public class JobUtils {
* @param workspaceName
* @return
*/
public
static
String
startFlow
(
String
flowName
,
String
workspaceName
){
//请求的路径
String
requestUrl
=
REQUEST_PREFIX
+
"127.0.0.1"
+
":"
+
SERVER_PORT
+
REQUEST_FLOW_START
;
Map
<
String
,
String
>
map
=
new
HashMap
<>(
5
);
map
.
put
(
"flowName"
,
flowName
);
map
.
put
(
"workspaceName"
,
workspaceName
);
//发送请求 添加任务
String
addRequestResult
=
HttpUtil
.
post
(
requestUrl
,
"param="
+
JSON
.
toJSONString
(
map
,
WriteClassName
));
log
.
info
(
"--------------------开始接口调用成功,结果为:{}------------------------"
,
addRequestResult
);
return
addRequestResult
;
}
/**
* 暂停工作流
* @param flowName
* @param workspaceName
* @return
*/
public
static
String
stopFlow
(
String
flowName
,
String
workspaceName
){
//请求的路径
String
requestUrl
=
REQUEST_PREFIX
+
"127.0.0.1"
+
":"
+
SERVER_PORT
+
REQUEST_FLOW_STOP
;
...
...
@@ -105,7 +154,140 @@ public class JobUtils {
//发送请求 添加任务
String
addRequestResult
=
HttpUtil
.
post
(
requestUrl
,
"param="
+
JSON
.
toJSONString
(
map
,
WriteClassName
));
//String addRequestResult = HttpUtil.post(requestUrl, JSON.toJSONString(pluginPackage, WriteClassName))
log
.
info
(
"--------------------暂停成功,结果为:{}------------------------"
,
addRequestResult
);
log
.
info
(
"--------------------暂停接口调用成功,结果为:{}------------------------"
,
addRequestResult
);
return
addRequestResult
;
}
/**
* 删除工作流
* @param flowName
* @param workspaceName
* @return
*/
public
static
String
deleteFlow
(
String
flowName
,
String
workspaceName
){
//请求的路径
String
requestUrl
=
REQUEST_PREFIX
+
"127.0.0.1"
+
":"
+
SERVER_PORT
+
REQUEST_FLOW_DELETE
;
Map
<
String
,
String
>
map
=
new
HashMap
<>(
5
);
map
.
put
(
"flowName"
,
flowName
);
map
.
put
(
"workspaceName"
,
workspaceName
);
//发送请求 添加任务
String
addRequestResult
=
HttpUtil
.
post
(
requestUrl
,
"param="
+
JSON
.
toJSONString
(
map
,
WriteClassName
));
log
.
info
(
"--------------------删除接口调用成功,结果为:{}------------------------"
,
addRequestResult
);
return
addRequestResult
;
}
/**
* 删除工作流
* @param flowName
* @param workspaceName
* @return
*/
public
static
String
repealSchedule
(
String
flowName
,
String
workspaceName
){
//请求的路径
String
requestUrl
=
REQUEST_PREFIX
+
"127.0.0.1"
+
":"
+
SERVER_PORT
+
REQUEST_FLOW_REPEAL
;
Map
<
String
,
String
>
map
=
new
HashMap
<>(
5
);
map
.
put
(
"flowName"
,
flowName
);
map
.
put
(
"workspaceName"
,
workspaceName
);
//发送请求 添加任务
String
addRequestResult
=
HttpUtil
.
post
(
requestUrl
,
"param="
+
JSON
.
toJSONString
(
map
,
WriteClassName
));
log
.
info
(
"--------------------删除接口调用成功,结果为:{}------------------------"
,
addRequestResult
);
return
addRequestResult
;
}
/**
* 重新开始某次调度
* @param runId
* @return
*/
public
static
String
startSchedule
(
String
runId
){
//请求的路径
String
requestUrl
=
REQUEST_PREFIX
+
"127.0.0.1"
+
":"
+
SERVER_PORT
+
REQUEST_FLOW_REREPEAL
;
Map
<
String
,
Object
>
map
=
new
HashMap
<>(
5
);
map
.
put
(
"runId"
,
runId
);
//发送请求 添加任务
String
addRequestResult
=
HttpUtil
.
post
(
requestUrl
,
map
);
log
.
info
(
"--------------------重新开始调度接口调用成功,结果为:{}------------------------"
,
addRequestResult
);
return
addRequestResult
;
}
/**
* 杀死任务
* @param runId
* @param flowName
* @param nodeName
* @return
*/
public
static
String
killJob
(
String
runId
,
String
flowName
,
String
nodeName
){
//请求的路径
String
requestUrl
=
REQUEST_PREFIX
+
"127.0.0.1"
+
":"
+
SERVER_PORT
+
REQUEST_FLOW_KILL_JOB
;
Map
<
String
,
String
>
map
=
new
HashMap
<>(
5
);
map
.
put
(
"runId"
,
runId
);
map
.
put
(
"flowName"
,
flowName
);
map
.
put
(
"nodeName"
,
nodeName
);
//发送请求 添加任务
String
addRequestResult
=
HttpUtil
.
post
(
requestUrl
,
"param="
+
JSON
.
toJSONString
(
map
,
WriteClassName
));
log
.
info
(
"--------------------杀死任务接口调用成功,结果为:{}------------------------"
,
addRequestResult
);
return
addRequestResult
;
}
/**
* 杀死工作流
* @param runId
* @return
*/
public
static
String
killFlow
(
String
runId
){
//请求的路径
String
requestUrl
=
REQUEST_PREFIX
+
"127.0.0.1"
+
":"
+
SERVER_PORT
+
REQUEST_FLOW_KILL_FLOW
;
Map
<
String
,
Object
>
map
=
new
HashMap
<>(
2
);
map
.
put
(
"runId"
,
runId
);
//发送请求 添加任务
String
addRequestResult
=
HttpUtil
.
post
(
requestUrl
,
map
);
log
.
info
(
"--------------------杀死工作流接口调用成功,结果为:{}------------------------"
,
addRequestResult
);
return
addRequestResult
;
}
/**
* 重跑节点
* @param runId
* @param runState
* @param workspaceName
* @param flowName
* @param nodeName
* @return
*/
public
static
String
reRunJob
(
String
runId
,
String
runState
,
String
workspaceName
,
String
flowName
,
String
nodeName
){
//请求的路径
String
requestUrl
=
REQUEST_PREFIX
+
"127.0.0.1"
+
":"
+
SERVER_PORT
+
REQUEST_FLOW_RERUNJOB
;
Map
<
String
,
String
>
map
=
new
HashMap
<>(
5
);
map
.
put
(
"runId"
,
runId
);
map
.
put
(
"runState"
,
runState
);
map
.
put
(
"workspaceName"
,
workspaceName
);
map
.
put
(
"flowName"
,
flowName
);
map
.
put
(
"nodeName"
,
nodeName
);
//发送请求 添加任务
String
addRequestResult
=
HttpUtil
.
post
(
requestUrl
,
"param="
+
JSON
.
toJSONString
(
map
,
WriteClassName
));
log
.
info
(
"--------------------重跑节点接口调用成功,结果为:{}------------------------"
,
addRequestResult
);
return
addRequestResult
;
}
/**
* 手动置为成功
* @param runId
* @param workspaceName
* @param flowName
* @param nodeName
* @return
*/
public
static
String
makeSuccess
(
String
runId
,
String
workspaceName
,
String
flowName
,
String
nodeName
){
//请求的路径
String
requestUrl
=
REQUEST_PREFIX
+
"127.0.0.1"
+
":"
+
SERVER_PORT
+
REQUEST_FLOW_MAKESUCCESS
;
Map
<
String
,
String
>
map
=
new
HashMap
<>(
5
);
map
.
put
(
"runId"
,
runId
);
map
.
put
(
"workspaceName"
,
workspaceName
);
map
.
put
(
"flowName"
,
flowName
);
map
.
put
(
"nodeName"
,
nodeName
);
//发送请求 添加任务
String
addRequestResult
=
HttpUtil
.
post
(
requestUrl
,
"param="
+
JSON
.
toJSONString
(
map
,
WriteClassName
));
log
.
info
(
"--------------------手动置为成功接口调用成功,结果为:{}------------------------"
,
addRequestResult
);
return
addRequestResult
;
}
...
...
@@ -121,7 +303,7 @@ public class JobUtils {
String
requestUrl
=
REQUEST_PREFIX
+
"127.0.0.1"
+
":"
+
SERVER_PORT
+
REQUEST_WORKSPACE_ADD
;
//发送请求 添加任务
String
addRequestResult
=
HttpUtil
.
post
(
requestUrl
,
"workspaceName="
+
workspaceName
);
log
.
info
(
"--------------------
添加任务完成,添加
结果为:{}------------------------"
,
addRequestResult
);
log
.
info
(
"--------------------
创建工作空间接口调用成功,
结果为:{}------------------------"
,
addRequestResult
);
return
addRequestResult
;
}
...
...
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