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
7d4d61ff
Commit
7d4d61ff
authored
Dec 16, 2019
by
huangfusuper
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
完善调度中心与插件端的通讯数据交换!添加任务回调,传递任务信息
parent
c55f178b
Show whitespace changes
Inline
Side-by-side
Showing
23 changed files
with
204 additions
and
120 deletions
+204
-120
JobController.java
...dmin/src/main/java/com/byit/controller/JobController.java
+4
-6
WorkRoulette.java
...h-admin-core/src/main/java/com/byit/job/WorkRoulette.java
+0
-3
JavaBeanJobTask.java
...min-core/src/main/java/com/byit/task/JavaBeanJobTask.java
+11
-22
JobScheduleHelper.java
...core/src/main/java/com/byit/thread/JobScheduleHelper.java
+2
-12
AdminSenPluginDto.java
...mon/src/main/java/com/byit/job/dto/AdminSenPluginDto.java
+39
-0
PluginBeanJobInfo.java
...mon/src/main/java/com/byit/job/dto/PluginBeanJobInfo.java
+1
-1
PluginJobRunResultDto.java
...src/main/java/com/byit/job/dto/PluginJobRunResultDto.java
+39
-0
IpUtil.java
...-core-common/src/main/java/com/byit/job/utils/IpUtil.java
+7
-6
SourceObj2TargetObjUtil.java
...main/java/com/byit/job/utils/SourceObj2TargetObjUtil.java
+1
-2
RunJobServerHandler.java
...lugin/src/main/java/com/byit/rpc/RunJobServerHandler.java
+12
-43
RunJobThread.java
...ector-plugin/src/main/java/com/byit/rpc/RunJobThread.java
+57
-0
JobUtils.java
...exector-plugin/src/main/java/com/byit/utils/JobUtils.java
+1
-1
JobOperating.java
...r-api/src/main/java/com/byit/plugin/api/JobOperating.java
+1
-1
Mains.java
...nt/byit-demo-client/src/main/java/com/byit/job/Mains.java
+1
-1
flow_version.sql
doc/sql/myth_job/flow_version.sql
+1
-1
job_flight_schedule.sql
doc/sql/myth_job/job_flight_schedule.sql
+6
-6
job_flow.sql
doc/sql/myth_job/job_flow.sql
+1
-1
job_flow_nodes.sql
doc/sql/myth_job/job_flow_nodes.sql
+5
-5
job_flow_snapshot.sql
doc/sql/myth_job/job_flow_snapshot.sql
+1
-1
job_lock.sql
doc/sql/myth_job/job_lock.sql
+7
-1
job_log.sql
doc/sql/myth_job/job_log.sql
+1
-1
job_logglue.sql
doc/sql/myth_job/job_logglue.sql
+1
-1
job_read_ahead.sql
doc/sql/myth_job/job_read_ahead.sql
+5
-5
No files found.
byit-myth-admin/src/main/java/com/byit/controller/JobController.java
View file @
7d4d61ff
package
com
.
byit
.
controller
;
import
com.byit.job.WorkRoulette
;
import
com.byit.job.model.MythJobFlightSchedule
;
import
com.byit.job.dto.PluginJobRunResultDto
;
import
com.byit.job.model.MythJobReadAhead
;
import
com.byit.job.
model
.PluginBeanJobInfo
;
import
com.byit.job.
dto
.PluginBeanJobInfo
;
import
com.byit.job.utils.SourceObj2TargetObjUtil
;
import
com.byit.job.vo.ReturnResult
;
import
com.byit.service.JobReadAheadService
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.web.bind.annotation.*
;
import
javax.sql.DataSource
;
import
java.util.List
;
/**
...
...
@@ -32,8 +30,8 @@ public class JobController {
}
@PostMapping
(
value
=
"callbackRes"
)
public
String
callbackRes
(
@RequestBody
ReturnResult
<
String
>
result
){
System
.
out
.
println
(
result
.
getCode
()+
"-----"
+
result
.
getMsg
());
public
String
callbackRes
(
@RequestBody
PluginJobRunResultDto
pluginJobRunResultDto
){
System
.
out
.
println
(
pluginJobRunResultDto
.
getReturnResult
().
getCode
()+
"-----"
+
pluginJobRunResultDto
.
getReturnResult
()
.
getMsg
());
return
"好的,我知道你执行成功了"
;
}
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/job/WorkRoulette.java
View file @
7d4d61ff
package
com
.
byit
.
job
;
import
com.byit.job.model.MythJobFlightSchedule
;
import
com.byit.job.model.PluginBeanJobInfo
;
import
com.byit.task.JavaBeanJobTask
;
import
io.netty.util.HashedWheelTimer
;
import
io.netty.util.TimerTask
;
...
...
byit-myth-core/myth-admin-core/src/main/java/com/byit/task/JavaBeanJobTask.java
View file @
7d4d61ff
...
...
@@ -2,7 +2,7 @@ package com.byit.task;
import
cn.hutool.http.HttpUtil
;
import
com.alibaba.fastjson.JSON
;
import
com.
alibaba.fastjson.JSONObject
;
import
com.
byit.job.dto.AdminSenPluginDto
;
import
com.byit.job.exceptions.plugin.PluginException
;
import
com.byit.job.model.MythJobFlightSchedule
;
import
com.byit.job.utils.IpUtil
;
...
...
@@ -33,30 +33,19 @@ public class JavaBeanJobTask implements TimerTask {
String
url
=
IpUtil
.
electiveUrl
(
mythJobFlightSchedule
.
getPluginUrls
(
),
rpcInvokerRouter
);
String
jobHandelName
=
mythJobFlightSchedule
.
getExecutorHandler
(
);
String
param
=
mythJobFlightSchedule
.
getJobParam
(
);
JSONObject
jsonObject
=
new
JSONObject
();
jsonObject
.
put
(
"jobHandelName"
,
jobHandelName
);
jsonObject
.
put
(
"callbackMethod"
,
"http://127.0.0.1:8080/job/callbackRes"
);
jsonObject
.
put
(
"param"
,
param
);
String
result
=
HttpUtil
.
post
(
url
,
JSON
.
toJSONString
(
jsonObject
),
10
*
1000
);
log
.
info
(
"---------------{}------------"
,
result
);
String
runId
=
mythJobFlightSchedule
.
getRunID
();
AdminSenPluginDto
adminSenPluginDto
=
new
AdminSenPluginDto
();
adminSenPluginDto
.
setCallbackUrl
(
"http://127.0.0.1:8080/job/callbackRes"
);
adminSenPluginDto
.
setJobHandelName
(
jobHandelName
);
adminSenPluginDto
.
setJobParam
(
param
);
adminSenPluginDto
.
setRunId
(
runId
);
String
result
=
HttpUtil
.
post
(
url
,
JSON
.
toJSONString
(
adminSenPluginDto
),
10
*
1000
);
log
.
debug
(
"---------------{}------------"
,
result
);
}
catch
(
PluginException
ignored
){
log
.
error
(
"-------------
-------{},{}-----------------
"
,
ignored
.
getIEnum
().
getCode
(),
ignored
.
getIEnum
().
getMsg
());
log
.
error
(
"-------------
通讯异常:{},{}
"
,
ignored
.
getIEnum
().
getCode
(),
ignored
.
getIEnum
().
getMsg
());
}
catch
(
Exception
e
){
e
.
printStackTrace
();
}
}
public
static
void
main
(
String
[]
args
)
{
JSONObject
jsonObject
=
new
JSONObject
();
jsonObject
.
put
(
"heartbeat"
,
"PENG"
);
//jsonObject.put("jobHandelName","addJob");
//jsonObject.put("param","asd");
long
start
=
System
.
currentTimeMillis
(
);
HttpUtil
.
post
(
"http://10.0.55.200:8888"
,
JSON
.
toJSONString
(
jsonObject
));
System
.
out
.
println
((
System
.
currentTimeMillis
(
)-
start
)
);
}
}
byit-myth-core/myth-admin-core/src/main/java/com/byit/thread/JobScheduleHelper.java
View file @
7d4d61ff
package
com
.
byit
.
thread
;
import
cn.hutool.core.collection.CollectionUtil
;
import
com.alibaba.fastjson.JSON
;
import
com.alibaba.fastjson.annotation.JSONField
;
import
com.alibaba.fastjson.util.TypeUtils
;
import
com.byit.job.WorkRoulette
;
import
com.byit.job.model.MythJobFlightSchedule
;
import
com.byit.job.model.MythJobReadAhead
;
import
com.byit.service.JobFlightScheduleService
;
import
com.byit.service.JobReadAheadService
;
import
com.byit.task.JavaBeanJobTask
;
import
lombok.AllArgsConstructor
;
import
lombok.Data
;
import
lombok.NoArgsConstructor
;
import
lombok.extern.slf4j.Slf4j
;
import
org.apache.commons.lang3.StringUtils
;
import
org.springframework.beans.BeanUtils
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.beans.factory.annotation.Value
;
import
org.springframework.stereotype.Component
;
import
javax.sql.DataSource
;
...
...
@@ -26,7 +19,6 @@ import java.sql.PreparedStatement;
import
java.sql.SQLException
;
import
java.util.ArrayList
;
import
java.util.List
;
import
java.util.Map
;
import
java.util.concurrent.TimeUnit
;
/**
...
...
@@ -115,9 +107,10 @@ public class JobScheduleHelper{
/**
* 需要去检验当前任务的上级节点是否已经执行成功,没有执行,或者处于暂停状态则跳过该任务
* 大概思路,根据任务流id,从任务流执行回溯表查询该任务流的所有节点,查看上级节点是否已经执行成功
* //TODO 需要修改 判断父节点是否执行完毕 注意 父节点是一个集合
*/
if
(
StringUtils
.
isNotBlank
(
mythJobInfo
.
getParentId
())){
log
.
info
(
"任务:{}"
,
mythJobInfo
);
log
.
debug
(
"任务:{}"
,
mythJobInfo
);
MythJobFlightSchedule
mythJobFlightSchedule
=
new
MythJobFlightSchedule
();
BeanUtils
.
copyProperties
(
mythJobInfo
,
mythJobFlightSchedule
);
mythJobFlightSchedules
.
add
(
mythJobFlightSchedule
);
...
...
@@ -323,7 +316,4 @@ public class JobScheduleHelper{
this
.
dataSource
=
dataSource
;
}
public
static
void
main
(
String
[]
args
)
{
System
.
out
.
println
(
System
.
currentTimeMillis
()+
30
*
1000
);
}
}
byit-myth-core/myth-core-common/src/main/java/com/byit/job/dto/AdminSenPluginDto.java
0 → 100644
View file @
7d4d61ff
package
com
.
byit
.
job
.
dto
;
import
lombok.AllArgsConstructor
;
import
lombok.Data
;
import
lombok.NoArgsConstructor
;
import
lombok.ToString
;
/**
* @program: byit-myth-job->AdminSenPluginDto
* @description: 调度中心,通讯插件任务的节点承载
* @author: huangfu
* @date: 2019/12/16 11:19
**/
@Data
@AllArgsConstructor
@NoArgsConstructor
@ToString
public
class
AdminSenPluginDto
{
/**
* 运行标识
*/
private
String
runId
;
/**
* 插件端的任务标识
*/
private
String
jobHandelName
;
/**
* 任务的参数
*/
private
String
jobParam
;
/**
* 回调URL
*/
private
String
callbackUrl
;
/**
* 检测是否有心跳参数,有心跳参数则为测试参数,且为PENG的话,服务端回复 PONG
*/
private
String
heartbeat
;
}
byit-myth-core/myth-core-common/src/main/java/com/byit/job/
model
/PluginBeanJobInfo.java
→
byit-myth-core/myth-core-common/src/main/java/com/byit/job/
dto
/PluginBeanJobInfo.java
View file @
7d4d61ff
package
com
.
byit
.
job
.
model
;
package
com
.
byit
.
job
.
dto
;
import
lombok.AllArgsConstructor
;
import
lombok.Data
;
...
...
byit-myth-core/myth-core-common/src/main/java/com/byit/job/dto/PluginJobRunResultDto.java
0 → 100644
View file @
7d4d61ff
package
com
.
byit
.
job
.
dto
;
import
com.byit.job.vo.ReturnResult
;
import
lombok.AllArgsConstructor
;
import
lombok.Data
;
import
lombok.NoArgsConstructor
;
import
lombok.ToString
;
import
java.util.Date
;
/**
* @program: byit-myth-job->PluginJobRunResult
* @description: 这个是插件端的执行情况通知到调度中心的数据承载
* @author: huangfu
* @date: 2019/12/16 11:04
**/
@Data
@NoArgsConstructor
@AllArgsConstructor
@ToString
public
class
PluginJobRunResultDto
{
/**
* 任务的运行标识
*/
private
String
jobRunId
;
/**
* 任务的执行结果
*/
private
ReturnResult
<
String
>
returnResult
;
/**
* 调度时间
*/
private
Date
startTime
;
/**
* 结束时间
*/
private
Date
endTime
;
}
byit-myth-core/myth-core-common/src/main/java/com/byit/job/utils/IpUtil.java
View file @
7d4d61ff
...
...
@@ -4,6 +4,7 @@ import cn.hutool.core.collection.CollectionUtil;
import
cn.hutool.http.HttpUtil
;
import
com.alibaba.fastjson.JSON
;
import
com.alibaba.fastjson.JSONObject
;
import
com.byit.job.dto.AdminSenPluginDto
;
import
com.byit.job.enums.plugin.PluginEnum
;
import
com.byit.job.exceptions.plugin.PluginException
;
import
com.byit.rpc.remoting.invoker.route.RpcLoadBalance
;
...
...
@@ -27,11 +28,11 @@ public class IpUtil {
*/
public
static
boolean
survivalTest
(
String
url
){
try
{
log
.
info
(
"------------------开始发送心跳包--------------------"
);
JSONObject
jsonObject
=
new
JSONObject
(
);
jsonObject
.
put
(
"heartbeat"
,
"PENG"
);
String
heartbeatRes
=
HttpUtil
.
post
(
url
,
JSON
.
toJSONString
(
jsonObject
),
2
*
1000
);
log
.
info
(
"------------------接收心跳包--------------------"
);
log
.
debug
(
"------------------开始发送心跳包--------------------"
);
AdminSenPluginDto
adminSenPluginDto
=
new
AdminSenPluginDto
(
);
adminSenPluginDto
.
setHeartbeat
(
"PENG"
);
String
heartbeatRes
=
HttpUtil
.
post
(
url
,
JSON
.
toJSONString
(
adminSenPluginDto
),
2
*
1000
);
log
.
debug
(
"------------------接收心跳包--------------------"
);
return
"PONG"
.
equals
(
heartbeatRes
);
}
catch
(
Exception
e
){
log
.
error
(
"--------------------{},服务不可用------------------"
,
url
);
...
...
@@ -47,7 +48,7 @@ public class IpUtil {
* @return
*/
public
static
String
electiveUrl
(
String
urlsStr
,
RpcLoadBalance
rpcLoadBalance
)
throws
PluginException
{
log
.
info
(
"-----------------开始根据路由规则选取服务地址--------------------"
);
log
.
debug
(
"-----------------开始根据路由规则选取服务地址--------------------"
);
String
[]
urls
=
urlsStr
.
split
(
","
);
TreeSet
<
String
>
treeSet
=
new
TreeSet
<>(
Arrays
.
asList
(
urls
));
boolean
flag
=
false
;
...
...
byit-myth-core/myth-core-common/src/main/java/com/byit/job/utils/SourceObj2TargetObjUtil.java
View file @
7d4d61ff
package
com
.
byit
.
job
.
utils
;
import
com.byit.job.model.MythJobReadAhead
;
import
com.byit.job.
model
.PluginBeanJobInfo
;
import
com.byit.job.
dto
.PluginBeanJobInfo
;
import
java.text.ParseException
;
import
java.util.Date
;
import
java.util.UUID
;
...
...
byit-myth-executor/myth-exector-plugin/src/main/java/com/byit/rpc/RunJobServerHandler.java
View file @
7d4d61ff
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.vo.ReturnResult
;
import
com.byit.utils.JobUtils
;
import
com.byit.job.dto.AdminSenPluginDto
;
import
io.netty.buffer.ByteBuf
;
import
io.netty.buffer.Unpooled
;
import
io.netty.channel.ChannelHandlerContext
;
...
...
@@ -28,7 +25,7 @@ public class RunJobServerHandler extends SimpleChannelInboundHandler<FullHttpReq
/**
* LinkedBlockingQueue 不指定容量就变成了无界队列
*/
private
static
final
ThreadPoolExecutor
jobTriggerPool
=
new
ThreadPoolExecutor
(
private
static
final
ThreadPoolExecutor
JOB_TRIGGER_POOL
=
new
ThreadPoolExecutor
(
50
,
200
,
60L
,
...
...
@@ -40,20 +37,21 @@ public class RunJobServerHandler extends SimpleChannelInboundHandler<FullHttpReq
private
static
final
String
PONG
=
"PONG"
;
@Override
protected
void
channelRead0
(
ChannelHandlerContext
ctx
,
FullHttpRequest
req
)
throws
Exception
{
log
.
info
(
"----------------------有请求过来了------------------------"
);
log
.
debug
(
"----------------------有请求过来了------------------------"
);
String
heartbeatResponseBody
=
"FAILURE"
;
if
(
req
instanceof
HttpRequest
){
JSONObject
jsonObject
=
analysisParam
(
req
.
content
(
));
if
(
null
==
jsonObject
){
if
(
req
!=
null
){
//解析 调度中心 的数据对象
AdminSenPluginDto
adminSenPluginDto
=
analysisParam
(
req
.
content
(
));
if
(
null
==
adminSenPluginDto
){
throw
new
Exception
(
"核心参数 为 null"
);
}
//检测是否有心跳参数,有心跳参数则为测试参数,且为PENG的话,服务端回复 PONG
String
heartbeat
=
(
String
)(
jsonObject
.
get
(
"heartbeat"
)
);
String
heartbeat
=
adminSenPluginDto
.
getHeartbeat
(
);
if
(
null
==
heartbeat
){
jobTriggerPool
.
execute
(
new
RunJobThread
(
jsonObject
));
JOB_TRIGGER_POOL
.
execute
(
new
RunJobThread
(
adminSenPluginDto
));
heartbeatResponseBody
=
"SUCCESS"
;
}
else
if
(
PENG
.
equals
(
heartbeat
)){
log
.
info
(
"----------调度平台心跳检测-------------"
);
log
.
debug
(
"----------调度平台心跳检测-------------"
);
heartbeatResponseBody
=
PONG
;
}
//----------------------------------消息发送-----------------------------------
...
...
@@ -75,8 +73,8 @@ public class RunJobServerHandler extends SimpleChannelInboundHandler<FullHttpReq
* @param byteBuf
* @return
*/
private
JSONObject
analysisParam
(
ByteBuf
byteBuf
){
return
JSON
.
parseObject
(
byteBuf
.
toString
(
CharsetUtil
.
UTF_8
)
);
private
AdminSenPluginDto
analysisParam
(
ByteBuf
byteBuf
){
return
JSON
.
parseObject
(
byteBuf
.
toString
(
CharsetUtil
.
UTF_8
)
,
AdminSenPluginDto
.
class
);
}
/**
...
...
@@ -91,36 +89,7 @@ public class RunJobServerHandler extends SimpleChannelInboundHandler<FullHttpReq
ctx
.
close
();
}
}
@Slf4j
class
RunJobThread
implements
Runnable
{
private
JSONObject
jsonObject
;
public
RunJobThread
(
JSONObject
jsonObject
)
{
this
.
jsonObject
=
jsonObject
;
}
@Override
public
void
run
()
{
String
jobHandelName
=
((
String
)(
jsonObject
.
get
(
"jobHandelName"
)));
String
param
=
((
String
)(
jsonObject
.
get
(
"param"
)));
ReturnResult
<
String
>
stringReturnResult
=
runJob
(
jobHandelName
,
param
);
String
callbackMethod
=
(
String
)
jsonObject
.
get
(
"callbackMethod"
);
String
post
=
cn
.
hutool
.
http
.
HttpUtil
.
post
(
callbackMethod
,
JSON
.
toJSONString
(
stringReturnResult
));
log
.
info
(
"---------服务器端:{}:{}-------------"
,
callbackMethod
,
post
);
}
private
ReturnResult
<
String
>
runJob
(
String
jobHandlerName
,
String
param
){
Class
<?
extends
IJobHandler
>
jobClass
=
JobUtils
.
jobCache
.
get
(
jobHandlerName
);
try
{
IJobHandler
iJobHandler
=
jobClass
.
newInstance
(
);
return
iJobHandler
.
execute
(
param
);
}
catch
(
Exception
e
)
{
e
.
printStackTrace
(
);
}
return
null
;
}
}
byit-myth-executor/myth-exector-plugin/src/main/java/com/byit/rpc/RunJobThread.java
0 → 100644
View file @
7d4d61ff
package
com
.
byit
.
rpc
;
import
com.alibaba.fastjson.JSON
;
import
com.byit.job.dto.AdminSenPluginDto
;
import
com.byit.job.dto.PluginJobRunResultDto
;
import
com.byit.job.handler.interfaces.IJobHandler
;
import
com.byit.job.vo.ReturnResult
;
import
com.byit.utils.JobUtils
;
import
lombok.extern.slf4j.Slf4j
;
import
java.util.Date
;
/**
* @program: byit-myth-job->RunJobThread
* @description: 任务线程
* @author: huangfu
* @date: 2019/12/16 11:40
**/
@Slf4j
public
class
RunJobThread
implements
Runnable
{
private
AdminSenPluginDto
adminSenPluginDto
;
public
RunJobThread
(
AdminSenPluginDto
adminSenPluginDto
)
{
this
.
adminSenPluginDto
=
adminSenPluginDto
;
}
@Override
public
void
run
()
{
//创建回复对象
PluginJobRunResultDto
pluginJobRunResultDto
=
new
PluginJobRunResultDto
();
pluginJobRunResultDto
.
setStartTime
(
new
Date
());
//设定运行标识
pluginJobRunResultDto
.
setJobRunId
(
adminSenPluginDto
.
getRunId
());
//运行任务
ReturnResult
<
String
>
stringReturnResult
=
runJob
(
adminSenPluginDto
.
getJobHandelName
(),
adminSenPluginDto
.
getJobParam
());
//设置运行结果
pluginJobRunResultDto
.
setReturnResult
(
stringReturnResult
);
//获取回调通知URL
String
callbackUrl
=
adminSenPluginDto
.
getCallbackUrl
();
//设置结束时间
pluginJobRunResultDto
.
setEndTime
(
new
Date
());
cn
.
hutool
.
http
.
HttpUtil
.
post
(
callbackUrl
,
JSON
.
toJSONString
(
pluginJobRunResultDto
));
log
.
info
(
"---------服务器端:{}-------------"
,
pluginJobRunResultDto
);
}
private
ReturnResult
<
String
>
runJob
(
String
jobHandlerName
,
String
param
){
Class
<?
extends
IJobHandler
>
jobClass
=
JobUtils
.
jobCache
.
get
(
jobHandlerName
);
try
{
IJobHandler
iJobHandler
=
jobClass
.
newInstance
(
);
return
iJobHandler
.
execute
(
param
);
}
catch
(
Exception
e
)
{
e
.
printStackTrace
(
);
}
return
null
;
}
}
byit-myth-executor/myth-exector-plugin/src/main/java/com/byit/utils/JobUtils.java
View file @
7d4d61ff
...
...
@@ -5,7 +5,7 @@ import com.alibaba.fastjson.JSON;
import
com.byit.job.enums.plugin.PluginEnum
;
import
com.byit.job.exceptions.plugin.PluginException
;
import
com.byit.job.handler.interfaces.IJobHandler
;
import
com.byit.job.
model
.PluginBeanJobInfo
;
import
com.byit.job.
dto
.PluginBeanJobInfo
;
import
lombok.extern.slf4j.Slf4j
;
import
org.apache.commons.lang3.StringUtils
;
...
...
byit-myth-executor/myth-executor-api/src/main/java/com/byit/plugin/api/JobOperating.java
View file @
7d4d61ff
package
com
.
byit
.
plugin
.
api
;
import
com.byit.job.
model
.PluginBeanJobInfo
;
import
com.byit.job.
dto
.PluginBeanJobInfo
;
/**
* @program: byit-myth-job->JobOperating
...
...
demo-client/byit-demo-client/src/main/java/com/byit/job/Mains.java
View file @
7d4d61ff
package
com
.
byit
.
job
;
import
com.byit.job.
model
.PluginBeanJobInfo
;
import
com.byit.job.
dto
.PluginBeanJobInfo
;
import
com.byit.launcher.JobRunServerLauncher
;
import
com.byit.rpc.remoting.invoker.route.LoadBalance
;
import
com.byit.utils.JobUtils
;
...
...
doc/sql/myth_job/flow_version.sql
View file @
7d4d61ff
...
...
@@ -11,7 +11,7 @@
Target Server Version : 50721
File Encoding : 65001
Date: 1
3/12/2019 10:20:21
Date: 1
6/12/2019 10:58:04
*/
SET
NAMES
utf8mb4
;
...
...
doc/sql/myth_job/job_flight_schedule.sql
View file @
7d4d61ff
...
...
@@ -11,7 +11,7 @@
Target Server Version : 50721
File Encoding : 65001
Date: 1
3/12/2019 10:20:27
Date: 1
6/12/2019 10:58:10
*/
SET
NAMES
utf8mb4
;
...
...
@@ -28,10 +28,10 @@ CREATE TABLE `job_flight_schedule` (
`job_cron`
varchar
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'任务周期调度cron表达式'
,
`job_desc`
varchar
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'任务详情'
,
`plugin_urls`
varchar
(
512
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'插件方的url集合'
,
`routing_strategy`
char
(
1
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'路由策略:1随机(默认),2轮询,3最近最少使用,4最近最久未使用算法,5哈希算法'
,
`blocking_strategy`
char
(
1
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'阻塞策略:1丢弃,2阻塞等待(默认)'
,
`callback_token`
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'回调时的身份认证'
,
`gateway_token`
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'请求调度中心的token'
,
`routing_strategy`
varchar
(
36
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'路由策略:1随机(默认),2轮询,3最近最少使用,4最近最久未使用算法,5哈希算法'
,
`blocking_strategy`
varchar
(
36
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'阻塞策略:1丢弃,2阻塞等待(默认)'
,
`callback_token`
var
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'回调时的身份认证'
,
`gateway_token`
var
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'请求调度中心的token'
,
`request_host`
varchar
(
64
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'调度中心主机名'
,
`request_port`
varchar
(
16
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'调度中心端口号'
,
`job_param`
varchar
(
512
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'任务参数'
,
...
...
@@ -49,7 +49,7 @@ CREATE TABLE `job_flight_schedule` (
`author`
varchar
(
128
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'任务负责人'
,
`add_time`
datetime
(
0
)
NULL
DEFAULT
NULL
COMMENT
'任务节点添加时间'
,
`update_time`
datetime
(
0
)
NULL
DEFAULT
NULL
COMMENT
'任务节点修改时间'
,
`run_id`
varchar
(
36
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
,
`run_id`
varchar
(
36
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'运行标识'
,
PRIMARY
KEY
(
`id`
)
USING
BTREE
)
ENGINE
=
InnoDB
CHARACTER
SET
=
utf8
COLLATE
=
utf8_general_ci
ROW_FORMAT
=
Dynamic
;
...
...
doc/sql/myth_job/job_flow.sql
View file @
7d4d61ff
...
...
@@ -11,7 +11,7 @@
Target Server Version : 50721
File Encoding : 65001
Date: 1
3/12/2019 10:20:34
Date: 1
6/12/2019 10:58:16
*/
SET
NAMES
utf8mb4
;
...
...
doc/sql/myth_job/job_flow_nodes.sql
View file @
7d4d61ff
...
...
@@ -11,7 +11,7 @@
Target Server Version : 50721
File Encoding : 65001
Date: 1
3/12/2019 10:20:42
Date: 1
6/12/2019 10:58:24
*/
SET
NAMES
utf8mb4
;
...
...
@@ -28,10 +28,10 @@ CREATE TABLE `job_flow_nodes` (
`job_cron`
varchar
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'任务周期调度cron表达式'
,
`job_desc`
varchar
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'任务详情'
,
`plugin_urls`
varchar
(
512
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'插件方的url集合'
,
`routing_strategy`
char
(
1
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'路由策略:1随机(默认),2轮询,3最近最少使用,4最近最久未使用算法,5哈希算法'
,
`blocking_strategy`
char
(
1
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'阻塞策略:1丢弃,2阻塞等待(默认)'
,
`callback_token`
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'回调时的身份认证'
,
`gateway_token`
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'请求调度中心的token'
,
`routing_strategy`
varchar
(
36
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'路由策略:1随机(默认),2轮询,3最近最少使用,4最近最久未使用算法,5哈希算法'
,
`blocking_strategy`
varchar
(
36
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'阻塞策略:1丢弃,2阻塞等待(默认)'
,
`callback_token`
var
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'回调时的身份认证'
,
`gateway_token`
var
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'请求调度中心的token'
,
`request_host`
varchar
(
64
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'调度中心主机名'
,
`request_port`
varchar
(
16
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'调度中心端口号'
,
`job_param`
varchar
(
512
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'任务参数'
,
...
...
doc/sql/myth_job/job_flow_snapshot.sql
View file @
7d4d61ff
...
...
@@ -11,7 +11,7 @@
Target Server Version : 50721
File Encoding : 65001
Date: 1
3/12/2019 10:20:47
Date: 1
6/12/2019 10:58:29
*/
SET
NAMES
utf8mb4
;
...
...
doc/sql/myth_job/job_lock.sql
View file @
7d4d61ff
...
...
@@ -11,7 +11,7 @@
Target Server Version : 50721
File Encoding : 65001
Date: 1
3/12/2019 10:20:5
5
Date: 1
6/12/2019 10:58:3
5
*/
SET
NAMES
utf8mb4
;
...
...
@@ -27,4 +27,10 @@ CREATE TABLE `job_lock` (
PRIMARY
KEY
(
`lock_name`
)
USING
BTREE
)
ENGINE
=
InnoDB
CHARACTER
SET
=
utf8
COLLATE
=
utf8_general_ci
ROW_FORMAT
=
Dynamic
;
-- ----------------------------
-- Records of job_lock
-- ----------------------------
INSERT
INTO
`job_lock`
VALUES
(
'scanning_job_flight_schedule'
,
'扫描job_flight_schedule表,添加到任务轮里面的表'
);
INSERT
INTO
`job_lock`
VALUES
(
'scanning_job_read_ahead'
,
'扫描job_read_ahead表,添加到排期表的行锁'
);
SET
FOREIGN_KEY_CHECKS
=
1
;
doc/sql/myth_job/job_log.sql
View file @
7d4d61ff
...
...
@@ -11,7 +11,7 @@
Target Server Version : 50721
File Encoding : 65001
Date: 1
3/12/2019 10:21:0
0
Date: 1
6/12/2019 10:58:4
0
*/
SET
NAMES
utf8mb4
;
...
...
doc/sql/myth_job/job_logglue.sql
View file @
7d4d61ff
...
...
@@ -11,7 +11,7 @@
Target Server Version : 50721
File Encoding : 65001
Date: 1
3/12/2019 10:21:04
Date: 1
6/12/2019 10:58:45
*/
SET
NAMES
utf8mb4
;
...
...
doc/sql/myth_job/job_read_ahead.sql
View file @
7d4d61ff
...
...
@@ -11,7 +11,7 @@
Target Server Version : 50721
File Encoding : 65001
Date: 1
3/12/2019 10:21:09
Date: 1
6/12/2019 10:58:50
*/
SET
NAMES
utf8mb4
;
...
...
@@ -28,10 +28,10 @@ CREATE TABLE `job_read_ahead` (
`job_cron`
varchar
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'任务周期调度cron表达式'
,
`job_desc`
varchar
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'任务详情'
,
`plugin_urls`
varchar
(
512
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'插件方的url集合'
,
`routing_strategy`
char
(
1
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'路由策略:1随机(默认),2轮询,3最近最少使用,4最近最久未使用算法,5哈希算法'
,
`blocking_strategy`
char
(
1
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'阻塞策略:1丢弃,2阻塞等待(默认)'
,
`callback_token`
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'回调时的身份认证'
,
`gateway_token`
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'请求调度中心的token'
,
`routing_strategy`
varchar
(
36
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'路由策略:1随机(默认),2轮询,3最近最少使用,4最近最久未使用算法,5哈希算法'
,
`blocking_strategy`
varchar
(
36
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NOT
NULL
COMMENT
'阻塞策略:1丢弃,2阻塞等待(默认)'
,
`callback_token`
var
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'回调时的身份认证'
,
`gateway_token`
var
char
(
255
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'请求调度中心的token'
,
`request_host`
varchar
(
64
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'调度中心主机名'
,
`request_port`
varchar
(
16
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'调度中心端口号'
,
`job_param`
varchar
(
512
)
CHARACTER
SET
utf8
COLLATE
utf8_general_ci
NULL
DEFAULT
NULL
COMMENT
'任务参数'
,
...
...
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