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
818a7e1e
Commit
818a7e1e
authored
Apr 07, 2020
by
huangfusuper
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
对重跑节点的单独处理
parent
ad58eced
Hide whitespace changes
Inline
Side-by-side
Showing
3 changed files
with
112 additions
and
5 deletions
+112
-5
NodeNameEnum.java
...admin-core/src/main/java/com/byit/enums/NodeNameEnum.java
+1
-0
SpecialFlowThreadRunHelper.java
...va/com/byit/thread/helper/SpecialFlowThreadRunHelper.java
+102
-0
TaskThreadRunHelper.java
...main/java/com/byit/thread/helper/TaskThreadRunHelper.java
+9
-5
No files found.
byit-myth-core/myth-admin-core/src/main/java/com/byit/enums/NodeNameEnum.java
View file @
818a7e1e
...
@@ -7,6 +7,7 @@ package com.byit.enums;
...
@@ -7,6 +7,7 @@ package com.byit.enums;
public
enum
NodeNameEnum
{
public
enum
NodeNameEnum
{
START_NODE
(
"start"
,
"开始节点"
),
START_NODE
(
"start"
,
"开始节点"
),
END_NODE
(
"end"
,
"结束节点"
),
;
;
private
String
nodeName
;
private
String
nodeName
;
private
String
details
;
private
String
details
;
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/thread/helper/SpecialFlowThreadRunHelper.java
0 → 100644
View file @
818a7e1e
package
com
.
byit
.
thread
.
helper
;
import
cn.hutool.core.collection.CollectionUtil
;
import
com.byit.enums.FlowPropertyEnum
;
import
com.byit.enums.ScheduleEnum
;
import
com.byit.model.JobTask
;
import
com.byit.model.JobTaskSchedule
;
import
com.byit.service.JobTaskService
;
import
com.byit.service.mapservice.RunRecordingAndJobTaskService
;
import
com.byit.thread.BaseThreadRunHelper
;
import
lombok.extern.slf4j.Slf4j
;
import
org.apache.commons.lang3.StringUtils
;
import
javax.sql.DataSource
;
import
java.util.ArrayList
;
import
java.util.List
;
/**
* 特殊情况的工作流扫描
* @author huangfu
*/
@Slf4j
public
class
SpecialFlowThreadRunHelper
extends
BaseThreadRunHelper
{
/**
* 读取任务节点的预读
*/
private
static
final
long
PRE_READ_MS
=
7000
;
private
static
final
String
LOCK_NAME
=
"special_flow_lock"
;
private
final
DataSource
dataSource
;
private
final
JobTaskService
jobTaskService
;
/**
* 运行实例和task的组合操作
*/
private
final
RunRecordingAndJobTaskService
runRecordingAndJobTaskService
;
public
SpecialFlowThreadRunHelper
(
DataSource
dataSource
,
JobTaskService
jobTaskService
,
RunRecordingAndJobTaskService
runRecordingAndJobTaskService
)
{
this
.
dataSource
=
dataSource
;
this
.
jobTaskService
=
jobTaskService
;
this
.
runRecordingAndJobTaskService
=
runRecordingAndJobTaskService
;
}
@Override
public
Long
start
()
{
long
nowTime
=
System
.
currentTimeMillis
();
//开始寻找此时 不是暂停状态,而且七秒内即将运行的任务 而且还不是暂停的节点
List
<
JobTask
>
repairs
=
jobTaskService
.
findJobTaskByTriggerNextTimeLessThanEqual
(
nowTime
+
PRE_READ_MS
,
ScheduleEnum
.
REPAIR
.
getCode
());
List
<
JobTask
>
repeats
=
jobTaskService
.
findJobTaskByTriggerNextTimeLessThanEqual
(
nowTime
+
PRE_READ_MS
,
ScheduleEnum
.
REPEAT
.
getCode
());
List
<
JobTask
>
jobTasks
=
new
ArrayList
<>(
8
);
if
(
CollectionUtil
.
isNotEmpty
(
repairs
)){
jobTasks
.
addAll
(
repairs
);
}
if
(
CollectionUtil
.
isNotEmpty
(
repeats
)){
jobTasks
.
addAll
(
repeats
);
}
if
(
CollectionUtil
.
isNotEmpty
(
jobTasks
))
{
List
<
JobTaskSchedule
>
jobTaskSchedules
=
new
ArrayList
<>(
15
);
for
(
JobTask
jobTask
:
jobTasks
)
{
//获取上级节点
String
nodeDepend
=
jobTask
.
getNodeDepend
();
if
(
FlowPropertyEnum
.
IS_INNER
.
getCode
().
equals
(
jobTask
.
getIsVirtual
()))
{
//虚节点处理操作
innerNodeOperating
(
jobTask
);
}
else
{
//判断上级节点是否存在 不存在 直接跑,存在则进行判断
if
(
StringUtils
.
isNotBlank
(
nodeDepend
)){
String
[]
split
=
nodeDepend
.
split
(
","
);
}
else
{
}
}
}
}
return
null
;
}
/**
* 内嵌节点处理操作
* @param thisJobTask 当前的任务节点
*/
private
void
innerNodeOperating
(
JobTask
thisJobTask
){
try
{
runRecordingAndJobTaskService
.
saveRunRecordingAndTask
(
thisJobTask
);
}
catch
(
Exception
e
)
{
log
.
error
(
"--------------------虚节点处理出现异常{}------------------"
,
e
.
getMessage
());
}
}
@Override
public
DataSource
getDataSource
()
{
return
dataSource
;
}
@Override
public
String
getLockName
()
{
return
LOCK_NAME
;
}
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/thread/helper/TaskThreadRunHelper.java
View file @
818a7e1e
package
com
.
byit
.
thread
.
helper
;
package
com
.
byit
.
thread
.
helper
;
import
cn.hutool.core.collection.CollectionUtil
;
import
cn.hutool.core.collection.CollectionUtil
;
import
com.byit.enums.FlowPropertyEnum
;
import
com.byit.enums.*
;
import
com.byit.enums.NodeNameEnum
;
import
com.byit.enums.NodePropertyEnum
;
import
com.byit.enums.NodeRunStatusPropertyEnum
;
import
com.byit.job.exceptions.BusinessException
;
import
com.byit.job.exceptions.BusinessException
;
import
com.byit.model.JobTask
;
import
com.byit.model.JobTask
;
import
com.byit.model.JobTaskRunLog
;
import
com.byit.model.JobTaskRunLog
;
...
@@ -140,8 +137,15 @@ public class TaskThreadRunHelper extends BaseThreadRunHelper {
...
@@ -140,8 +137,15 @@ public class TaskThreadRunHelper extends BaseThreadRunHelper {
}
}
}
}
String
runId
;
if
(
ScheduleEnum
.
REPEAT
.
getCode
()
.
equals
(
thisJobTask
.
getScheduleType
())
){
runId
=
thisJobTask
.
getReRunId
();
}
else
{
runId
=
thisJobTask
.
getRunId
();
}
//这里返回的是上级节点的日志执行情况 把运行中的数据给过滤掉了
//这里返回的是上级节点的日志执行情况 把运行中的数据给过滤掉了
List
<
JobTaskRunLog
>
jobTaskRunLogList
=
jobTaskRunLogService
.
findJobTaskRunLogNotEndNodeByRunCodeCount
(
dependIdByNodeId
,
thisJobTask
.
getRunId
()
);
List
<
JobTaskRunLog
>
jobTaskRunLogList
=
jobTaskRunLogService
.
findJobTaskRunLogNotEndNodeByRunCodeCount
(
dependIdByNodeId
,
runId
);
if
(
CollectionUtil
.
isNotEmpty
(
jobTaskRunLogList
))
{
if
(
CollectionUtil
.
isNotEmpty
(
jobTaskRunLogList
))
{
if
(
CollectionUtil
.
isEmpty
(
dependIdByNodeId
)){
if
(
CollectionUtil
.
isEmpty
(
dependIdByNodeId
)){
throw
new
BusinessException
(
NodeRunStatusPropertyEnum
.
NODE_RELY_ERROR
.
getMsg
());
throw
new
BusinessException
(
NodeRunStatusPropertyEnum
.
NODE_RELY_ERROR
.
getMsg
());
...
...
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