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
844a728f
Commit
844a728f
authored
Dec 10, 2019
by
huangfusuper
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
添加基本的任务7秒预读扫描,和扫描后删除
parent
1d4d02ea
Show whitespace changes
Inline
Side-by-side
Showing
25 changed files
with
591 additions
and
83 deletions
+591
-83
DataSourceConf.java
...yth-admin/src/main/java/com/byit/conf/DataSourceConf.java
+0
-3
JobController.java
...dmin/src/main/java/com/byit/controller/JobController.java
+6
-6
MythJobAutoConfigure.java
...ore/src/main/java/com/byit/conf/MythJobAutoConfigure.java
+0
-14
JobFlightScheduleMapper.java
...e/src/main/java/com/byit/dao/JobFlightScheduleMapper.java
+22
-0
JobFlowNodesMapper.java
...n-core/src/main/java/com/byit/dao/JobFlowNodesMapper.java
+19
-0
JobInfoMapper.java
...-admin-core/src/main/java/com/byit/dao/JobInfoMapper.java
+8
-2
JobFlightScheduleService.java
.../main/java/com/byit/service/JobFlightScheduleService.java
+21
-0
JobReadAheadService.java
...e/src/main/java/com/byit/service/JobReadAheadService.java
+9
-4
JobFlightScheduleServiceImpl.java
...a/com/byit/service/impl/JobFlightScheduleServiceImpl.java
+25
-0
JobReadAheadServiceImpl.java
...n/java/com/byit/service/impl/JobReadAheadServiceImpl.java
+11
-6
JobScheduleHelper.java
...core/src/main/java/com/byit/thread/JobScheduleHelper.java
+38
-17
JobFlightScheduleMapper.xml
...ore/src/main/resources/mapper/JobFlightScheduleMapper.xml
+73
-0
JobFlowNodesMapper.xml
...min-core/src/main/resources/mapper/JobFlowNodesMapper.xml
+47
-0
MythJobInfoMapper.xml
...dmin-core/src/main/resources/mapper/MythJobInfoMapper.xml
+12
-8
GlueJobHandler.java
...c/main/java/com/byit/job/handler/impl/GlueJobHandler.java
+1
-1
ScriptJobHandler.java
...main/java/com/byit/job/handler/impl/ScriptJobHandler.java
+1
-1
IJobHandler.java
...ain/java/com/byit/job/handler/interfaces/IJobHandler.java
+1
-1
MythJobFlightSchedule.java
...c/main/java/com/byit/job/model/MythJobFlightSchedule.java
+8
-16
MythJobFlowNodes.java
...on/src/main/java/com/byit/job/model/MythJobFlowNodes.java
+141
-0
MythJobReadAhead.java
...on/src/main/java/com/byit/job/model/MythJobReadAhead.java
+144
-0
ReturnResult.java
...re-common/src/main/java/com/byit/job/vo/ReturnResult.java
+1
-1
RunJobServerHandler.java
...lugin/src/main/java/com/byit/rpc/RunJobServerHandler.java
+1
-1
DemoJob.java
.../byit-demo-client/src/main/java/com/byit/job/DemoJob.java
+1
-1
DemoJob1.java
...byit-demo-client/src/main/java/com/byit/job/DemoJob1.java
+1
-1
任务调度插件开发流程说明.docx
doc/任务调度插件开发流程说明.docx
+0
-0
No files found.
byit-myth-admin/src/main/java/com/byit/conf/DataSourceConf.java
View file @
844a728f
package
com
.
byit
.
conf
;
import
com.byit.service.JobInfoService
;
import
lombok.Data
;
import
org.springframework.beans.factory.InitializingBean
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.stereotype.Component
;
import
javax.annotation.Resource
;
...
...
byit-myth-admin/src/main/java/com/byit/controller/JobController.java
View file @
844a728f
package
com
.
byit
.
controller
;
import
com.byit.job.WorkRoulette
;
import
com.byit.job.model.MythJob
Info
;
import
com.byit.job.model.MythJob
ReadAhead
;
import
com.byit.job.model.PluginBeanJobInfo
;
import
com.byit.job.
model
.ReturnResult
;
import
com.byit.service.Job
Info
Service
;
import
com.byit.job.
vo
.ReturnResult
;
import
com.byit.service.Job
ReadAhead
Service
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.web.bind.annotation.*
;
...
...
@@ -23,7 +23,7 @@ public class JobController {
@Autowired
private
DataSource
dataSource
;
@Autowired
private
Job
InfoService
jobInfo
Service
;
private
Job
ReadAheadService
jobReadAhead
Service
;
@PostMapping
(
value
=
"addJob"
)
public
String
addJob
(
@RequestBody
PluginBeanJobInfo
pluginBeanJobInfo
){
WorkRoulette
.
addJob
(
pluginBeanJobInfo
);
...
...
@@ -36,7 +36,7 @@ public class JobController {
}
@GetMapping
(
value
=
"getJobInfo"
)
public
List
<
MythJob
Info
>
getJobInfo
(){
return
job
InfoService
.
findMythJobInfo
ByTriggerNextTime
(
11
);
public
List
<
MythJob
ReadAhead
>
getJobInfo
(){
return
job
ReadAheadService
.
findMythJobReadAhead
ByTriggerNextTime
(
11
);
}
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/conf/MythJobAutoConfigure.java
View file @
844a728f
package
com
.
byit
.
conf
;
import
com.byit.service.impl.JobInfoServiceImpl
;
import
com.byit.service.JobInfoService
;
import
org.mybatis.spring.annotation.MapperScan
;
import
org.springframework.beans.factory.annotation.Value
;
import
org.springframework.boot.autoconfigure.AutoConfigureBefore
;
import
org.springframework.boot.autoconfigure.condition.ConditionalOnBean
;
import
org.springframework.boot.autoconfigure.condition.ConditionalOnClass
;
import
org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration
;
import
org.springframework.boot.autoconfigure.jdbc.DataSourceBuilder
;
import
org.springframework.boot.autoconfigure.jdbc.DataSourceProperties
;
import
org.springframework.boot.context.properties.ConfigurationProperties
;
import
org.springframework.boot.context.properties.EnableConfigurationProperties
;
import
org.springframework.context.annotation.*
;
import
org.springframework.jdbc.datasource.DriverManagerDataSource
;
import
javax.sql.DataSource
;
/**
* @program: byit-myth-job->AppConf
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/dao/JobFlightScheduleMapper.java
0 → 100644
View file @
844a728f
package
com
.
byit
.
dao
;
import
com.byit.job.model.MythJobFlightSchedule
;
import
org.apache.ibatis.annotations.Param
;
import
org.springframework.stereotype.Repository
;
import
java.util.List
;
/**
* @program: byit-myth-job->JobFlightScheduleMapper
* @description: 任务排期表 这个就是预读五秒的数据,要添加进任务调度轮的数据
* @author: huangfu
* @date: 2019/12/10 15:06
**/
@Repository
public
interface
JobFlightScheduleMapper
{
/**
* 批量向排期表添加任务数据
* @param mythJobFlightSchedules
*/
void
addJobFlightSchedules
(
@Param
(
"mythJobFlightSchedule"
)
List
<
MythJobFlightSchedule
>
mythJobFlightSchedules
);
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/dao/JobFlowNodesMapper.java
0 → 100644
View file @
844a728f
package
com
.
byit
.
dao
;
import
com.byit.job.model.MythJobFlowNodes
;
import
org.springframework.stereotype.Repository
;
/**
* @program: byit-myth-job->JobFlowNodes
* @description: 任务流节点的dao
* @author: huangfu
* @date: 2019/12/10 14:07
**/
@Repository
public
interface
JobFlowNodesMapper
{
/**
* 修改任务节点
* @param mythJobFlowNodes
*/
void
updateMythJobFlowNodes
(
MythJobFlowNodes
mythJobFlowNodes
);
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/dao/JobInfoMapper.java
View file @
844a728f
package
com
.
byit
.
dao
;
import
com.byit.job.model.MythJob
Info
;
import
com.byit.job.model.MythJob
ReadAhead
;
import
org.apache.ibatis.annotations.Param
;
import
org.springframework.stereotype.Repository
;
...
...
@@ -19,5 +19,11 @@ public interface JobInfoMapper {
* @param maxNextTime
* @return
*/
List
<
MythJobInfo
>
findMythJobInfoByTriggerNextTime
(
@Param
(
"maxNextTime"
)
long
maxNextTime
);
List
<
MythJobReadAhead
>
findMythJobReadAheadByTriggerNextTime
(
@Param
(
"maxNextTime"
)
long
maxNextTime
);
/**
* 删除已经被调度的任务根据ID
* @param id
*/
void
removeMythJobReadAheadById
(
@Param
(
"id"
)
String
id
);
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/JobFlightScheduleService.java
0 → 100644
View file @
844a728f
package
com
.
byit
.
service
;
import
com.byit.job.model.MythJobFlightSchedule
;
import
org.apache.ibatis.annotations.Param
;
import
java.util.List
;
/**
* @program: byit-myth-job->JobFlightScheduleService
* @description: 任务排期表
* @author: huangfu
* @date: 2019/12/10 16:33
**/
public
interface
JobFlightScheduleService
{
/**
* 批量向排期表添加任务数据
* @param mythJobFlightSchedules
*/
void
addJobFlightSchedules
(
List
<
MythJobFlightSchedule
>
mythJobFlightSchedules
);
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/Job
Info
Service.java
→
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/Job
ReadAhead
Service.java
View file @
844a728f
package
com
.
byit
.
service
;
import
com.byit.job.model.MythJob
Info
;
import
com.byit.job.model.MythJob
ReadAhead
;
import
java.util.List
;
/**
* @program: byit-myth-job->Job
Info
Service
* @program: byit-myth-job->Job
ReadAhead
Service
* @description: 任务节点操作
* @author: huangfu
* @date: 2019/12/9 15:17
**/
public
interface
Job
Info
Service
{
public
interface
Job
ReadAhead
Service
{
/**
* 查询7秒内要执行的数据
* @param maxNextTime
* @return
*/
List
<
MythJob
Info
>
findMythJobInfo
ByTriggerNextTime
(
long
maxNextTime
);
List
<
MythJob
ReadAhead
>
findMythJobReadAhead
ByTriggerNextTime
(
long
maxNextTime
);
/**
* 删除已经被调度的任务根据ID
* @param id
*/
void
removeMythJobReadAheadById
(
String
id
);
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/impl/JobFlightScheduleServiceImpl.java
0 → 100644
View file @
844a728f
package
com
.
byit
.
service
.
impl
;
import
com.byit.dao.JobFlightScheduleMapper
;
import
com.byit.job.model.MythJobFlightSchedule
;
import
com.byit.service.JobFlightScheduleService
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.stereotype.Service
;
import
java.util.List
;
/**
* @program: byit-myth-job->JobFlightScheduleServiceImpl
* @description: 任务排期表
* @author: huangfu
* @date: 2019/12/10 16:36
**/
@Service
public
class
JobFlightScheduleServiceImpl
implements
JobFlightScheduleService
{
@Autowired
private
JobFlightScheduleMapper
jobFlightScheduleMapper
;
@Override
public
void
addJobFlightSchedules
(
List
<
MythJobFlightSchedule
>
mythJobFlightSchedules
)
{
jobFlightScheduleMapper
.
addJobFlightSchedules
(
mythJobFlightSchedules
);
}
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/impl/Job
Info
ServiceImpl.java
→
byit-myth-core/myth-admin-core/src/main/java/com/byit/service/impl/Job
ReadAhead
ServiceImpl.java
View file @
844a728f
package
com
.
byit
.
service
.
impl
;
import
com.byit.dao.JobInfoMapper
;
import
com.byit.job.model.MythJob
Info
;
import
com.byit.service.Job
Info
Service
;
import
com.byit.job.model.MythJob
ReadAhead
;
import
com.byit.service.Job
ReadAhead
Service
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.stereotype.Service
;
...
...
@@ -10,19 +10,24 @@ import org.springframework.stereotype.Service;
import
java.util.List
;
/**
* @program: byit-myth-job->Job
Info
ServiceImpl
* @program: byit-myth-job->Job
ReadAhead
ServiceImpl
* @description: 任务节点操作的实现类
* @author: huangfu
* @date: 2019/12/9 15:18
**/
@Service
@Slf4j
public
class
Job
InfoServiceImpl
implements
JobInfo
Service
{
public
class
Job
ReadAheadServiceImpl
implements
JobReadAhead
Service
{
@Autowired
private
JobInfoMapper
jobInfoMapper
;
@Override
public
List
<
MythJobInfo
>
findMythJobInfoByTriggerNextTime
(
long
maxNextTime
)
{
return
jobInfoMapper
.
findMythJobInfoByTriggerNextTime
(
maxNextTime
);
public
List
<
MythJobReadAhead
>
findMythJobReadAheadByTriggerNextTime
(
long
maxNextTime
)
{
return
jobInfoMapper
.
findMythJobReadAheadByTriggerNextTime
(
maxNextTime
);
}
@Override
public
void
removeMythJobReadAheadById
(
String
id
)
{
jobInfoMapper
.
removeMythJobReadAheadById
(
id
);
}
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/thread/JobScheduleHelper.java
View file @
844a728f
package
com
.
byit
.
thread
;
import
cn.hutool.core.collection.CollectionUtil
;
import
cn.hutool.core.util.ObjectUtil
;
import
com.byit.job.model.MythJobInfo
;
import
com.byit.service.JobInfoService
;
import
com.byit.job.model.MythJobFlightSchedule
;
import
com.byit.job.model.MythJobReadAhead
;
import
com.byit.service.JobFlightScheduleService
;
import
com.byit.service.JobReadAheadService
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.beans.BeanUtils
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.stereotype.Component
;
import
javax.annotation.Resource
;
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.concurrent.TimeUnit
;
...
...
@@ -29,7 +31,10 @@ import java.util.concurrent.TimeUnit;
public
class
JobScheduleHelper
{
private
DataSource
dataSource
;
@Autowired
private
JobInfoService
jobInfoService
;
private
JobReadAheadService
jobReadAheadService
;
@Autowired
private
JobFlightScheduleService
jobFlightScheduleService
;
/**
* 读取任务节点的预读
...
...
@@ -70,15 +75,15 @@ public class JobScheduleHelper{
}
}
log
.
info
(
"---------------------init myth-job admin jobInfoThread success------------------------"
);
boolean
preReadSuc
=
false
;
while
(!
jobInfoThreadToStop
)
{
//开始去扫描任务节点
long
start
=
System
.
currentTimeMillis
();
Connection
conn
=
null
;
Boolean
connAutoCommit
=
null
;
PreparedStatement
preparedStatement
=
null
;
preReadSuc
=
true
;
boolean
preReadSuc
=
true
;
try
{
conn
=
dataSource
.
getConnection
();
...
...
@@ -92,17 +97,32 @@ public class JobScheduleHelper{
//行锁已经加上 后续处理
long
nowTime
=
System
.
currentTimeMillis
();
//开始寻找此时 不是暂停状态,而且七秒内即将运行的任务
List
<
MythJobInfo
>
mythJobInfoByTriggerNextTime
=
jobInfoService
.
findMythJobInfoByTriggerNextTime
(
nowTime
+
PRE_READ_MS
);
log
.
info
(
"----------------{}----------------"
,
mythJobInfoByTriggerNextTime
);
if
(
CollectionUtil
.
isNotEmpty
(
mythJobInfoByTriggerNextTime
)){
mythJobInfoByTriggerNextTime
.
forEach
(
e
->
System
.
out
.
println
(
e
));
List
<
MythJobReadAhead
>
mythJobReadAheadByTriggerNextTime
=
jobReadAheadService
.
findMythJobReadAheadByTriggerNextTime
(
nowTime
+
PRE_READ_MS
);
if
(
CollectionUtil
.
isNotEmpty
(
mythJobReadAheadByTriggerNextTime
)){
List
<
MythJobFlightSchedule
>
mythJobFlightSchedules
=
new
ArrayList
<>(
15
);
mythJobReadAheadByTriggerNextTime
.
forEach
(
mythJobInfo
->
{
log
.
info
(
"任务:{}"
,
mythJobInfo
);
/**
* 需要去检验当前任务的上级节点是否已经执行成功,没有执行,或者处于暂停状态则跳过该任务
* 大概思路,根据任务流id,从任务流执行回溯表查询该任务流的所有节点,查看上级节点是否已经执行成功
*/
{
MythJobFlightSchedule
mythJobFlightSchedule
=
new
MythJobFlightSchedule
();
BeanUtils
.
copyProperties
(
mythJobInfo
,
mythJobFlightSchedule
);
mythJobFlightSchedules
.
add
(
mythJobFlightSchedule
);
jobReadAheadService
.
removeMythJobReadAheadById
(
mythJobInfo
.
getId
());
}
});
jobFlightScheduleService
.
addJobFlightSchedules
(
mythJobFlightSchedules
);
}
else
{
log
.
info
(
"-------------------空轮转-------------------"
);
preReadSuc
=
false
;
}
}
catch
(
Exception
e
){
if
(!
jobInfoThreadToStop
){
e
.
printStackTrace
();
log
.
error
(
"------------------
----{}-----
-----------------"
,
e
.
getMessage
());
log
.
error
(
"------------------
扫描job_flow_nodes出现异常:{}
-----------------"
,
e
.
getMessage
());
}
}
finally
{
//commit
...
...
@@ -111,7 +131,7 @@ public class JobScheduleHelper{
conn
.
commit
();
}
catch
(
SQLException
e
){
if
(!
jobInfoThreadToStop
){
log
.
error
(
"----------------
------{}---
-------------------"
,
e
.
getMessage
());
log
.
error
(
"----------------
提交行锁错误:{}
-------------------"
,
e
.
getMessage
());
}
}
//设置提交状态恢复原来的值
...
...
@@ -119,7 +139,7 @@ public class JobScheduleHelper{
conn
.
setAutoCommit
(
connAutoCommit
);
}
catch
(
SQLException
e
){
if
(!
jobInfoThreadToStop
){
log
.
error
(
"------------------
----{}----
------------------"
,
e
.
getMessage
());
log
.
error
(
"------------------
设置为自动提交出错:{}
------------------"
,
e
.
getMessage
());
}
}
//关闭连接
...
...
@@ -127,7 +147,7 @@ public class JobScheduleHelper{
conn
.
close
();
}
catch
(
SQLException
e
){
if
(!
jobInfoThreadToStop
){
log
.
error
(
"----------------
------{}---
-------------------"
,
e
.
getMessage
());
log
.
error
(
"----------------
关闭数据库连接:{}
-------------------"
,
e
.
getMessage
());
}
}
}
...
...
@@ -138,7 +158,7 @@ public class JobScheduleHelper{
preparedStatement
.
close
();
}
catch
(
SQLException
e
)
{
if
(!
jobInfoThreadToStop
){
log
.
error
(
"---------------
-------
{}----------------------"
,
e
.
getMessage
());
log
.
error
(
"---------------
关闭执行器出错:
{}----------------------"
,
e
.
getMessage
());
}
}
}
...
...
@@ -149,6 +169,7 @@ public class JobScheduleHelper{
if
(
cost
<
1000
){
// 预读期:成功-每秒扫描一次;失败-跳过这段时间
try
{
//这个睡眠是空轮转时,延长睡眠时间,奖励CUP的使用频率
TimeUnit
.
MILLISECONDS
.
sleep
((
preReadSuc
?
1000
:
PRE_READ_MS
)
-
System
.
currentTimeMillis
()%
1000
);
}
catch
(
InterruptedException
e
)
{
if
(!
jobInfoThreadToStop
)
{
...
...
@@ -170,6 +191,6 @@ public class JobScheduleHelper{
}
public
static
void
main
(
String
[]
args
)
{
System
.
out
.
println
(
System
.
currentTimeMillis
()+
7
000
);
System
.
out
.
println
(
System
.
currentTimeMillis
()+
30
*
1
000
);
}
}
byit-myth-core/myth-admin-core/src/main/resources/mapper/JobFlightScheduleMapper.xml
0 → 100644
View file @
844a728f
<?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.dao.JobFlightScheduleMapper"
>
<insert
id=
"addJobFlightSchedules"
parameterType=
"com.byit.job.model.MythJobFlightSchedule"
>
INSERT INTO job_flight_schedule
(
id,
job_name,
executor_handler,
job_cron,
job_desc,
plugin_urls,
routing_strategy,
blocking_strategy,
callback_token,
gateway_token,
request_host,
request_port,
job_param,
job_type,
taskflow_id,
parent_id,
taskflow_is_cron,
alarm_email,
executor_timeout,
source_urls,
glue_source,
glue_remark,
glue_updatetime,
trigger_next_time,
author,
add_time,
update_time,
run_id
)
VALUES
<foreach
collection=
"mythJobFlightSchedule"
separator=
","
item=
"item"
index=
"index"
>
(
#{item.id},
#{item.jobName},
#{item.executorHandler},
#{item.jobCron},
#{item.jobDesc},
#{item.pluginUrls},
#{item.routingStrategy},
#{item.blockingStrategy},
#{item.callbackToken},
#{item.gatewayToken},
#{item.requestHost},
#{item.requestPort},
#{item.jobParam},
#{item.jobType},
#{item.taskFlowId},
#{item.parentId},
#{item.taskFlowIsCron},
#{item.alarmEmail},
#{item.executorTimeout},
#{item.sourceUrls},
#{item.glueSource},
#{item.glueRemark},
#{item.glueUpdateTime},
#{item.triggerNextTime},
#{item.author},
#{item.addTime},
#{item.updateTime},
#{item.runID}
)
</foreach>
</insert>
</mapper>
\ No newline at end of file
byit-myth-core/myth-admin-core/src/main/resources/mapper/JobFlowNodesMapper.xml
0 → 100644
View file @
844a728f
<?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.dao.JobFlowNodesMapper"
>
<update
id=
"updateMythJobFlowNodes"
parameterType=
"com.byit.job.model.MythJobFlowNodes"
>
UPDATE job_read_ahead
<trim
prefix=
"SET"
suffixOverrides=
","
>
<if
test=
"jobName != null and jobName != ''"
>
job_name=#{jobName},
</if>
<if
test=
"executorHandler != null and executorHandler != ''"
>
executor_handler=#{executorHandler},
</if>
<if
test=
"jobCron != null and jobCron != ''"
>
job_cron=#{jobCron},
</if>
<if
test=
"jobDesc != null and jobDesc != ''"
>
job_desc=#{jobDesc},
</if>
<if
test=
"pluginUrls != null and pluginUrls != ''"
>
plugin_urls=#{pluginUrls},
</if>
<if
test=
"routingStrategy != null and routingStrategy != ''"
>
routing_strategy=#{routingStrategy},
</if>
<if
test=
"blockingStrategy != null and blockingStrategy != ''"
>
blocking_strategy=#{blockingStrategy},
</if>
<if
test=
"callbackToken != null and callbackToken != ''"
>
callback_token=#{callbackToken},
</if>
<if
test=
"gatewayToken != null and gatewayToken != ''"
>
gateway_token = #{gatewayToken},
</if>
<if
test=
"requestHost != null and requestHost != ''"
>
request_host =#{requestHost},
</if>
<if
test=
"requestPort != null and requestPort != ''"
>
request_port =#{requestPort},
</if>
<if
test=
"jobParam != null and jobParam != ''"
>
job_param =#{jobParam},
</if>
<if
test=
"jobType != null and jobType != ''"
>
job_type=#{jobType},
</if>
<if
test=
"taskFlowId != null and taskFlowId != ''"
>
taskflow_id =#{taskFlowId},
</if>
<if
test=
"parentId != null and parentId != ''"
>
parent_id =#{parentId},
</if>
<if
test=
"taskFlowIsCron != null and taskFlowIsCron != ''"
>
taskflow_is_cron =#{taskFlowIsCron},
</if>
<if
test=
"alarmEmail != null and alarmEmail != ''"
>
alarm_email =#{alarmEmail},
</if>
<if
test=
"executorTimeout != null"
>
executor_timeout =#{executorTimeout},
</if>
<if
test=
"sourceUrls != null and sourceUrls != ''"
>
source_urls =#{sourceUrls},
</if>
<if
test=
"glueSource != null and glueSource != ''"
>
glue_source =#{glueSource},
</if>
<if
test=
"glueRemark != null and glueRemark != ''"
>
glue_remark =#{glueRemark},
</if>
<if
test=
"glueUpdateTime != null"
>
glue_updatetime =#{glueUpdateTime},
</if>
<if
test=
"triggerNextTime != null"
>
trigger_next_time =#{triggerNextTime},
</if>
<if
test=
"author != null and author != ''"
>
author=#{author},
</if>
<if
test=
"addTime != null"
>
add_time =#{addTime},
</if>
<if
test=
"updateTime != null"
>
update_time =#{updateTime},
</if>
<if
test=
"removalMark != null"
>
removal_mark =#{removalMark},
</if>
<if
test=
"repeatTimes != null"
>
repeat_times =#{repeatTimes},
</if>
</trim>
WHERE id=#{id}
</update>
</mapper>
\ No newline at end of file
byit-myth-core/myth-admin-core/src/main/resources/mapper/MythJobInfoMapper.xml
View file @
844a728f
...
...
@@ -30,11 +30,10 @@
t.author,
t.add_time,
t.update_time,
t.removal_mark,
t.repeat_times
t.run_id
</sql>
<resultMap
id=
"MythJobInfo"
type=
"com.byit.job.model.MythJob
Info
"
>
<resultMap
id=
"MythJobInfo"
type=
"com.byit.job.model.MythJob
ReadAhead
"
>
<id
column=
"id"
property=
"id"
/>
<result
column=
"job_name"
property=
"jobName"
/>
<result
column=
"executor_handler"
property=
"executorHandler"
/>
...
...
@@ -63,14 +62,18 @@
<result
column=
"author"
property=
"author"
/>
<result
column=
"add_time"
property=
"addTime"
/>
<result
column=
"update_time"
property=
"updateTime"
/>
<result
column=
"removal_mark"
property=
"removalMark"
/>
<result
column=
"repeat_times"
property=
"repeatTimes"
/>
<result
column=
"run_id"
property=
"runID"
/>
</resultMap>
<select
id=
"findMythJob
Info
ByTriggerNextTime"
resultMap=
"MythJobInfo"
>
<select
id=
"findMythJob
ReadAhead
ByTriggerNextTime"
resultMap=
"MythJobInfo"
>
SELECT
<include
refid=
"Base_Column_List"
/>
FROM job_
info
AS t
FROM job_
read_ahead
AS t
WHERE t.trigger_status=1
and
t.trigger_next_time
<![CDATA[ <= ]]>
#{maxNextTime}
AND
t.trigger_next_time
<![CDATA[ <= ]]>
#{maxNextTime}
</select>
<delete
id=
"removeMythJobReadAheadById"
>
DELETE FROM job_read_ahead WHERE id=#{id}
</delete>
</mapper>
\ No newline at end of file
byit-myth-core/myth-core-common/src/main/java/com/byit/job/handler/impl/GlueJobHandler.java
View file @
844a728f
...
...
@@ -2,7 +2,7 @@ package com.byit.job.handler.impl;
import
com.byit.job.handler.BaseJobHandler
;
import
com.byit.job.handler.interfaces.IJobHandler
;
import
com.byit.job.
model
.ReturnResult
;
import
com.byit.job.
vo
.ReturnResult
;
/**
* @program: byit-myth-job->GlueJobHandler
...
...
byit-myth-core/myth-core-common/src/main/java/com/byit/job/handler/impl/ScriptJobHandler.java
View file @
844a728f
...
...
@@ -2,7 +2,7 @@ package com.byit.job.handler.impl;
import
com.byit.job.enums.GlueTypeEnum
;
import
com.byit.job.handler.BaseJobHandler
;
import
com.byit.job.
model
.ReturnResult
;
import
com.byit.job.
vo
.ReturnResult
;
/**
* @program: byit-myth-job->ScriptJobHandler
...
...
byit-myth-core/myth-core-common/src/main/java/com/byit/job/handler/interfaces/IJobHandler.java
View file @
844a728f
package
com
.
byit
.
job
.
handler
.
interfaces
;
import
com.byit.job.
model
.ReturnResult
;
import
com.byit.job.
vo
.ReturnResult
;
/**
* @program: byit-myth-job->IJobHandler
...
...
byit-myth-core/myth-core-common/src/main/java/com/byit/job/model/MythJob
Info
.java
→
byit-myth-core/myth-core-common/src/main/java/com/byit/job/model/MythJob
FlightSchedule
.java
View file @
844a728f
...
...
@@ -8,16 +8,16 @@ import lombok.ToString;
import
java.util.Date
;
/**
* @program: byit-myth-job->MythJob
Info
* @description: 任务
节点映射实体
* @program: byit-myth-job->MythJob
FlightSchedule
* @description: 任务
排期表
* @author: huangfu
* @date: 2019/12/
9 10:18
* @date: 2019/12/
10 15:05
**/
@Data
@AllArgsConstructor
@NoArgsConstructor
@ToString
public
class
MythJob
Info
{
public
class
MythJob
FlightSchedule
{
/**
* 任务节点的id
*/
...
...
@@ -100,16 +100,12 @@ public class MythJobInfo {
/**
* 任务的超时时间
*/
private
int
executorTimeout
;
private
Integer
executorTimeout
;
/**
* 脚本文件地址(或远程文件服务器地址)支持多个,逗号分割
*/
private
String
sourceUrls
;
/**
* 调度状态:0-暂停,1-运行
*/
private
String
triggerStatus
;
/**
*调度任务源码
*/
private
String
glueSource
;
...
...
@@ -124,7 +120,7 @@ public class MythJobInfo {
/**
* 任务下次调度时间
*/
private
l
ong
triggerNextTime
;
private
L
ong
triggerNextTime
;
/**
* 任务创建者
*/
...
...
@@ -138,11 +134,7 @@ public class MythJobInfo {
*/
private
Date
updateTime
;
/**
* 删除标志
*/
private
String
removalMark
;
/**
* 重复次数 -1遵循cron表达式解析
* 运行标识
*/
private
int
repeatTimes
;
private
String
runID
;
}
byit-myth-core/myth-core-common/src/main/java/com/byit/job/model/MythJobFlowNodes.java
0 → 100644
View file @
844a728f
package
com
.
byit
.
job
.
model
;
import
lombok.AllArgsConstructor
;
import
lombok.Data
;
import
lombok.NoArgsConstructor
;
import
lombok.ToString
;
import
java.util.Date
;
/**
* @program: byit-myth-job->MythJobFlowNodes
* @description: 任务流节点,它作为源数据表而言,所有数据都不应该参与修改操作
* @author: huangfu
* @date: 2019/12/10 14:01
**/
@Data
@AllArgsConstructor
@NoArgsConstructor
@ToString
public
class
MythJobFlowNodes
{
/**
* 任务节点的id
*/
private
String
id
;
/**
* 任务节点的名字
*/
private
String
jobName
;
/**
* 插件任务调度的key,调度中心会根据这个key找到对应的插件端任务,执行
*/
private
String
executorHandler
;
/**
* 任务周期调度cron表达式
*/
private
String
jobCron
;
/**
* 任务详情,展示在调度中心平台的备注
*/
private
String
jobDesc
;
/**
* 插件地址的集合
*/
private
String
pluginUrls
;
/**
* 路由策略
* 1随机(默认)
* 2轮询
* 3最近最少使用
* 4最近最久未使用算法
* 5哈希算法
*/
private
String
routingStrategy
;
/**
* 阻塞策略:
* 1丢弃
* 2阻塞等待(默认)
*/
private
String
blockingStrategy
;
/**
* 回调时的身份认证
*/
private
String
callbackToken
;
/**
* 请求调度中心的token
*/
private
String
gatewayToken
;
/**
* 调度中心主机名
*/
private
String
requestHost
;
/**
* 调度中心端口号
*/
private
String
requestPort
;
/**
* 任务参数
*/
private
String
jobParam
;
/**
* 任务类型,java,python,php,script,sql.shell
*/
private
String
jobType
;
/**
* 所属任务流的id
*/
private
String
taskFlowId
;
/**
* 上级节点
*/
private
String
parentId
;
/**
* 是否跟随任务流的时间设置?1不跟随(默认),2跟随
*/
private
String
taskFlowIsCron
;
/**
* 报警邮件
*/
private
String
alarmEmail
;
/**
* 任务的超时时间
*/
private
Integer
executorTimeout
;
/**
* 脚本文件地址(或远程文件服务器地址)支持多个,逗号分割
*/
private
String
sourceUrls
;
/**
*调度任务源码
*/
private
String
glueSource
;
/**
* 源码备注
*/
private
String
glueRemark
;
/**
* 源码的修改时间
*/
private
Date
glueUpdateTime
;
/**
* 任务创建者
*/
private
String
author
;
/**
* 任务节点添加时间
*/
private
Date
addTime
;
/**
* 任务节点修改时间
*/
private
Date
updateTime
;
/**
* 删除标志 1正常 2删除
*/
private
String
removalMark
;
/**
* 重复次数 -1永久运行
*/
private
Integer
repeatTimes
;
}
byit-myth-core/myth-core-common/src/main/java/com/byit/job/model/MythJobReadAhead.java
0 → 100644
View file @
844a728f
package
com
.
byit
.
job
.
model
;
import
lombok.AllArgsConstructor
;
import
lombok.Data
;
import
lombok.NoArgsConstructor
;
import
lombok.ToString
;
import
java.util.Date
;
/**
* @program: byit-myth-job->MythJobReadAhead
* @description: 任务预读表映射实体,他作为任务流节点的快照,其实所有的修改操作都应该在这个类上做
* @author: huangfu
* @date: 2019/12/9 10:18
**/
@Data
@AllArgsConstructor
@NoArgsConstructor
@ToString
public
class
MythJobReadAhead
{
/**
* 任务节点的id
*/
private
String
id
;
/**
* 任务节点的名字
*/
private
String
jobName
;
/**
* 插件任务调度的key,调度中心会根据这个key找到对应的插件端任务,执行
*/
private
String
executorHandler
;
/**
* 任务周期调度cron表达式
*/
private
String
jobCron
;
/**
* 任务详情,展示在调度中心平台的备注
*/
private
String
jobDesc
;
/**
* 插件地址的集合
*/
private
String
pluginUrls
;
/**
* 路由策略
* 1随机(默认)
* 2轮询
* 3最近最少使用
* 4最近最久未使用算法
* 5哈希算法
*/
private
String
routingStrategy
;
/**
* 阻塞策略:
* 1丢弃
* 2阻塞等待(默认)
*/
private
String
blockingStrategy
;
/**
* 回调时的身份认证
*/
private
String
callbackToken
;
/**
* 请求调度中心的token
*/
private
String
gatewayToken
;
/**
* 调度中心主机名
*/
private
String
requestHost
;
/**
* 调度中心端口号
*/
private
String
requestPort
;
/**
* 任务参数
*/
private
String
jobParam
;
/**
* 任务类型,java,python,php,script,sql.shell
*/
private
String
jobType
;
/**
* 所属任务流的id
*/
private
String
taskFlowId
;
/**
* 上级节点
*/
private
String
parentId
;
/**
* 是否跟随任务流的时间设置?1不跟随(默认),2跟随
*/
private
String
taskFlowIsCron
;
/**
* 报警邮件
*/
private
String
alarmEmail
;
/**
* 任务的超时时间
*/
private
Integer
executorTimeout
;
/**
* 脚本文件地址(或远程文件服务器地址)支持多个,逗号分割
*/
private
String
sourceUrls
;
/**
* 调度状态:0-暂停,1-运行
*/
private
String
triggerStatus
;
/**
*调度任务源码
*/
private
String
glueSource
;
/**
* 源码备注
*/
private
String
glueRemark
;
/**
* 源码的修改时间
*/
private
Date
glueUpdateTime
;
/**
* 任务下次调度时间
*/
private
Long
triggerNextTime
;
/**
* 任务创建者
*/
private
String
author
;
/**
* 任务节点添加时间
*/
private
Date
addTime
;
/**
* 任务节点修改时间
*/
private
Date
updateTime
;
/**
* 任务流的唯一运行标识
*/
private
String
runID
;
}
byit-myth-core/myth-core-common/src/main/java/com/byit/job/
model
/ReturnResult.java
→
byit-myth-core/myth-core-common/src/main/java/com/byit/job/
vo
/ReturnResult.java
View file @
844a728f
package
com
.
byit
.
job
.
model
;
package
com
.
byit
.
job
.
vo
;
import
com.byit.job.enums.JobResultEnum
;
import
lombok.AllArgsConstructor
;
...
...
byit-myth-executor/myth-exector-plugin/src/main/java/com/byit/rpc/RunJobServerHandler.java
View file @
844a728f
...
...
@@ -3,7 +3,7 @@ package com.byit.rpc;
import
com.alibaba.fastjson.JSON
;
import
com.alibaba.fastjson.JSONObject
;
import
com.byit.job.handler.interfaces.IJobHandler
;
import
com.byit.job.
model
.ReturnResult
;
import
com.byit.job.
vo
.ReturnResult
;
import
com.byit.utils.JobUtils
;
import
io.netty.buffer.ByteBuf
;
import
io.netty.buffer.Unpooled
;
...
...
demo-client/byit-demo-client/src/main/java/com/byit/job/DemoJob.java
View file @
844a728f
...
...
@@ -2,7 +2,7 @@ package com.byit.job;
import
com.byit.annotations.JobHandler
;
import
com.byit.job.handler.BaseJobHandler
;
import
com.byit.job.
model
.ReturnResult
;
import
com.byit.job.
vo
.ReturnResult
;
/**
* @program: byit-myth-job->DemoJob
...
...
demo-client/byit-demo-client/src/main/java/com/byit/job/DemoJob1.java
View file @
844a728f
...
...
@@ -2,7 +2,7 @@ package com.byit.job;
import
com.byit.annotations.JobHandler
;
import
com.byit.job.handler.BaseJobHandler
;
import
com.byit.job.
model
.ReturnResult
;
import
com.byit.job.
vo
.ReturnResult
;
/**
* @program: byit-myth-job->DemoJob
...
...
doc/任务调度插件开发流程说明.docx
View file @
844a728f
No preview for this file type
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