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
db9669db
Commit
db9669db
authored
May 13, 2020
by
huangfusuper
Browse files
Options
Browse Files
Download
Plain Diff
Merge remote-tracking branch 'origin/developer' into developer
parents
5e574d01
185bea57
Hide whitespace changes
Inline
Side-by-side
Showing
4 changed files
with
58 additions
and
65 deletions
+58
-65
ApiFlowServiceImpl.java
...c/main/java/com/byit/service/impl/ApiFlowServiceImpl.java
+57
-57
RunRecordingMapper.java
...ore/src/main/java/com/byit/mapper/RunRecordingMapper.java
+0
-2
RunRecordingAndJobTaskServiceImpl.java
...ce/mapservice/impl/RunRecordingAndJobTaskServiceImpl.java
+1
-1
RunRecordingMapper.xml
...min-core/src/main/resources/mapper/RunRecordingMapper.xml
+0
-5
No files found.
byit-myth-admin/src/main/java/com/byit/service/impl/ApiFlowServiceImpl.java
View file @
db9669db
...
@@ -482,25 +482,22 @@ public class ApiFlowServiceImpl implements ApiFlowService {
...
@@ -482,25 +482,22 @@ public class ApiFlowServiceImpl implements ApiFlowService {
//获取当前时间
//获取当前时间
String
reRunId
=
UUID
.
randomUUID
().
toString
().
replace
(
"-"
,
""
);
String
reRunId
=
UUID
.
randomUUID
().
toString
().
replace
(
"-"
,
""
);
Long
triggerTime
=
System
.
currentTimeMillis
();
Long
triggerTime
=
System
.
currentTimeMillis
();
JobTask
jobTask
=
new
JobTask
();
WaitingTask
waitingTask
=
new
WaitingTask
();
BeanUtils
.
copyProperties
(
jobTaskRunLog
,
jobTask
);
BeanUtils
.
copyProperties
(
jobTaskRunLog
,
waitingTask
);
jobTask
.
setTriggerTime
(
triggerTime
);
waitingTask
.
setTriggerTime
(
triggerTime
);
jobTask
.
setTriggerStatus
(
"1"
);
waitingTask
.
setReRunId
(
jobTaskRunLog
.
getRunId
());
jobTask
.
setRunId
(
reRunId
);
jobTask
.
setReRunId
(
jobTaskRunLog
.
getRunId
());
//设置为重跑
//设置为重跑
jobTask
.
setScheduleType
(
ScheduleTypeEnum
.
REPEAT
.
getCode
());
waitingTask
.
setScheduleType
(
ScheduleTypeEnum
.
REPEAT
.
getCode
());
jobTask
.
setOperator
(
userName
);
waitingTask
.
setOperator
(
userName
);
jobTask
.
setLogId
(
null
);
//查询当前节点的依赖节点
//查询当前节点的依赖节点
List
<
Integer
>
dependNodeIdList
=
nodeDependencyMapper
.
findDependIdByNodeId
(
node
.
getNodeId
());
List
<
Integer
>
dependNodeIdList
=
nodeDependencyMapper
.
findDependIdByNodeId
(
node
.
getNodeId
());
if
(
dependNodeIdList
!=
null
&&
dependNodeIdList
.
size
()
>
0
){
if
(
dependNodeIdList
!=
null
&&
dependNodeIdList
.
size
()
>
0
){
job
Task
.
setNodeDepend
(
Joiner
.
on
(
","
).
join
(
dependNodeIdList
));
waiting
Task
.
setNodeDepend
(
Joiner
.
on
(
","
).
join
(
dependNodeIdList
));
}
}
List
<
JobTask
>
job
TaskList
=
new
ArrayList
<>();
List
<
WaitingTask
>
waiting
TaskList
=
new
ArrayList
<>();
jobTaskList
.
add
(
job
Task
);
waitingTaskList
.
add
(
waiting
Task
);
//校验通过,开始设置重跑
//校验通过,开始设置重跑
//判断重跑机制(单节点重跑,节点及下游重跑
//判断重跑机制(单节点重跑,节点及下游重跑
...
@@ -508,17 +505,26 @@ public class ApiFlowServiceImpl implements ApiFlowService {
...
@@ -508,17 +505,26 @@ public class ApiFlowServiceImpl implements ApiFlowService {
//查询依赖本节点的节点,并添加到集合中
//查询依赖本节点的节点,并添加到集合中
List
<
Integer
>
subNodeIdList
=
nodeDependencyMapper
.
findSubNodeList
(
node
.
getNodeId
());
List
<
Integer
>
subNodeIdList
=
nodeDependencyMapper
.
findSubNodeList
(
node
.
getNodeId
());
if
(
subNodeIdList
!=
null
&&
subNodeIdList
.
size
()
>
0
){
if
(
subNodeIdList
!=
null
&&
subNodeIdList
.
size
()
>
0
){
addDependNode
(
r
eRunId
,
runInfo
.
getRunId
(),
triggerTime
,
job
TaskList
,
subNodeIdList
,
userName
);
addDependNode
(
r
unInfo
.
getRunId
(),
triggerTime
,
waiting
TaskList
,
subNodeIdList
,
userName
);
}
}
}
}
WaitingRecord
waitingRecord
=
new
WaitingRecord
();
RunRecording
newRunRecording
=
buildRunRecording
(
runRecording
,
triggerTime
,
reRunId
,
flow
,
userName
);
RunRecording
newRunRecording
=
buildRunRecording
(
runRecording
,
triggerTime
,
reRunId
,
flow
,
userName
);
runRecording
.
setFlowStatus
(
RunRecordingEnum
.
FLOW_STATUS_RUN_ING
.
getCode
());
BeanUtils
.
copyProperties
(
newRunRecording
,
waitingRecord
);
runRecording
.
setFlowRunResult
(
null
);
Integer
order
=
waitingRecordMapper
.
findOrderByFlowId
(
flow
.
getFlowId
());
runRecordingMapper
.
updateState
(
runRecording
);
if
(
order
==
null
){
order
=
0
;
}
waitingRecord
.
setWaitOrder
(++
order
);
Integer
waitId
=
waitingRecordMapper
.
insertSelective
(
waitingRecord
);
//生成新的工作流实例
//生成新的工作流实例
runRecordingMapper
.
saveRunRecording
(
newRunRecording
);
runRecordingMapper
.
saveRunRecording
(
newRunRecording
);
jobTaskMapper
.
saveJobTasks
(
jobTaskList
);
waitingTaskList
.
forEach
(
waiting
->
{
waiting
.
setWaitId
(
waitingRecord
.
getWaitId
());
waitingTaskMapper
.
insertSelective
(
waiting
);
});
}
}
private
RunRecording
buildRunRecording
(
RunRecording
runRecording
,
Long
triggerTime
,
String
reRunId
,
Flow
flow
,
String
userName
)
{
private
RunRecording
buildRunRecording
(
RunRecording
runRecording
,
Long
triggerTime
,
String
reRunId
,
Flow
flow
,
String
userName
)
{
...
@@ -567,72 +573,66 @@ public class ApiFlowServiceImpl implements ApiFlowService {
...
@@ -567,72 +573,66 @@ public class ApiFlowServiceImpl implements ApiFlowService {
Long
triggerTime
=
System
.
currentTimeMillis
();
Long
triggerTime
=
System
.
currentTimeMillis
();
String
reRunId
=
UUID
.
randomUUID
().
toString
().
replace
(
"-"
,
""
);
String
reRunId
=
UUID
.
randomUUID
().
toString
().
replace
(
"-"
,
""
);
List
<
JobTask
>
jobTaskList
=
new
ArrayList
<>();
List
<
JobTaskRunLogWithBLOBs
>
jobTaskRunLogList
=
jobTaskRunLogMapper
.
findJobTaskRunLogWithBLOBsByFlowIdAndRunId
(
flow
.
getFlowId
(),
runInfo
.
getRunId
());
List
<
JobTaskRunLogWithBLOBs
>
jobTaskRunLogList
=
jobTaskRunLogMapper
.
findJobTaskRunLogWithBLOBsByFlowIdAndRunId
(
flow
.
getFlowId
(),
runInfo
.
getRunId
());
ValidationUtil
.
dataNotNull
(
jobTaskRunLogList
,
"该工作流没有在调度上的任务"
);
ValidationUtil
.
dataNotNull
(
jobTaskRunLogList
,
"该工作流没有在调度上的任务"
);
String
userName
=
currentUserUtils
.
account
();
String
userName
=
currentUserUtils
.
account
();
WaitingRecord
waitingRecord
=
new
WaitingRecord
();
//生成新的实例
RunRecording
newRunRecording
=
buildRunRecording
(
runRecording
,
triggerTime
,
reRunId
,
flow
,
userName
);
BeanUtils
.
copyProperties
(
newRunRecording
,
waitingRecord
);
Integer
order
=
waitingRecordMapper
.
findOrderByFlowId
(
flow
.
getFlowId
());
if
(
order
==
null
){
order
=
0
;
}
waitingRecord
.
setWaitOrder
(++
order
);
Integer
waitId
=
waitingRecordMapper
.
insertSelective
(
waitingRecord
);
//生成新的工作流实例
runRecordingMapper
.
saveRunRecording
(
newRunRecording
);
jobTaskRunLogList
.
forEach
(
jobTaskRunLog
->
{
jobTaskRunLogList
.
forEach
(
jobTaskRunLog
->
{
JobTask
jobTask
=
new
JobTask
();
WaitingTask
waitingTask
=
new
WaitingTask
();
BeanUtils
.
copyProperties
(
jobTaskRunLog
,
jobTask
);
BeanUtils
.
copyProperties
(
jobTaskRunLog
,
waitingTask
);
jobTask
.
setTriggerTime
(
triggerTime
);
waitingTask
.
setTriggerTime
(
triggerTime
);
jobTask
.
setTriggerStatus
(
"1"
);
waitingTask
.
setReRunId
(
jobTaskRunLog
.
getRunId
());
jobTask
.
setRunId
(
reRunId
);
jobTask
.
setReRunId
(
runInfo
.
getRunId
());
//设置为重跑
//设置为重跑
job
Task
.
setScheduleType
(
ScheduleTypeEnum
.
REPEAT
.
getCode
());
waiting
Task
.
setScheduleType
(
ScheduleTypeEnum
.
REPEAT
.
getCode
());
job
Task
.
setOperator
(
userName
);
waiting
Task
.
setOperator
(
userName
);
jobTask
.
setLogId
(
null
);
waitingTask
.
setReRunId
(
jobTaskRunLog
.
getRunId
()
);
//查询当前节点的依赖节点
//查询当前节点的依赖节点
List
<
Integer
>
dependNodeIdList
=
nodeDependencyMapper
.
findDependIdByNodeId
(
jobTaskRunLog
.
getNodeId
());
List
<
Integer
>
dependNodeIdList
=
nodeDependencyMapper
.
findDependIdByNodeId
(
jobTaskRunLog
.
getNodeId
());
if
(
dependNodeIdList
!=
null
&&
dependNodeIdList
.
size
()
>
0
){
if
(
dependNodeIdList
!=
null
&&
dependNodeIdList
.
size
()
>
0
){
job
Task
.
setNodeDepend
(
Joiner
.
on
(
","
).
join
(
dependNodeIdList
));
waiting
Task
.
setNodeDepend
(
Joiner
.
on
(
","
).
join
(
dependNodeIdList
));
}
}
jobTaskList
.
add
(
jobTask
);
waitingTask
.
setWaitId
(
waitingRecord
.
getWaitId
());
waitingTaskMapper
.
insertSelective
(
waitingTask
);
});
});
//生成运行实例
RunRecording
newRunRecording
=
buildRunRecording
(
runRecording
,
triggerTime
,
reRunId
,
flow
,
userName
);
runRecording
.
setFlowStatus
(
RunRecordingEnum
.
FLOW_STATUS_RUN_ING
.
getCode
());
runRecording
.
setFlowRunResult
(
null
);
runRecordingMapper
.
updateState
(
runRecording
);
//将工作流下的节点改为运行中
jobTaskRunLogMapper
.
updateByFlowIdAndRunId
(
flow
.
getFlowId
(),
runInfo
.
getRunId
());
//生成新的工作流实例
runRecordingMapper
.
saveRunRecording
(
newRunRecording
);
jobTaskMapper
.
saveJobTasks
(
jobTaskList
);
}
}
}
}
private
void
addDependNode
(
String
r
eRunId
,
String
runId
,
Long
triggerTime
,
List
<
JobTask
>
job
TaskList
,
List
<
Integer
>
subNodeIdList
,
String
userName
)
{
private
void
addDependNode
(
String
r
unId
,
Long
triggerTime
,
List
<
WaitingTask
>
waiting
TaskList
,
List
<
Integer
>
subNodeIdList
,
String
userName
)
{
subNodeIdList
.
forEach
(
childNodeId
->
{
subNodeIdList
.
forEach
(
childNodeId
->
{
JobTaskRunLog
jobTaskRunLog
=
jobTaskRunLogMapper
.
findByRunIdAndNodeId
(
runId
,
childNodeId
);
JobTaskRunLog
jobTaskRunLog
=
jobTaskRunLogMapper
.
findByRunIdAndNodeId
(
runId
,
childNodeId
);
JobTask
jobTask
=
new
JobTask
();
WaitingTask
waitingTask
=
new
WaitingTask
();
BeanUtils
.
copyProperties
(
jobTaskRunLog
,
jobTask
);
BeanUtils
.
copyProperties
(
jobTaskRunLog
,
waitingTask
);
jobTask
.
setTriggerTime
(
triggerTime
);
waitingTask
.
setTriggerTime
(
triggerTime
);
jobTask
.
setTriggerStatus
(
"1"
);
waitingTask
.
setReRunId
(
runId
);
jobTask
.
setRunId
(
reRunId
);
jobTask
.
setReRunId
(
runId
);
//设置为重跑
//设置为重跑
jobTask
.
setScheduleType
(
ScheduleTypeEnum
.
REPEAT
.
getCode
());
waitingTask
.
setScheduleType
(
ScheduleTypeEnum
.
REPEAT
.
getCode
());
jobTask
.
setOperator
(
userName
);
waitingTask
.
setOperator
(
userName
);
jobTask
.
setLogId
(
null
);
//查询当前节点的依赖节点
//查询当前节点的依赖节点
List
<
Integer
>
dependNodeIdList
=
nodeDependencyMapper
.
findDependIdByNodeId
(
jobTaskRunLog
.
getNodeId
());
List
<
Integer
>
dependNodeIdList
=
nodeDependencyMapper
.
findDependIdByNodeId
(
jobTaskRunLog
.
getNodeId
());
if
(
dependNodeIdList
!=
null
&&
dependNodeIdList
.
size
()
>
0
){
if
(
dependNodeIdList
!=
null
&&
dependNodeIdList
.
size
()
>
0
){
job
Task
.
setNodeDepend
(
Joiner
.
on
(
","
).
join
(
dependNodeIdList
));
waiting
Task
.
setNodeDepend
(
Joiner
.
on
(
","
).
join
(
dependNodeIdList
));
}
}
//将节点对应的日志改为运行中
waitingTaskList
.
add
(
waitingTask
);
jobTaskRunLog
.
setRunCode
(
NodeRunStatusPropertyEnum
.
RUN_ING
.
getCode
());
jobTaskRunLogMapper
.
updateJobTaskRunLog
(
jobTaskRunLog
);
jobTaskList
.
add
(
jobTask
);
//查询依赖于当前节点的下级节点
//查询依赖于当前节点的下级节点
List
<
Integer
>
childNodeIdList
=
nodeDependencyMapper
.
findSubNodeList
(
childNodeId
);
List
<
Integer
>
childNodeIdList
=
nodeDependencyMapper
.
findSubNodeList
(
childNodeId
);
if
(
childNodeIdList
!=
null
&&
childNodeIdList
.
size
()
>
0
){
if
(
childNodeIdList
!=
null
&&
childNodeIdList
.
size
()
>
0
){
addDependNode
(
r
eRunId
,
runId
,
triggerTime
,
job
TaskList
,
childNodeIdList
,
userName
);
addDependNode
(
r
unId
,
triggerTime
,
waiting
TaskList
,
childNodeIdList
,
userName
);
}
}
});
});
}
}
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/RunRecordingMapper.java
View file @
db9669db
...
@@ -206,8 +206,6 @@ public interface RunRecordingMapper {
...
@@ -206,8 +206,6 @@ public interface RunRecordingMapper {
*/
*/
RunRecording
findNewStatus
(
Integer
flowId
);
RunRecording
findNewStatus
(
Integer
flowId
);
int
updateState
(
RunRecording
runRecording
);
Integer
findFlowNum
(
@Param
(
"startTime"
)
Date
startTime
,
Integer
findFlowNum
(
@Param
(
"startTime"
)
Date
startTime
,
@Param
(
"endTime"
)
Date
endTime
,
@Param
(
"endTime"
)
Date
endTime
,
@Param
(
"flowIdList"
)
List
<
Integer
>
flowIdList
);
@Param
(
"flowIdList"
)
List
<
Integer
>
flowIdList
);
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/mapservice/impl/RunRecordingAndJobTaskServiceImpl.java
View file @
db9669db
...
@@ -241,7 +241,7 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
...
@@ -241,7 +241,7 @@ public class RunRecordingAndJobTaskServiceImpl implements RunRecordingAndJobTask
target
.
setTriggerStatus
(
"1"
);
target
.
setTriggerStatus
(
"1"
);
target
.
setRunId
(
waitingRecord
.
getRunId
());
target
.
setRunId
(
waitingRecord
.
getRunId
());
BeanUtils
.
copyProperties
(
waitingTask
,
target
);
BeanUtils
.
copyProperties
(
waitingTask
,
target
);
target
.
setReRunId
(
runRecording
.
getReRunId
());
target
.
setReRunId
(
waitingTask
.
getReRunId
());
target
.
setId
(
null
);
target
.
setId
(
null
);
target
.
setTriggerTime
(
0L
);
target
.
setTriggerTime
(
0L
);
return
target
;
return
target
;
...
...
byit-myth-core/myth-admin-core/src/main/resources/mapper/RunRecordingMapper.xml
View file @
db9669db
...
@@ -606,9 +606,5 @@
...
@@ -606,9 +606,5 @@
where run_id = #{runId,jdbcType=VARCHAR}
where run_id = #{runId,jdbcType=VARCHAR}
and flow_name = #{flowName,jdbcType=VARCHAR}
and flow_name = #{flowName,jdbcType=VARCHAR}
</update>
</update>
<update
id=
"updateState"
parameterType=
"com.byit.model.RunRecording"
>
update run_recording set flow_status = #{flowStatus},flow_run_result = #{flowRunResult}
where recording_id = #{recordingId,jdbcType=INTEGER}
</update>
</mapper>
</mapper>
\ No newline at end of file
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