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
f7fa9088
Commit
f7fa9088
authored
Jan 08, 2020
by
huangfusuper
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
【添加代码】自动扫描即将执行的工作流
parent
22c83f51
Hide whitespace changes
Inline
Side-by-side
Showing
19 changed files
with
417 additions
and
9 deletions
+417
-9
MythJobScheduler.java
...in-core/src/main/java/com/byit/conf/MythJobScheduler.java
+6
-5
FlowMapper.java
...-admin-core/src/main/java/com/byit/mapper/FlowMapper.java
+9
-0
JobTaskMapper.java
...min-core/src/main/java/com/byit/mapper/JobTaskMapper.java
+7
-0
NodeMapper.java
...-admin-core/src/main/java/com/byit/mapper/NodeMapper.java
+11
-0
JobTask.java
...myth-admin-core/src/main/java/com/byit/model/JobTask.java
+8
-1
RunRecording.java
...admin-core/src/main/java/com/byit/model/RunRecording.java
+1
-1
FlowService.java
...dmin-core/src/main/java/com/byit/service/FlowService.java
+12
-0
JobTaskService.java
...n-core/src/main/java/com/byit/service/JobTaskService.java
+8
-0
NodeService.java
...dmin-core/src/main/java/com/byit/service/NodeService.java
+15
-0
RunNodeServer.java
...in-core/src/main/java/com/byit/service/RunNodeServer.java
+23
-0
FlowServiceImpl.java
.../src/main/java/com/byit/service/impl/FlowServiceImpl.java
+10
-0
JobTaskServiceImpl.java
...c/main/java/com/byit/service/impl/JobTaskServiceImpl.java
+5
-0
NodeServiceImpl.java
.../src/main/java/com/byit/service/impl/NodeServiceImpl.java
+6
-0
RunNodeServiceImpl.java
...c/main/java/com/byit/service/impl/RunNodeServiceImpl.java
+63
-0
EmailScanHelper.java
...n-core/src/main/java/com/byit/thread/EmailScanHelper.java
+7
-1
FlowScanHelper.java
...in-core/src/main/java/com/byit/thread/FlowScanHelper.java
+182
-0
FlowMapper.xml
.../myth-admin-core/src/main/resources/mapper/FlowMapper.xml
+7
-0
JobTaskMapper.xml
...th-admin-core/src/main/resources/mapper/JobTaskMapper.xml
+24
-1
NodeMapper.xml
.../myth-admin-core/src/main/resources/mapper/NodeMapper.xml
+13
-0
No files found.
byit-myth-core/myth-admin-core/src/main/java/com/byit/conf/MythJobScheduler.java
View file @
f7fa9088
package
com
.
byit
.
conf
;
package
com
.
byit
.
conf
;
import
com.byit.thread.EmailScanHelper
;
import
com.byit.thread.*
;
import
com.byit.thread.JobScheduleHelper
;
import
com.byit.thread.LogScanHelper
;
import
com.byit.thread.RunRecordingScanHelper
;
import
lombok.extern.slf4j.Slf4j
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.beans.factory.DisposableBean
;
import
org.springframework.beans.factory.DisposableBean
;
import
org.springframework.beans.factory.InitializingBean
;
import
org.springframework.beans.factory.InitializingBean
;
...
@@ -23,13 +20,15 @@ public class MythJobScheduler implements InitializingBean, DisposableBean {
...
@@ -23,13 +20,15 @@ public class MythJobScheduler implements InitializingBean, DisposableBean {
private
final
LogScanHelper
logScanHelper
;
private
final
LogScanHelper
logScanHelper
;
private
final
RunRecordingScanHelper
runRecordingScanHelper
;
private
final
RunRecordingScanHelper
runRecordingScanHelper
;
private
final
EmailScanHelper
emailScanHelper
;
private
final
EmailScanHelper
emailScanHelper
;
private
final
FlowScanHelper
flowScanHelper
;
@Autowired
@Autowired
public
MythJobScheduler
(
JobScheduleHelper
jobScheduleHelper
,
LogScanHelper
logScanHelper
,
RunRecordingScanHelper
runRecordingScanHelper
,
EmailScanHelper
emailScanHelper
)
{
public
MythJobScheduler
(
JobScheduleHelper
jobScheduleHelper
,
LogScanHelper
logScanHelper
,
RunRecordingScanHelper
runRecordingScanHelper
,
EmailScanHelper
emailScanHelper
,
FlowScanHelper
flowScanHelper
)
{
this
.
jobScheduleHelper
=
jobScheduleHelper
;
this
.
jobScheduleHelper
=
jobScheduleHelper
;
this
.
logScanHelper
=
logScanHelper
;
this
.
logScanHelper
=
logScanHelper
;
this
.
runRecordingScanHelper
=
runRecordingScanHelper
;
this
.
runRecordingScanHelper
=
runRecordingScanHelper
;
this
.
emailScanHelper
=
emailScanHelper
;
this
.
emailScanHelper
=
emailScanHelper
;
this
.
flowScanHelper
=
flowScanHelper
;
}
}
/**
/**
...
@@ -42,6 +41,7 @@ public class MythJobScheduler implements InitializingBean, DisposableBean {
...
@@ -42,6 +41,7 @@ public class MythJobScheduler implements InitializingBean, DisposableBean {
this
.
logScanHelper
.
doStop
();
this
.
logScanHelper
.
doStop
();
this
.
runRecordingScanHelper
.
doStop
();
this
.
runRecordingScanHelper
.
doStop
();
this
.
emailScanHelper
.
doStop
();
this
.
emailScanHelper
.
doStop
();
this
.
flowScanHelper
.
doStop
();
}
}
/**
/**
...
@@ -55,5 +55,6 @@ public class MythJobScheduler implements InitializingBean, DisposableBean {
...
@@ -55,5 +55,6 @@ public class MythJobScheduler implements InitializingBean, DisposableBean {
this
.
logScanHelper
.
start
();
this
.
logScanHelper
.
start
();
this
.
runRecordingScanHelper
.
start
();
this
.
runRecordingScanHelper
.
start
();
this
.
emailScanHelper
.
start
();
this
.
emailScanHelper
.
start
();
this
.
flowScanHelper
.
start
();
}
}
}
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/FlowMapper.java
View file @
f7fa9088
...
@@ -3,7 +3,16 @@ package com.byit.mapper;
...
@@ -3,7 +3,16 @@ package com.byit.mapper;
import
com.byit.model.Flow
;
import
com.byit.model.Flow
;
import
org.apache.ibatis.annotations.Param
;
import
org.apache.ibatis.annotations.Param
;
import
java.util.List
;
public
interface
FlowMapper
{
public
interface
FlowMapper
{
/**
* 查询半个小时内即将要执行的工作流
* @param triggerNextTime
* @return
*/
List
<
Flow
>
findHalfAnHourFlow
(
@Param
(
"triggerNextTime"
)
Long
triggerNextTime
);
int
deleteById
(
Integer
flowId
);
int
deleteById
(
Integer
flowId
);
int
insertSelective
(
Flow
record
);
int
insertSelective
(
Flow
record
);
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/JobTaskMapper.java
View file @
f7fa9088
...
@@ -40,6 +40,13 @@ public interface JobTaskMapper {
...
@@ -40,6 +40,13 @@ public interface JobTaskMapper {
int
saveJobTask
(
JobTask
jobTask
);
int
saveJobTask
(
JobTask
jobTask
);
/**
/**
* 批量保存
* @param jobTasks
* @return
*/
int
saveJobTasks
(
@Param
(
"jobTasks"
)
List
<
JobTask
>
jobTasks
);
/**
* 修改任务表
* 修改任务表
* @param jobTask
* @param jobTask
* @return
* @return
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/mapper/NodeMapper.java
View file @
f7fa9088
...
@@ -5,7 +5,18 @@ import org.apache.ibatis.annotations.Param;
...
@@ -5,7 +5,18 @@ import org.apache.ibatis.annotations.Param;
import
java.util.List
;
import
java.util.List
;
/**
* @author HUANGFU
*/
public
interface
NodeMapper
{
public
interface
NodeMapper
{
/**
* 跟怒工作流ID和版本名称查询所有的节点
* @param flowId
* @param versionName
* @return
*/
List
<
Node
>
findNodeByFlowIdAndVersionName
(
@Param
(
"flowId"
)
Integer
flowId
,
@Param
(
"versionName"
)
String
versionName
);
int
deleteById
(
Integer
nodeId
);
int
deleteById
(
Integer
nodeId
);
int
insertSelective
(
Node
record
);
int
insertSelective
(
Node
record
);
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/model/JobTask.java
View file @
f7fa9088
...
@@ -3,13 +3,20 @@ package com.byit.model;
...
@@ -3,13 +3,20 @@ package com.byit.model;
import
io.swagger.annotations.ApiModel
;
import
io.swagger.annotations.ApiModel
;
import
io.swagger.annotations.ApiModelProperty
;
import
io.swagger.annotations.ApiModelProperty
;
import
java.io.Serializable
;
import
java.io.Serializable
;
import
lombok.AllArgsConstructor
;
import
lombok.Builder
;
import
lombok.Data
;
import
lombok.Data
;
import
lombok.NoArgsConstructor
;
/**
/**
*
*
*/
*/
@ApiModel
@ApiModel
@Data
@Data
@AllArgsConstructor
@NoArgsConstructor
@Builder
public
class
JobTask
implements
Serializable
{
public
class
JobTask
implements
Serializable
{
/**
/**
*/
*/
...
@@ -119,7 +126,7 @@ public class JobTask implements Serializable {
...
@@ -119,7 +126,7 @@ public class JobTask implements Serializable {
private
String
routingStrategy
;
private
String
routingStrategy
;
/**
/**
* 运行标识
* 运行标识
-----
*/
*/
@ApiModelProperty
(
"运行标识"
)
@ApiModelProperty
(
"运行标识"
)
private
String
runId
;
private
String
runId
;
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/model/RunRecording.java
View file @
f7fa9088
...
@@ -52,7 +52,7 @@ public class RunRecording implements Serializable {
...
@@ -52,7 +52,7 @@ public class RunRecording implements Serializable {
/**
/**
* 执行结果 1 成功 2 失败 3 补批成功 4 补批失败 5.kill
* 执行结果 1 成功 2 失败 3 补批成功 4 补批失败 5.kill
*/
*/
@ApiModelProperty
(
"执行结果 1 成功 2 失败 3 补批成功 4 补批失败 5.kill"
)
@ApiModelProperty
(
"执行结果
0未执行
1 成功 2 失败 3 补批成功 4 补批失败 5.kill"
)
private
String
flowRunResult
;
private
String
flowRunResult
;
/**
/**
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/FlowService.java
View file @
f7fa9088
...
@@ -3,6 +3,9 @@ package com.byit.service;
...
@@ -3,6 +3,9 @@ package com.byit.service;
import
com.byit.model.Flow
;
import
com.byit.model.Flow
;
import
com.byit.model.vo.FlowVo
;
import
com.byit.model.vo.FlowVo
;
import
com.byit.model.vo.NodeVo
;
import
com.byit.model.vo.NodeVo
;
import
org.apache.ibatis.annotations.Param
;
import
java.util.List
;
/**
/**
* @description: 工作流业务逻辑接口
* @description: 工作流业务逻辑接口
...
@@ -10,6 +13,13 @@ import com.byit.model.vo.NodeVo;
...
@@ -10,6 +13,13 @@ import com.byit.model.vo.NodeVo;
* @create: 2019-12-23 17:33
* @create: 2019-12-23 17:33
*/
*/
public
interface
FlowService
{
public
interface
FlowService
{
/**
* 查询半个小时内即将要执行的工作流
* @param triggerNextTime
* @return
*/
List
<
Flow
>
findHalfAnHourFlow
(
Long
triggerNextTime
);
/**
/**
* 保存工作流信息
* 保存工作流信息
...
@@ -30,6 +40,8 @@ public interface FlowService {
...
@@ -30,6 +40,8 @@ public interface FlowService {
*/
*/
void
updateFlow
(
FlowVo
flowVo
);
void
updateFlow
(
FlowVo
flowVo
);
int
updateByIdSelective
(
Flow
record
);
/**
/**
* 真实删除当前表工作流信息
* 真实删除当前表工作流信息
* @param flowId
* @param flowId
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/JobTaskService.java
View file @
f7fa9088
package
com
.
byit
.
service
;
package
com
.
byit
.
service
;
import
com.byit.model.JobTask
;
import
com.byit.model.JobTask
;
import
org.apache.ibatis.annotations.Param
;
import
java.util.List
;
import
java.util.List
;
...
@@ -25,6 +26,13 @@ public interface JobTaskService {
...
@@ -25,6 +26,13 @@ public interface JobTaskService {
void
addMythJobTask
(
JobTask
jobTask
);
void
addMythJobTask
(
JobTask
jobTask
);
/**
/**
* 批量保存
* @param jobTasks
* @return
*/
int
saveJobTasks
(
List
<
JobTask
>
jobTasks
);
/**
* 删除已经被调度的任务根据ID
* 删除已经被调度的任务根据ID
* @param id
* @param id
*/
*/
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/NodeService.java
View file @
f7fa9088
package
com
.
byit
.
service
;
package
com
.
byit
.
service
;
import
com.byit.model.Node
;
import
org.apache.ibatis.annotations.Param
;
import
java.util.List
;
/**
/**
* @description: 任务节点逻辑处理接口
* @description: 任务节点逻辑处理接口
* @author: gml
* @author: gml
* @create: 2019-12-24 10:38
* @create: 2019-12-24 10:38
*/
*/
public
interface
NodeService
{
public
interface
NodeService
{
/**
* 跟怒工作流ID和版本名称查询所有的节点
* @param flowId
* @param versionName
* @return
*/
List
<
Node
>
findNodeByFlowIdAndVersionName
(
Integer
flowId
,
String
versionName
);
}
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/RunNodeServer.java
0 → 100644
View file @
f7fa9088
package
com
.
byit
.
service
;
import
com.byit.model.Flow
;
import
com.byit.model.Node
;
import
org.springframework.stereotype.Service
;
import
org.springframework.transaction.annotation.Propagation
;
import
org.springframework.transaction.annotation.Transactional
;
import
java.util.List
;
/**
*封装事务
* @author huangfu
*/
public
interface
RunNodeServer
{
/**
* 保存执行记录task
* @param flow
* @param nodes
*/
void
saveRunRecAndTask
(
Flow
flow
,
List
<
Node
>
nodes
);
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/impl/FlowServiceImpl.java
View file @
f7fa9088
...
@@ -44,6 +44,11 @@ public class FlowServiceImpl implements FlowService {
...
@@ -44,6 +44,11 @@ public class FlowServiceImpl implements FlowService {
private
NodeDependencyMapper
nodeDependencyMapper
;
private
NodeDependencyMapper
nodeDependencyMapper
;
@Override
@Override
public
List
<
Flow
>
findHalfAnHourFlow
(
Long
triggerNextTime
)
{
return
flowMapper
.
findHalfAnHourFlow
(
triggerNextTime
);
}
@Override
public
Flow
saveJobFlow
(
FlowVo
flowVo
)
{
public
Flow
saveJobFlow
(
FlowVo
flowVo
)
{
ValidationUtil
.
dataNotNull
(
flowVo
.
getFlowId
(),
"工作流Id不允许为空!"
);
ValidationUtil
.
dataNotNull
(
flowVo
.
getFlowId
(),
"工作流Id不允许为空!"
);
...
@@ -197,6 +202,11 @@ public class FlowServiceImpl implements FlowService {
...
@@ -197,6 +202,11 @@ public class FlowServiceImpl implements FlowService {
}
}
@Override
@Override
public
int
updateByIdSelective
(
Flow
record
)
{
return
flowMapper
.
updateByIdSelective
(
record
);
}
@Override
public
void
deleteFlow
(
Integer
flowId
)
{
public
void
deleteFlow
(
Integer
flowId
)
{
}
}
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/impl/JobTaskServiceImpl.java
View file @
f7fa9088
...
@@ -46,6 +46,11 @@ public class JobTaskServiceImpl implements JobTaskService {
...
@@ -46,6 +46,11 @@ public class JobTaskServiceImpl implements JobTaskService {
jobTaskMapper
.
saveJobTask
(
jobTask
);
jobTaskMapper
.
saveJobTask
(
jobTask
);
}
}
@Override
public
int
saveJobTasks
(
List
<
JobTask
>
jobTasks
)
{
return
jobTaskMapper
.
saveJobTasks
(
jobTasks
);
}
/**
/**
* 根据id删除一个任务
* 根据id删除一个任务
* @param id 任务的id
* @param id 任务的id
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/impl/NodeServiceImpl.java
View file @
f7fa9088
package
com
.
byit
.
service
.
impl
;
package
com
.
byit
.
service
.
impl
;
import
com.byit.mapper.NodeMapper
;
import
com.byit.mapper.NodeMapper
;
import
com.byit.model.Node
;
import
com.byit.service.NodeService
;
import
com.byit.service.NodeService
;
import
org.springframework.stereotype.Service
;
import
org.springframework.stereotype.Service
;
import
javax.annotation.Resource
;
import
javax.annotation.Resource
;
import
java.util.List
;
/**
/**
* @description: 任务节点逻辑处理实现类
* @description: 任务节点逻辑处理实现类
...
@@ -17,4 +19,8 @@ public class NodeServiceImpl implements NodeService {
...
@@ -17,4 +19,8 @@ public class NodeServiceImpl implements NodeService {
@Resource
@Resource
private
NodeMapper
nodeMapper
;
private
NodeMapper
nodeMapper
;
@Override
public
List
<
Node
>
findNodeByFlowIdAndVersionName
(
Integer
flowId
,
String
versionName
)
{
return
nodeMapper
.
findNodeByFlowIdAndVersionName
(
flowId
,
versionName
);
}
}
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/impl/RunNodeServiceImpl.java
0 → 100644
View file @
f7fa9088
package
com
.
byit
.
service
.
impl
;
import
com.byit.model.Flow
;
import
com.byit.model.JobTask
;
import
com.byit.model.Node
;
import
com.byit.model.RunRecording
;
import
com.byit.service.JobTaskService
;
import
com.byit.service.RunNodeServer
;
import
com.byit.service.RunRecordingService
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.beans.BeanUtils
;
import
org.springframework.stereotype.Service
;
import
org.springframework.transaction.annotation.Propagation
;
import
org.springframework.transaction.annotation.Transactional
;
import
java.util.ArrayList
;
import
java.util.List
;
import
java.util.UUID
;
/**
* @author huangfu
*/
@Service
@Transactional
(
propagation
=
Propagation
.
REQUIRED
,
rollbackFor
=
Exception
.
class
)
@Slf4j
public
class
RunNodeServiceImpl
implements
RunNodeServer
{
private
final
RunRecordingService
runRecordingService
;
private
final
JobTaskService
jobTaskService
;
public
RunNodeServiceImpl
(
RunRecordingService
runRecordingService
,
JobTaskService
jobTaskService
)
{
this
.
runRecordingService
=
runRecordingService
;
this
.
jobTaskService
=
jobTaskService
;
}
@Override
public
void
saveRunRecAndTask
(
Flow
flow
,
List
<
Node
>
nodes
)
{
log
.
info
(
"---------saveRunRecAndTask start------【保存工作流:{}和节点:{}】-----------------------"
,
flow
,
nodes
);
String
runId
=
UUID
.
randomUUID
().
toString
().
replace
(
"-"
,
""
);
log
.
info
(
"-------------【开始保存运行记录runId为:{}】------------------"
,
runId
);
RunRecording
build
=
new
RunRecording
();
BeanUtils
.
copyProperties
(
flow
,
build
);
build
.
setRunId
(
runId
);
build
.
setDispatchIp
(
"127.0.0.1"
);
build
.
setFlowVersionName
(
flow
.
getVersionName
());
build
.
setTriggerTime
(
flow
.
getTriggerNextTime
());
runRecordingService
.
saveRunRecording
(
build
);
boolean
isScheduleFollow
=
"1"
.
equals
(
flow
.
getScheduleFollow
());
log
.
info
(
"-------------【开始保存节点信息,任务是否为跟随工作流,{}】---------------"
,
isScheduleFollow
);
List
<
JobTask
>
jobTasks
=
new
ArrayList
<>(
32
);
nodes
.
forEach
(
node
->
{
JobTask
jobTask
=
new
JobTask
();
BeanUtils
.
copyProperties
(
node
,
jobTask
);
if
(
isScheduleFollow
){
jobTask
.
setTriggerTime
(
flow
.
getTriggerNextTime
());
}
jobTask
.
setRunId
(
runId
);
jobTasks
.
add
(
jobTask
);
});
jobTaskService
.
saveJobTasks
(
jobTasks
);
log
.
info
(
"-------saveRunRecAndTask end-----------【运行结束】-----------------"
);
}
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/thread/EmailScanHelper.java
View file @
f7fa9088
...
@@ -46,6 +46,8 @@ public class EmailScanHelper {
...
@@ -46,6 +46,8 @@ public class EmailScanHelper {
dateAligned
(
5000
);
dateAligned
(
5000
);
log
.
info
(
"--------------【com.byit.thread.EmailScanHelper.startScanNotSentEmailFlow】init success---------------"
);
log
.
info
(
"--------------【com.byit.thread.EmailScanHelper.startScanNotSentEmailFlow】init success---------------"
);
while
(!
emailThreadIsStop
){
while
(!
emailThreadIsStop
){
//是否需要睡眠
boolean
isSleep
=
false
;
Connection
conn
=
null
;
Connection
conn
=
null
;
Boolean
connAutoCommit
=
null
;
Boolean
connAutoCommit
=
null
;
PreparedStatement
preparedStatement
=
null
;
PreparedStatement
preparedStatement
=
null
;
...
@@ -67,7 +69,7 @@ public class EmailScanHelper {
...
@@ -67,7 +69,7 @@ public class EmailScanHelper {
emailAlarmService
.
sendEmail
(
emailAlarmVo
);
emailAlarmService
.
sendEmail
(
emailAlarmVo
);
});
});
}
else
{
}
else
{
dateAligned
(
20000
)
;
isSleep
=
true
;
}
}
}
catch
(
Exception
e
){
}
catch
(
Exception
e
){
e
.
printStackTrace
();
e
.
printStackTrace
();
...
@@ -111,6 +113,10 @@ public class EmailScanHelper {
...
@@ -111,6 +113,10 @@ public class EmailScanHelper {
}
}
}
}
}
}
if
(
isSleep
){
dateAligned
(
20000
);
}
}
}
});
});
emailThread
.
setDaemon
(
true
);
emailThread
.
setDaemon
(
true
);
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/thread/FlowScanHelper.java
0 → 100644
View file @
f7fa9088
package
com
.
byit
.
thread
;
import
cn.hutool.core.collection.CollectionUtil
;
import
com.byit.model.Flow
;
import
com.byit.model.JobTask
;
import
com.byit.model.Node
;
import
com.byit.model.RunRecording
;
import
com.byit.service.*
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.beans.BeanUtils
;
import
org.springframework.stereotype.Component
;
import
javax.sql.DataSource
;
import
java.sql.Connection
;
import
java.sql.PreparedStatement
;
import
java.sql.SQLException
;
import
java.util.ArrayList
;
import
java.util.List
;
import
java.util.UUID
;
import
java.util.concurrent.TimeUnit
;
/**
* 扫描任务流表的线程
* 目的:扫描即将要执行的任务流表 半个小时
* @author huangfu
*/
@Component
@Slf4j
public
class
FlowScanHelper
{
private
final
Long
PRE_TEST_TIME
=
System
.
currentTimeMillis
()+
TimeUnit
.
HOURS
.
toMillis
(
30
);
private
final
FlowService
flowService
;
private
final
DataSource
dataSource
;
private
final
NodeService
nodeService
;
private
final
JobTaskService
jobTaskService
;
private
final
RunRecordingService
runRecordingService
;
private
final
RunNodeServer
runNodeServer
;
private
Thread
flowThread
;
private
volatile
boolean
flowThreadIsStop
=
false
;
public
FlowScanHelper
(
FlowService
flowService
,
DataSource
dataSource
,
NodeService
nodeService
,
JobTaskService
jobTaskService
,
RunRecordingService
runRecordingService
,
RunNodeServer
runNodeServer
)
{
this
.
flowService
=
flowService
;
this
.
dataSource
=
dataSource
;
this
.
nodeService
=
nodeService
;
this
.
jobTaskService
=
jobTaskService
;
this
.
runRecordingService
=
runRecordingService
;
this
.
runNodeServer
=
runNodeServer
;
}
public
void
start
(){
flowThreadStart
();
}
/**
* 运行扫描工作流的线程
*/
private
void
flowThreadStart
(){
flowThread
=
new
Thread
(()
->{
dateAligned
(
5000
);
log
.
info
(
"--------------【com.byit.thread.FlowScanHelper#flowThreadStart】init success---------------"
);
while
(!
flowThreadIsStop
){
Connection
conn
=
null
;
Boolean
connAutoCommit
=
null
;
PreparedStatement
preparedStatement
=
null
;
//是否需要睡眠
boolean
isSleep
=
false
;
try
{
/**
* 添加行锁
*/
conn
=
dataSource
.
getConnection
();
connAutoCommit
=
conn
.
getAutoCommit
();
conn
.
setAutoCommit
(
false
);
preparedStatement
=
conn
.
prepareStatement
(
"SELECT * FROM JOB_LOCK WHERE LOCK_NAME = 'flow_lock' FOR UPDATE "
);
preparedStatement
.
execute
();
List
<
Flow
>
halfAnHourFlow
=
flowService
.
findHalfAnHourFlow
(
PRE_TEST_TIME
);
if
(
CollectionUtil
.
isNotEmpty
(
halfAnHourFlow
))
{
halfAnHourFlow
.
forEach
(
flow
->
{
log
.
debug
(
"-----------------【工作流{}的执行次数大于0,放行】-------------------------"
,
flow
.
getFlowName
());
String
versionName
=
flow
.
getVersionName
();
Integer
flowId
=
flow
.
getFlowId
();
List
<
Node
>
nodeByFlowIdAndVersionName
=
nodeService
.
findNodeByFlowIdAndVersionName
(
flowId
,
versionName
);
runNodeServer
.
saveRunRecAndTask
(
flow
,
nodeByFlowIdAndVersionName
);
flow
.
setRepeatCount
(
flow
.
getRepeatCount
()-
1
);
flowService
.
updateByIdSelective
(
flow
);
});
}
else
{
isSleep
=
true
;
}
}
catch
(
Exception
e
){
e
.
printStackTrace
();
}
finally
{
//释放资源
if
(
conn
!=
null
){
try
{
conn
.
commit
();
}
catch
(
SQLException
e
)
{
if
(!
flowThreadIsStop
){
log
.
error
(
"--------------------【提交行锁出错】---------------------"
);
}
}
}
try
{
if
(
conn
!=
null
){
conn
.
setAutoCommit
(
connAutoCommit
);
}
}
catch
(
SQLException
e
)
{
if
(!
flowThreadIsStop
){
log
.
error
(
"--------------------【恢复自动提交出错】---------------------"
);
}
}
if
(
preparedStatement
!=
null
){
try
{
preparedStatement
.
close
();
}
catch
(
SQLException
e
)
{
if
(!
flowThreadIsStop
){
log
.
error
(
"--------------------【关闭执行器出错】---------------------"
);
}
}
}
try
{
conn
.
close
();
}
catch
(
SQLException
e
)
{
if
(!
flowThreadIsStop
){
log
.
error
(
"--------------------【关闭链接出错】---------------------"
);
}
}
}
if
(
isSleep
){
try
{
log
.
info
(
"------------------【未扫描到要执行的工作流】-------------------"
);
TimeUnit
.
HOURS
.
sleep
(
10
);
}
catch
(
InterruptedException
e
)
{
log
.
warn
(
"----------------------【扫描工作流的线程被关闭了】----------------------------"
);
}
}
}
});
flowThread
.
setDaemon
(
true
);
flowThread
.
setName
(
"myth-job#【FlowScanHelper】#flowThreadStart"
);
flowThread
.
start
();
}
public
void
doStop
(){
this
.
flowThreadIsStop
=
true
;
try
{
TimeUnit
.
SECONDS
.
sleep
(
1
);
}
catch
(
InterruptedException
e
)
{
e
.
printStackTrace
(
);
}
if
(
flowThread
.
getState
()
!=
Thread
.
State
.
TERMINATED
)
{
flowThread
.
interrupt
();
try
{
flowThread
.
join
();
}
catch
(
InterruptedException
e
)
{
e
.
printStackTrace
(
);
}
}
log
.
warn
(
"---------------【工作扫描线程被注销】-----------------------"
);
}
/**
* 对齐时钟。整秒运行
*/
private
void
dateAligned
(
long
waitTime
){
try
{
TimeUnit
.
MILLISECONDS
.
sleep
(
waitTime
-
System
.
currentTimeMillis
()%
1000
);
}
catch
(
InterruptedException
e
)
{
log
.
warn
(
"----------------【线程被中断】-----------------------"
);
}
}
}
byit-myth-core/myth-admin-core/src/main/resources/mapper/FlowMapper.xml
View file @
f7fa9088
...
@@ -33,6 +33,13 @@
...
@@ -33,6 +33,13 @@
author, add_time, start_up, principal, version_name, repeat_count, remaining_count,
author, add_time, start_up, principal, version_name, repeat_count, remaining_count,
schedule_follow, is_update
schedule_follow, is_update
</sql>
</sql>
<select
id=
"findHalfAnHourFlow"
resultMap=
"BaseResultMap"
>
select
<include
refid=
"Base_Column_List"
/>
from flow where trigger_next_time
<![CDATA[ <= ]]>
#{triggerNextTime,jdbcType=BIGINT} and remaining_count
<![CDATA[ <> ]]>
0
</select>
<select
id=
"getById"
parameterType=
"java.lang.Integer"
resultMap=
"BaseResultMap"
>
<select
id=
"getById"
parameterType=
"java.lang.Integer"
resultMap=
"BaseResultMap"
>
<!-- generated @mbg.generated date: 2019-12-31 -->
<!-- generated @mbg.generated date: 2019-12-31 -->
select
select
...
...
byit-myth-core/myth-admin-core/src/main/resources/mapper/JobTaskMapper.xml
View file @
f7fa9088
...
@@ -58,7 +58,7 @@
...
@@ -58,7 +58,7 @@
</select>
</select>
<!--根据id查询-->
<!--根据id查询-->
<select
id=
"findJobTaskById"
parameterType=
"java.lang.Integer"
resultMap=
"ResultMapWithBLOBs"
>
<select
id=
"findJobTaskById"
parameterType=
"java.lang.Integer"
resultMap=
"ResultMapWithBLOBs"
>
select
select
<include
refid=
"Base_Column_List"
/>
<include
refid=
"Base_Column_List"
/>
,
,
<include
refid=
"Blob_Column_List"
/>
<include
refid=
"Blob_Column_List"
/>
...
@@ -242,6 +242,29 @@
...
@@ -242,6 +242,29 @@
</if>
</if>
</trim>
</trim>
</insert>
</insert>
<insert
id=
"saveJobTasks"
parameterType=
"com.byit.model.JobTask"
>
insert into job_task (
id, node_id, block_strategy, plugin_token, failed_retry_count, flow_id, gateway_token,
job_type, handler_name, node_desc, node_name, map_flow_id, node_timeout, is_virtual,
plugin_urls, priority, failed_retry_interval, routing_strategy, run_id, run_param,
run_source_desc, script_urls, source_principal, trigger_time, trigger_status, version_name,
run_command,run_source
) values
<foreach
collection=
"jobTasks"
item=
"jobTask"
separator =
","
>
(
#{jobTask.id,jdbcType=INTEGER}, #{jobTask.nodeId,jdbcType=INTEGER},#{jobTask.blockStrategy,jdbcType=VARCHAR},#{jobTask.pluginToken,jdbcType=VARCHAR},
#{jobTask.failedRetryCount,jdbcType=INTEGER},#{jobTask.flowId,jdbcType=INTEGER}, #{jobTask.gatewayToken,jdbcType=VARCHAR},#{jobTask.jobType,jdbcType=VARCHAR},
#{jobTask.handlerName,jdbcType=VARCHAR},#{jobTask.nodeDesc,jdbcType=VARCHAR},#{jobTask.nodeName,jdbcType=VARCHAR},#{jobTask.mapFlowId,jdbcType=INTEGER},
#{jobTask.nodeTimeout,jdbcType=BIGINT}, #{jobTask.isVirtual,jdbcType=CHAR},#{jobTask.pluginUrls,jdbcType=VARCHAR},#{jobTask.priority,jdbcType=CHAR},
#{jobTask.failedRetryInterval,jdbcType=BIGINT},#{jobTask.routingStrategy,jdbcType=VARCHAR},#{jobTask.runId,jdbcType=VARCHAR},
#{jobTask.runParam,jdbcType=VARCHAR},#{jobTask.runSourceDesc,jdbcType=VARCHAR},#{jobTask.scriptUrls,jdbcType=VARCHAR},
#{jobTask.sourcePrincipal,jdbcType=VARCHAR},#{jobTask.triggerTime,jdbcType=BIGINT},#{jobTask.triggerStatus,jdbcType=CHAR},
#{jobTask.versionName,jdbcType=VARCHAR},#{jobTask.runCommand,jdbcType=VARCHAR},#{jobTask.runSource,jdbcType=LONGVARCHAR}
)
</foreach>
</insert>
<update
id=
"updateJobTask"
parameterType=
"com.byit.model.JobTask"
>
<update
id=
"updateJobTask"
parameterType=
"com.byit.model.JobTask"
>
<!-- generated @mbg.generated date: 2019-12-25 -->
<!-- generated @mbg.generated date: 2019-12-25 -->
update job_task
update job_task
...
...
byit-myth-core/myth-admin-core/src/main/resources/mapper/NodeMapper.xml
View file @
f7fa9088
...
@@ -47,10 +47,23 @@
...
@@ -47,10 +47,23 @@
routing_strategy, run_param, run_source_desc, script_urls, source_principal, source_update_time,
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
trigger_next_time, author, add_time, version_name, on_fork, run_command
</sql>
</sql>
<sql
id=
"Blob_Column_List"
>
<sql
id=
"Blob_Column_List"
>
<!-- generated @mbg.generated date: 2019-12-31 -->
<!-- generated @mbg.generated date: 2019-12-31 -->
run_source
run_source
</sql>
</sql>
<select
id=
"findNodeByFlowIdAndVersionName"
resultMap=
"ResultMapWithBLOBs"
>
select
<include
refid=
"Base_Column_List"
/>
,
<include
refid=
"Blob_Column_List"
/>
from node
where flow_id = #{flowId,jdbcType=INTEGER} AND version_name = #{versionName,jdbcType=VARCHAR}
</select>
<select
id=
"getById"
parameterType=
"java.lang.Integer"
resultMap=
"ResultMapWithBLOBs"
>
<select
id=
"getById"
parameterType=
"java.lang.Integer"
resultMap=
"ResultMapWithBLOBs"
>
<!-- generated @mbg.generated date: 2019-12-31 -->
<!-- generated @mbg.generated date: 2019-12-31 -->
select
select
...
...
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