Commit c55f178b by huangfusuper

添加任务调度执行器,优化回调接口!

parent 844a728f
package com.byit.controller; package com.byit.controller;
import com.byit.job.WorkRoulette; import com.byit.job.WorkRoulette;
import com.byit.job.model.MythJobFlightSchedule;
import com.byit.job.model.MythJobReadAhead; import com.byit.job.model.MythJobReadAhead;
import com.byit.job.model.PluginBeanJobInfo; import com.byit.job.model.PluginBeanJobInfo;
import com.byit.job.utils.SourceObj2TargetObjUtil;
import com.byit.job.vo.ReturnResult; import com.byit.job.vo.ReturnResult;
import com.byit.service.JobReadAheadService; import com.byit.service.JobReadAheadService;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
...@@ -21,18 +23,18 @@ import java.util.List; ...@@ -21,18 +23,18 @@ import java.util.List;
@RequestMapping("job") @RequestMapping("job")
public class JobController { public class JobController {
@Autowired @Autowired
private DataSource dataSource;
@Autowired
private JobReadAheadService jobReadAheadService; private JobReadAheadService jobReadAheadService;
@PostMapping(value = "addJob") @PostMapping(value = "addJob")
public String addJob(@RequestBody PluginBeanJobInfo pluginBeanJobInfo){ public String addJob(@RequestBody PluginBeanJobInfo pluginBeanJobInfo){
WorkRoulette.addJob(pluginBeanJobInfo); MythJobReadAhead mythJobReadAhead = SourceObj2TargetObjUtil.pluginBeanJobInfo2MythJobReadAhead(pluginBeanJobInfo);
jobReadAheadService.addOneMythJobReadAhead(mythJobReadAhead);
return "SUCCESS"; return "SUCCESS";
} }
@PostMapping(value = "callbackRes") @PostMapping(value = "callbackRes")
public void callbackRes(@RequestBody ReturnResult<String> result){ public String callbackRes(@RequestBody ReturnResult<String> result){
System.out.println(result.getCode()+"-----"+result.getMsg()); System.out.println(result.getCode()+"-----"+result.getMsg());
return "好的,我知道你执行成功了";
} }
@GetMapping(value = "getJobInfo") @GetMapping(value = "getJobInfo")
......
...@@ -15,8 +15,20 @@ import java.util.List; ...@@ -15,8 +15,20 @@ import java.util.List;
@Repository @Repository
public interface JobFlightScheduleMapper { public interface JobFlightScheduleMapper {
/** /**
* 查询预读五秒的数据
* @param maxNextTime
* @return
*/
List<MythJobFlightSchedule> findMythJobFlightScheduleByTriggerNextTime(@Param("maxNextTime")long maxNextTime);
/**
* 批量向排期表添加任务数据 * 批量向排期表添加任务数据
* @param mythJobFlightSchedules * @param mythJobFlightSchedules
*/ */
void addJobFlightSchedules(@Param("mythJobFlightSchedule") List<MythJobFlightSchedule> mythJobFlightSchedules); void addJobFlightSchedules(@Param("mythJobFlightSchedule") List<MythJobFlightSchedule> mythJobFlightSchedules);
/**
* 删除已经加载的节点
* @param id
*/
void deleteMythJobFlightScheduleById(@Param("id")String id);
} }
...@@ -22,8 +22,15 @@ public interface JobInfoMapper { ...@@ -22,8 +22,15 @@ public interface JobInfoMapper {
List<MythJobReadAhead> findMythJobReadAheadByTriggerNextTime(@Param("maxNextTime")long maxNextTime); List<MythJobReadAhead> findMythJobReadAheadByTriggerNextTime(@Param("maxNextTime")long maxNextTime);
/** /**
* 添加单个任务节点
* @param mythJobReadAhead
*/
void addOneMythJobReadAhead(MythJobReadAhead mythJobReadAhead);
/**
* 删除已经被调度的任务根据ID * 删除已经被调度的任务根据ID
* @param id * @param id
*/ */
void removeMythJobReadAheadById(@Param("id") String id); void removeMythJobReadAheadById(@Param("id") String id);
} }
package com.byit.job; package com.byit.job;
import com.byit.job.model.MythJobFlightSchedule;
import com.byit.job.model.PluginBeanJobInfo; import com.byit.job.model.PluginBeanJobInfo;
import com.byit.task.JavaBeanJobTask; import com.byit.task.JavaBeanJobTask;
import io.netty.util.HashedWheelTimer; import io.netty.util.HashedWheelTimer;
import io.netty.util.TimerTask;
import java.util.concurrent.ThreadFactory; import java.util.concurrent.ThreadFactory;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
...@@ -30,9 +32,8 @@ public class WorkRoulette { ...@@ -30,9 +32,8 @@ public class WorkRoulette {
}, 1, TimeUnit.SECONDS, 8, true, 0); }, 1, TimeUnit.SECONDS, 8, true, 0);
public static void addJob(PluginBeanJobInfo pluginBeanJobInfo) { public static void addJob(TimerTask timerTask) {
JavaBeanJobTask javaBeanJobTask = new JavaBeanJobTask(pluginBeanJobInfo); hashedWheelTimer.newTimeout(timerTask, TimeUnit.SECONDS.toNanos(20), TimeUnit.NANOSECONDS);
hashedWheelTimer.newTimeout(javaBeanJobTask, TimeUnit.SECONDS.toNanos(20), TimeUnit.NANOSECONDS);
} }
} }
...@@ -13,9 +13,23 @@ import java.util.List; ...@@ -13,9 +13,23 @@ import java.util.List;
**/ **/
public interface JobFlightScheduleService { public interface JobFlightScheduleService {
/** /**
* 查询预读五秒的数据
* @param maxNextTime
* @return
*/
List<MythJobFlightSchedule> findMythJobFlightScheduleByTriggerNextTime(long maxNextTime);
/**
* 批量向排期表添加任务数据 * 批量向排期表添加任务数据
* @param mythJobFlightSchedules * @param mythJobFlightSchedules
*/ */
void addJobFlightSchedules(List<MythJobFlightSchedule> mythJobFlightSchedules); void addJobFlightSchedules(List<MythJobFlightSchedule> mythJobFlightSchedules);
/**
* 删除已经加载的节点
* @param id
*/
void deleteMythJobFlightScheduleById(String id);
} }
...@@ -19,6 +19,12 @@ public interface JobReadAheadService { ...@@ -19,6 +19,12 @@ public interface JobReadAheadService {
List<MythJobReadAhead> findMythJobReadAheadByTriggerNextTime(long maxNextTime); List<MythJobReadAhead> findMythJobReadAheadByTriggerNextTime(long maxNextTime);
/** /**
* 添加单个任务节点
* @param mythJobReadAhead
*/
void addOneMythJobReadAhead(MythJobReadAhead mythJobReadAhead);
/**
* 删除已经被调度的任务根据ID * 删除已经被调度的任务根据ID
* @param id * @param id
*/ */
......
...@@ -18,8 +18,21 @@ import java.util.List; ...@@ -18,8 +18,21 @@ import java.util.List;
public class JobFlightScheduleServiceImpl implements JobFlightScheduleService { public class JobFlightScheduleServiceImpl implements JobFlightScheduleService {
@Autowired @Autowired
private JobFlightScheduleMapper jobFlightScheduleMapper; private JobFlightScheduleMapper jobFlightScheduleMapper;
@Override
public List<MythJobFlightSchedule> findMythJobFlightScheduleByTriggerNextTime(long maxNextTime) {
return jobFlightScheduleMapper.findMythJobFlightScheduleByTriggerNextTime(maxNextTime);
}
@Override @Override
public void addJobFlightSchedules(List<MythJobFlightSchedule> mythJobFlightSchedules) { public void addJobFlightSchedules(List<MythJobFlightSchedule> mythJobFlightSchedules) {
jobFlightScheduleMapper.addJobFlightSchedules(mythJobFlightSchedules); jobFlightScheduleMapper.addJobFlightSchedules(mythJobFlightSchedules);
} }
@Override
public void deleteMythJobFlightScheduleById(String id) {
jobFlightScheduleMapper.deleteMythJobFlightScheduleById(id);
}
} }
...@@ -27,6 +27,11 @@ public class JobReadAheadServiceImpl implements JobReadAheadService { ...@@ -27,6 +27,11 @@ public class JobReadAheadServiceImpl implements JobReadAheadService {
} }
@Override @Override
public void addOneMythJobReadAhead(MythJobReadAhead mythJobReadAhead) {
jobInfoMapper.addOneMythJobReadAhead(mythJobReadAhead);
}
@Override
public void removeMythJobReadAheadById(String id) { public void removeMythJobReadAheadById(String id) {
jobInfoMapper.removeMythJobReadAheadById(id); jobInfoMapper.removeMythJobReadAheadById(id);
} }
......
...@@ -4,7 +4,7 @@ import cn.hutool.http.HttpUtil; ...@@ -4,7 +4,7 @@ import cn.hutool.http.HttpUtil;
import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.JSONObject;
import com.byit.job.exceptions.plugin.PluginException; import com.byit.job.exceptions.plugin.PluginException;
import com.byit.job.model.PluginBeanJobInfo; import com.byit.job.model.MythJobFlightSchedule;
import com.byit.job.utils.IpUtil; import com.byit.job.utils.IpUtil;
import com.byit.rpc.remoting.invoker.route.LoadBalance; import com.byit.rpc.remoting.invoker.route.LoadBalance;
import com.byit.rpc.remoting.invoker.route.RpcLoadBalance; import com.byit.rpc.remoting.invoker.route.RpcLoadBalance;
...@@ -20,24 +20,23 @@ import lombok.extern.slf4j.Slf4j; ...@@ -20,24 +20,23 @@ import lombok.extern.slf4j.Slf4j;
**/ **/
@Slf4j @Slf4j
public class JavaBeanJobTask implements TimerTask { public class JavaBeanJobTask implements TimerTask {
private PluginBeanJobInfo pluginBeanJobInfo; private MythJobFlightSchedule mythJobFlightSchedule;
public JavaBeanJobTask(PluginBeanJobInfo pluginBeanJobInfo) { public JavaBeanJobTask( MythJobFlightSchedule mythJobFlightSchedule) {
this.pluginBeanJobInfo = pluginBeanJobInfo; this.mythJobFlightSchedule = mythJobFlightSchedule;
} }
@Override @Override
public void run(Timeout timeout) { public void run(Timeout timeout) {
//根据负责均衡方案获取对应IP //根据负责均衡方案获取对应IP
RpcLoadBalance rpcInvokerRouter = LoadBalance.match(pluginBeanJobInfo.getRoutingStrategy( ), LoadBalance.ROUND).rpcInvokerRouter; RpcLoadBalance rpcInvokerRouter = LoadBalance.match(mythJobFlightSchedule.getRoutingStrategy( ), LoadBalance.ROUND).rpcInvokerRouter;
try{ try{
String url = IpUtil.electiveUrl(pluginBeanJobInfo.getUrl( ), rpcInvokerRouter); String url = IpUtil.electiveUrl(mythJobFlightSchedule.getPluginUrls( ), rpcInvokerRouter);
String jobHandelName = pluginBeanJobInfo.getJobHandelName( ); String jobHandelName = mythJobFlightSchedule.getExecutorHandler( );
String param = pluginBeanJobInfo.getParam( ); String param = mythJobFlightSchedule.getJobParam( );
JSONObject jsonObject = new JSONObject(); JSONObject jsonObject = new JSONObject();
jsonObject.put("jobHandelName",jobHandelName); jsonObject.put("jobHandelName",jobHandelName);
jsonObject.put("callbackMethod","http://127.0.0.1:8080/callbackRes"); jsonObject.put("callbackMethod","http://127.0.0.1:8080/job/callbackRes");
jsonObject.put("param",param); jsonObject.put("param",param);
String result = HttpUtil.post(url, JSON.toJSONString(jsonObject),10*1000); String result = HttpUtil.post(url, JSON.toJSONString(jsonObject),10*1000);
log.info("---------------{}------------",result); log.info("---------------{}------------",result);
......
package com.byit.thread; package com.byit.thread;
import cn.hutool.core.collection.CollectionUtil; 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.MythJobFlightSchedule;
import com.byit.job.model.MythJobReadAhead; import com.byit.job.model.MythJobReadAhead;
import com.byit.service.JobFlightScheduleService; import com.byit.service.JobFlightScheduleService;
import com.byit.service.JobReadAheadService; 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 lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.BeanUtils; import org.springframework.beans.BeanUtils;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import javax.sql.DataSource; import javax.sql.DataSource;
...@@ -16,6 +26,7 @@ import java.sql.PreparedStatement; ...@@ -16,6 +26,7 @@ import java.sql.PreparedStatement;
import java.sql.SQLException; import java.sql.SQLException;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
/** /**
...@@ -35,7 +46,6 @@ public class JobScheduleHelper{ ...@@ -35,7 +46,6 @@ public class JobScheduleHelper{
@Autowired @Autowired
private JobFlightScheduleService jobFlightScheduleService; private JobFlightScheduleService jobFlightScheduleService;
/** /**
* 读取任务节点的预读 * 读取任务节点的预读
*/ */
...@@ -68,14 +78,14 @@ public class JobScheduleHelper{ ...@@ -68,14 +78,14 @@ public class JobScheduleHelper{
public void start(){ public void start(){
jobInfoThread = new Thread(()->{ jobInfoThread = new Thread(()->{
try { try {
TimeUnit.MILLISECONDS.sleep(5000 - System.currentTimeMillis()%1000 ); TimeUnit.MILLISECONDS.sleep(PRE_READ_MS - System.currentTimeMillis()%1000 );
} catch (InterruptedException e) { } catch (InterruptedException e) {
if (!jobInfoThreadToStop) { if (!jobInfoThreadToStop) {
log.error(e.getMessage(), e); log.error(e.getMessage(), e);
} }
} }
log.info("---------------------init myth-job admin jobInfoThread success------------------------"); log.info("---------------------init myth-job admin jobInfoThread success------------------------");
boolean preReadSuc = false; boolean preReadSuc;
while (!jobInfoThreadToStop) { while (!jobInfoThreadToStop) {
//开始去扫描任务节点 //开始去扫描任务节点
long start = System.currentTimeMillis(); long start = System.currentTimeMillis();
...@@ -92,7 +102,7 @@ public class JobScheduleHelper{ ...@@ -92,7 +102,7 @@ public class JobScheduleHelper{
//修改为不自动提交 //修改为不自动提交
conn.setAutoCommit(false); conn.setAutoCommit(false);
//先添加扫描任务节点的行锁 //先添加扫描任务节点的行锁
preparedStatement = conn.prepareStatement("SELECT * FROM JOB_LOCK WHERE LOCK_NAME = 'scanning_job_nodes' FOR UPDATE "); preparedStatement = conn.prepareStatement("SELECT * FROM JOB_LOCK WHERE LOCK_NAME = 'scanning_job_read_ahead' FOR UPDATE ");
preparedStatement.execute(); preparedStatement.execute();
//行锁已经加上 后续处理 //行锁已经加上 后续处理
long nowTime = System.currentTimeMillis(); long nowTime = System.currentTimeMillis();
...@@ -102,21 +112,25 @@ public class JobScheduleHelper{ ...@@ -102,21 +112,25 @@ public class JobScheduleHelper{
if(CollectionUtil.isNotEmpty(mythJobReadAheadByTriggerNextTime)){ if(CollectionUtil.isNotEmpty(mythJobReadAheadByTriggerNextTime)){
List<MythJobFlightSchedule> mythJobFlightSchedules = new ArrayList<>(15); List<MythJobFlightSchedule> mythJobFlightSchedules = new ArrayList<>(15);
mythJobReadAheadByTriggerNextTime.forEach(mythJobInfo -> { mythJobReadAheadByTriggerNextTime.forEach(mythJobInfo -> {
log.info("任务:{}", mythJobInfo);
/** /**
* 需要去检验当前任务的上级节点是否已经执行成功,没有执行,或者处于暂停状态则跳过该任务 * 需要去检验当前任务的上级节点是否已经执行成功,没有执行,或者处于暂停状态则跳过该任务
* 大概思路,根据任务流id,从任务流执行回溯表查询该任务流的所有节点,查看上级节点是否已经执行成功 * 大概思路,根据任务流id,从任务流执行回溯表查询该任务流的所有节点,查看上级节点是否已经执行成功
*/ */
{ if(StringUtils.isNotBlank(mythJobInfo.getParentId())){
log.info("任务:{}", mythJobInfo);
MythJobFlightSchedule mythJobFlightSchedule = new MythJobFlightSchedule(); MythJobFlightSchedule mythJobFlightSchedule = new MythJobFlightSchedule();
BeanUtils.copyProperties(mythJobInfo,mythJobFlightSchedule); BeanUtils.copyProperties(mythJobInfo,mythJobFlightSchedule);
mythJobFlightSchedules.add(mythJobFlightSchedule); mythJobFlightSchedules.add(mythJobFlightSchedule);
jobReadAheadService.removeMythJobReadAheadById(mythJobInfo.getId()); jobReadAheadService.removeMythJobReadAheadById(mythJobInfo.getId());
} }
}); });
if(CollectionUtil.isNotEmpty(mythJobFlightSchedules)) {
jobFlightScheduleService.addJobFlightSchedules(mythJobFlightSchedules); jobFlightScheduleService.addJobFlightSchedules(mythJobFlightSchedules);
}
}else{ }else{
log.info("-------------------空轮转-------------------"); //log.info("-------------------空轮转-------------------");
preReadSuc = false; preReadSuc = false;
} }
}catch (Exception e){ }catch (Exception e){
...@@ -182,8 +196,127 @@ public class JobScheduleHelper{ ...@@ -182,8 +196,127 @@ public class JobScheduleHelper{
//设置为守护线程 //设置为守护线程
jobInfoThread.setDaemon(true); jobInfoThread.setDaemon(true);
//设置名字 //设置名字
jobInfoThread.setName("myth-job,admin JobScheduleHelper#scheduleThread"); jobInfoThread.setName("myth-job,admin JobScheduleHelper#jobInfoThread");
jobInfoThread.start(); jobInfoThread.start();
/**
* 开始操作任务排期表:
* 1. 扫描五秒要执行的数据
* 2.加载时,创建任务日志,传入调度时间
* 3.加载到任务调度轮盘
* 4.删除数据
*/
scheduleThread = new Thread(() ->{
try {
TimeUnit.MILLISECONDS.sleep(4000 - System.currentTimeMillis()%1000 );
} catch (InterruptedException e) {
if (!scheduleThreadToStop) {
log.error(e.getMessage(), e);
}
}
log.info("---------------------init myth-job admin scheduleThread success------------------------");
/**
* 判断是否是空轮转
* 空轮转的话是需要休眠的
*/
boolean preReadSuc;
while (!scheduleThreadToStop){
//开始去扫描任务节点
long start = System.currentTimeMillis();
Connection conn = null;
Boolean connAutoCommit = null;
PreparedStatement preparedStatement = null;
preReadSuc = true;
try{
conn = dataSource.getConnection();
//记录当前的自动提交状态
connAutoCommit = conn.getAutoCommit();
//修改为不自动提交
conn.setAutoCommit(false);
//先添加扫描任务节点的行锁
preparedStatement = conn.prepareStatement("SELECT * FROM JOB_LOCK WHERE LOCK_NAME = 'scanning_job_flight_schedule' FOR UPDATE ");
preparedStatement.execute();
//行锁已经加上 后续处理
long nowTime = System.currentTimeMillis();
//查询所有符合条件的任务节点
List<MythJobFlightSchedule> mythJobFlightScheduleByTriggerNextTimes = jobFlightScheduleService.findMythJobFlightScheduleByTriggerNextTime(nowTime + SCHEDULE_READ_MS);
if(CollectionUtil.isNotEmpty(mythJobFlightScheduleByTriggerNextTimes)){
//循环遍历添加任务
mythJobFlightScheduleByTriggerNextTimes.forEach(mythJobFlightScheduleByTriggerNextTime ->{
if ("BEAN".equals(mythJobFlightScheduleByTriggerNextTime.getJobType())) {
JavaBeanJobTask javaBeanJobTask = new JavaBeanJobTask(mythJobFlightScheduleByTriggerNextTime);
jobFlightScheduleService.deleteMythJobFlightScheduleById(mythJobFlightScheduleByTriggerNextTime.getId());
WorkRoulette.addJob(javaBeanJobTask);
}
});
}else{
//log.info("-----------------任务排期空轮转-----------------");
preReadSuc = false;
}
}catch (Exception e){
if(!scheduleThreadToStop){
log.error("------------------扫描job_flight_schedule出现异常:{}-----------------",e.getMessage());
}
}finally {
//commit
if (conn!=null) {
try {
conn.commit();
}catch (SQLException e){
if(!scheduleThreadToStop){
log.error("----------------提交行锁错误:{}-------------------",e.getMessage());
}
}
//设置提交状态恢复原来的值
try {
conn.setAutoCommit(connAutoCommit);
}catch (SQLException e){
if(!scheduleThreadToStop){
log.error("------------------设置为自动提交出错:{}------------------",e.getMessage());
}
}
//关闭连接
try {
conn.close();
}catch (SQLException e){
if(!scheduleThreadToStop){
log.error("----------------关闭数据库连接:{}-------------------",e.getMessage());
}
}
}
//关闭执行器
if(null !=preparedStatement){
try {
preparedStatement.close();
} catch (SQLException e) {
if(!scheduleThreadToStop){
log.error("---------------关闭执行器出错:{}----------------------",e.getMessage());
}
}
}
//计算总的处理时间 将时间保持在一秒处理一次
long cost = System.currentTimeMillis()-start;
//不足一秒的等待 等够1秒为止
if(cost < 1000){
// 预读期:成功-每秒扫描一次;失败-跳过这段时间
try {
//这个睡眠是空轮转时,延长睡眠时间,奖励CUP的使用频率
TimeUnit.MILLISECONDS.sleep((preReadSuc?1000:SCHEDULE_READ_MS) - System.currentTimeMillis()%1000);
} catch (InterruptedException e) {
if (!scheduleThreadToStop) {
log.error("----------------------{}-------------------------",e.getMessage());
}
}
}
}
}
});
scheduleThread.setName("myth-job,admin JobScheduleHelper#scheduleThread");
scheduleThread.setDaemon(true);
scheduleThread.start();
} }
public void setDataSource(DataSource dataSource) { public void setDataSource(DataSource dataSource) {
...@@ -191,6 +324,6 @@ public class JobScheduleHelper{ ...@@ -191,6 +324,6 @@ public class JobScheduleHelper{
} }
public static void main(String[] args) { public static void main(String[] args) {
System.out.println(System.currentTimeMillis()+30*1000 ); System.out.println(System.currentTimeMillis()+30*1000);
} }
} }
<?xml version="1.0" encoding="UTF-8"?> <?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" > <!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd" >
<mapper namespace="com.byit.dao.JobFlightScheduleMapper"> <mapper namespace="com.byit.dao.JobFlightScheduleMapper">
<sql id="Base_Column_List">
t.id,
t.job_name,
t.executor_handler,
t.job_cron,
t.job_desc,
t.plugin_urls,
t.routing_strategy,
t.blocking_strategy,
t.callback_token,
t.gateway_token,
t.request_host,
t.request_port,
t.job_param,
t.job_type,
t.taskflow_id,
t.parent_id,
t.taskflow_is_cron,
t.alarm_email,
t.executor_timeout,
t.source_urls,
t.glue_source,
t.glue_remark,
t.glue_updatetime,
t.trigger_next_time,
t.author,
t.add_time,
t.update_time,
t.run_id
</sql>
<resultMap id="mythJobFlightSchedule" type="com.byit.job.model.MythJobFlightSchedule">
<id column="id" property="id"/>
<result column="job_name" property="jobName"/>
<result column="executor_handler" property="executorHandler"/>
<result column="job_cron" property="jobCron"/>
<result column="job_desc" property="jobDesc"/>
<result column="plugin_urls" property="pluginUrls"/>
<result column="routing_strategy" property="routingStrategy"/>
<result column="blocking_strategy" property="blockingStrategy"/>
<result column="callback_token" property="callbackToken"/>
<result column="gateway_token" property="gatewayToken"/>
<result column="request_host" property="requestHost"/>
<result column="request_port" property="requestPort"/>
<result column="job_param" property="jobParam"/>
<result column="job_type" property="jobType"/>
<result column="taskflow_id" property="taskFlowId"/>
<result column="parent_id" property="parentId"/>
<result column="taskflow_is_cron" property="taskFlowIsCron"/>
<result column="alarm_email" property="alarmEmail"/>
<result column="executor_timeout" property="executorTimeout"/>
<result column="source_urls" property="sourceUrls"/>
<result column="glue_source" property="glueSource"/>
<result column="glue_remark" property="glueRemark"/>
<result column="glue_updatetime" property="glueUpdateTime"/>
<result column="trigger_next_time" property="triggerNextTime"/>
<result column="author" property="author"/>
<result column="add_time" property="addTime"/>
<result column="update_time" property="updateTime"/>
<result column="run_id" property="runID"/>
</resultMap>
<select id="findMythJobFlightScheduleByTriggerNextTime" resultMap="mythJobFlightSchedule">
SELECT <include refid="Base_Column_List" />
FROM job_flight_schedule AS t
WHERE t.trigger_next_time<![CDATA[ <= ]]> #{maxNextTime}
</select>
<insert id="addJobFlightSchedules" parameterType="com.byit.job.model.MythJobFlightSchedule"> <insert id="addJobFlightSchedules" parameterType="com.byit.job.model.MythJobFlightSchedule">
INSERT INTO job_flight_schedule INSERT INTO job_flight_schedule
( (
...@@ -69,4 +139,7 @@ ...@@ -69,4 +139,7 @@
</insert> </insert>
<delete id="deleteMythJobFlightScheduleById">
DELETE FROM job_flight_schedule WHERE id= #{id}
</delete>
</mapper> </mapper>
\ No newline at end of file
...@@ -37,9 +37,10 @@ ...@@ -37,9 +37,10 @@
<if test="addTime != null">add_time =#{addTime},</if> <if test="addTime != null">add_time =#{addTime},</if>
<if test="updateTime != null">update_time =#{updateTime},</if> <if test="updateTime != null">update_time =#{updateTime},</if>
<if test="removalMark != null">removal_mark =#{removalMark},</if> <if test="removalMark != null and removalMark!=''">removal_mark =#{removalMark},</if>
<if test="repeatTimes != null">repeat_times =#{repeatTimes},</if> <if test="repeatTimes != null">repeat_times =#{repeatTimes},</if>
<if test="flowVersion != null and flowVersion!=''">flow_version =#{flowVersion},</if>
</trim> </trim>
WHERE id=#{id} WHERE id=#{id}
</update> </update>
......
...@@ -72,6 +72,19 @@ ...@@ -72,6 +72,19 @@
AND t.trigger_next_time<![CDATA[ <= ]]> #{maxNextTime} AND t.trigger_next_time<![CDATA[ <= ]]> #{maxNextTime}
</select> </select>
<insert id="addOneMythJobReadAhead" parameterType="com.byit.job.model.MythJobReadAhead">
insert into job_read_ahead (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,trigger_status,glue_source,glue_remark,
glue_updatetime,trigger_next_time,author,add_time,update_time,run_id)
values(#{id},#{jobName},#{executorHandler},#{jobCron},#{jobDesc},#{pluginUrls},#{routingStrategy},#{blockingStrategy},
#{callbackToken},#{gatewayToken},#{requestHost},#{requestPort},#{jobParam},#{jobType},#{taskFlowId},
#{parentId},#{taskFlowIsCron},#{alarmEmail},#{executorTimeout},#{sourceUrls},#{triggerStatus},#{glueSource},#{glueRemark},
#{glueUpdateTime},#{triggerNextTime},#{author},#{addTime},#{updateTime},#{runID})
</insert>
<delete id="removeMythJobReadAheadById"> <delete id="removeMythJobReadAheadById">
DELETE FROM job_read_ahead WHERE id=#{id} DELETE FROM job_read_ahead WHERE id=#{id}
</delete> </delete>
......
...@@ -41,6 +41,11 @@ ...@@ -41,6 +41,11 @@
<artifactId>byit-myth-rpc</artifactId> <artifactId>byit-myth-rpc</artifactId>
</dependency> </dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
</dependency>
</dependencies> </dependencies>
</project> </project>
\ No newline at end of file
...@@ -137,5 +137,8 @@ public class MythJobFlowNodes { ...@@ -137,5 +137,8 @@ public class MythJobFlowNodes {
* 重复次数 -1永久运行 * 重复次数 -1永久运行
*/ */
private Integer repeatTimes; private Integer repeatTimes;
/**
* 版本号
*/
private String flowVersion;
} }
...@@ -58,22 +58,9 @@ public class PluginBeanJobInfo { ...@@ -58,22 +58,9 @@ public class PluginBeanJobInfo {
*/ */
private String param; private String param;
private String alarmEmail;
/**
@Override * 负责人
public String toString() { */
return "PluginBeanJobInfo{" + private String author;
"jobHandelName='" + jobHandelName + '\'' +
", url='" + url + '\'' +
", mythCron='" + mythCron + '\'' +
", routingStrategy='" + routingStrategy + '\'' +
", blockingStrategy='" + blockingStrategy + '\'' +
", callbackToken='" + callbackToken + '\'' +
", gatewayToken='" + gatewayToken + '\'' +
", requestHost='" + requestHost + '\'' +
", requestPort=" + requestPort +
", param='" + param + '\'' +
'}';
}
} }
package com.byit.job.utils;
import com.byit.job.model.MythJobReadAhead;
import com.byit.job.model.PluginBeanJobInfo;
import java.text.ParseException;
import java.util.Date;
import java.util.UUID;
/**
* @program: byit-myth-job->SourceObj2TargetObjUtil
* @description: 这是一个转换的工具类,定义两个类之间的相互转换
* @author: huangfu
* @date: 2019/12/14 13:00
**/
public class SourceObj2TargetObjUtil {
/**
* 插件端的实体对象转换为对任务节点的对象
* @param pluginBeanJobInfo
* @return
*/
public static MythJobReadAhead pluginBeanJobInfo2MythJobReadAhead(PluginBeanJobInfo pluginBeanJobInfo){
MythJobReadAhead mythJobReadAhead = new MythJobReadAhead();
mythJobReadAhead.setId(UUID.randomUUID( ).toString().replace("-",""));
mythJobReadAhead.setExecutorHandler(pluginBeanJobInfo.getJobHandelName());
mythJobReadAhead.setPluginUrls(pluginBeanJobInfo.getUrl());
mythJobReadAhead.setJobCron(pluginBeanJobInfo.getMythCron());
mythJobReadAhead.setRoutingStrategy(pluginBeanJobInfo.getRoutingStrategy());
mythJobReadAhead.setBlockingStrategy(pluginBeanJobInfo.getBlockingStrategy());
mythJobReadAhead.setCallbackToken(pluginBeanJobInfo.getCallbackToken());
mythJobReadAhead.setGatewayToken(pluginBeanJobInfo.getGatewayToken());
mythJobReadAhead.setJobParam(pluginBeanJobInfo.getParam());
mythJobReadAhead.setAlarmEmail(pluginBeanJobInfo.getAlarmEmail());
mythJobReadAhead.setAuthor(pluginBeanJobInfo.getAuthor());
mythJobReadAhead.setAddTime(new Date());
mythJobReadAhead.setExecutorTimeout(30);
mythJobReadAhead.setJobDesc("这是一个测试任务,来自插件端的测试任务");
mythJobReadAhead.setJobName("测试任务");
mythJobReadAhead.setJobType("BEAN");
mythJobReadAhead.setTriggerStatus("1");
mythJobReadAhead.setParentId("1");
mythJobReadAhead.setTaskFlowIsCron("1");
//TODO 暂时随机分配 未来这个东西是线程自动添加的
mythJobReadAhead.setRunID(UUID.randomUUID( ).toString().replace("-",""));
mythJobReadAhead.setTriggerNextTime(System.currentTimeMillis()+60000);
return mythJobReadAhead;
}
}
...@@ -105,8 +105,8 @@ class RunJobThread implements Runnable{ ...@@ -105,8 +105,8 @@ class RunJobThread implements Runnable{
String param = ((String)(jsonObject.get("param"))); String param = ((String)(jsonObject.get("param")));
ReturnResult<String> stringReturnResult = runJob(jobHandelName, param); ReturnResult<String> stringReturnResult = runJob(jobHandelName, param);
String callbackMethod = (String)jsonObject.get("callbackMethod"); String callbackMethod = (String)jsonObject.get("callbackMethod");
cn.hutool.http.HttpUtil.post(callbackMethod,JSON.toJSONString(stringReturnResult)); String post = cn.hutool.http.HttpUtil.post(callbackMethod, JSON.toJSONString(stringReturnResult));
log.info("---------服务器端:{},{}-------------",stringReturnResult.getCode(),stringReturnResult.getMsg()); log.info("---------服务器端:{}:{}-------------",callbackMethod,post);
} }
private ReturnResult<String> runJob(String jobHandlerName,String param){ private ReturnResult<String> runJob(String jobHandlerName,String param){
......
...@@ -14,7 +14,7 @@ import com.byit.job.vo.ReturnResult; ...@@ -14,7 +14,7 @@ import com.byit.job.vo.ReturnResult;
public class DemoJob1 extends BaseJobHandler { public class DemoJob1 extends BaseJobHandler {
@Override @Override
public ReturnResult<String> execute(String s) throws Exception { public ReturnResult<String> execute(String s) throws Exception {
Thread.sleep(200000); Thread.sleep(10000);
System.out.println("--------------DemoJob1----------------"+s ); System.out.println("--------------DemoJob1----------------"+s );
return ReturnResult.SUCCESS; return ReturnResult.SUCCESS;
} }
......
...@@ -3,6 +3,7 @@ package com.byit.job; ...@@ -3,6 +3,7 @@ package com.byit.job;
import com.byit.job.model.PluginBeanJobInfo; import com.byit.job.model.PluginBeanJobInfo;
import com.byit.launcher.JobRunServerLauncher; import com.byit.launcher.JobRunServerLauncher;
import com.byit.rpc.remoting.invoker.route.LoadBalance; import com.byit.rpc.remoting.invoker.route.LoadBalance;
import com.byit.utils.JobUtils;
import java.io.IOException; import java.io.IOException;
...@@ -23,14 +24,15 @@ public class Mains { ...@@ -23,14 +24,15 @@ public class Mains {
String gatewayToken="asdsadsa"; String gatewayToken="asdsadsa";
String requestIP = "10.0.55.237"; String requestIP = "10.0.55.237";
String requestPort="8080"; String requestPort="8080";
String param="sadsadsa"; String param="不延迟任务";
String name="addJob"; String name="addJob";
PluginBeanJobInfo pluginBeanJobInfo = new PluginBeanJobInfo(name,plServerUrl,mythCron,routingStrategy,blockingStrategy,callbackToken,gatewayToken,requestIP,requestPort,param); PluginBeanJobInfo pluginBeanJobInfo = new PluginBeanJobInfo(name,plServerUrl,mythCron,routingStrategy,blockingStrategy,callbackToken,gatewayToken,requestIP,requestPort,param,"","");
System.out.println(pluginBeanJobInfo); System.out.println(pluginBeanJobInfo);
//PluginBeanJobInfo javaBeanJobInfo1 = new PluginBeanJobInfo(name,plServerUrl,mythCron,routingStrategy,blockingStrategy,callbackToken,gatewayToken,requestIP,requestPort,param); PluginBeanJobInfo javaBeanJobInfo1 = new PluginBeanJobInfo(name,plServerUrl,mythCron,routingStrategy,blockingStrategy,callbackToken,gatewayToken,requestIP,requestPort,param,"","");
//javaBeanJobInfo1.setJobHandelName("addJob1"); javaBeanJobInfo1.setJobHandelName("addJob1");
//JobUtils.addJob(javaBeanJobInfo1); javaBeanJobInfo1.setParam("延迟任务");
//JobUtils.addJob(pluginBeanJobInfo); JobUtils.addJob(javaBeanJobInfo1);
JobUtils.addJob(pluginBeanJobInfo);
} }
} }
/*
Navicat Premium Data Transfer
Source Server : 10.0.10.118
Source Server Type : MySQL
Source Server Version : 50721
Source Host : 10.0.10.118:3306
Source Schema : myth-job
Target Server Type : MySQL
Target Server Version : 50721
File Encoding : 65001
Date: 13/12/2019 10:20:21
*/
SET NAMES utf8mb4;
SET FOREIGN_KEY_CHECKS = 0;
-- ----------------------------
-- Table structure for flow_version
-- ----------------------------
DROP TABLE IF EXISTS `flow_version`;
CREATE TABLE `flow_version` (
`id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '任务流的版本的id',
`version_name` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '版本名称',
`node_count` int(13) NULL DEFAULT NULL COMMENT '当前版本的任务流的节点数',
`flow_desc` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '当前任务流的介绍',
`add_time` datetime(0) NULL DEFAULT NULL COMMENT '当前版本的添加时间',
`version_sign` varchar(30) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '当前版本的标志 this',
`flow_id` varchar(32) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '任务流的id',
`timeout` bigint(20) NULL DEFAULT NULL COMMENT '任务流的超时时间',
`trigger_next_time` bigint(13) NULL DEFAULT NULL COMMENT '此版本下次执行的时间',
`repeat_count` int(13) NULL DEFAULT NULL COMMENT '重复次数',
`remaining_count` varchar(13) 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;
SET FOREIGN_KEY_CHECKS = 1;
/*
Navicat Premium Data Transfer
Source Server : 10.0.10.118
Source Server Type : MySQL
Source Server Version : 50721
Source Host : 10.0.10.118:3306
Source Schema : myth-job
Target Server Type : MySQL
Target Server Version : 50721
File Encoding : 65001
Date: 13/12/2019 10:20:27
*/
SET NAMES utf8mb4;
SET FOREIGN_KEY_CHECKS = 0;
-- ----------------------------
-- Table structure for job_flight_schedule
-- ----------------------------
DROP TABLE IF EXISTS `job_flight_schedule`;
CREATE TABLE `job_flight_schedule` (
`id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '主键 任务id',
`job_name` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '任务名称',
`executor_handler` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '插件任务的调度key',
`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',
`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 '任务参数',
`job_type` varchar(32) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '任务类型,java,pyton,php,script,sql.shell',
`taskflow_id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '所属任务流的id',
`parent_id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '上级节点',
`taskflow_is_cron` char(1) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '是否跟随任务流的时间设置?1不跟随(默认),2跟随',
`alarm_email` varchar(256) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '报警邮件',
`executor_timeout` int(11) NULL DEFAULT NULL COMMENT '任务超时时间',
`source_urls` varchar(512) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '脚本文件地址(或远程文件服务器地址)支持多个,逗号分割',
`glue_source` longtext CHARACTER SET utf8 COLLATE utf8_general_ci NULL COMMENT '调度任务源码',
`glue_remark` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '源码备注',
`glue_updatetime` datetime(0) NULL DEFAULT NULL COMMENT '源码更新时间',
`trigger_next_time` bigint(20) NULL DEFAULT NULL COMMENT '下次调度时间',
`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,
PRIMARY KEY (`id`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8 COLLATE = utf8_general_ci ROW_FORMAT = Dynamic;
SET FOREIGN_KEY_CHECKS = 1;
/*
Navicat Premium Data Transfer
Source Server : 10.0.10.118
Source Server Type : MySQL
Source Server Version : 50721
Source Host : 10.0.10.118:3306
Source Schema : myth-job
Target Server Type : MySQL
Target Server Version : 50721
File Encoding : 65001
Date: 13/12/2019 10:20:34
*/
SET NAMES utf8mb4;
SET FOREIGN_KEY_CHECKS = 0;
-- ----------------------------
-- Table structure for job_flow
-- ----------------------------
DROP TABLE IF EXISTS `job_flow`;
CREATE TABLE `job_flow` (
`id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '任务流的id',
`flow_name` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '任务流的名字',
`add_time` datetime(0) NULL DEFAULT NULL COMMENT '任务流的创建时间',
`update_time` datetime(0) NULL DEFAULT NULL COMMENT '任务流的修改时间',
`remove_mark` char(1) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '删除标识 1正常 2删除',
PRIMARY KEY (`id`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8 COLLATE = utf8_general_ci ROW_FORMAT = Dynamic;
SET FOREIGN_KEY_CHECKS = 1;
/*
Navicat Premium Data Transfer
Source Server : 10.0.10.118
Source Server Type : MySQL
Source Server Version : 50721
Source Host : 10.0.10.118:3306
Source Schema : myth-job
Target Server Type : MySQL
Target Server Version : 50721
File Encoding : 65001
Date: 13/12/2019 10:20:42
*/
SET NAMES utf8mb4;
SET FOREIGN_KEY_CHECKS = 0;
-- ----------------------------
-- Table structure for job_flow_nodes
-- ----------------------------
DROP TABLE IF EXISTS `job_flow_nodes`;
CREATE TABLE `job_flow_nodes` (
`id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '主键 任务id',
`job_name` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '任务名称',
`executor_handler` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '插件任务的调度key',
`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',
`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 '任务参数',
`job_type` varchar(32) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '任务类型,java,pyton,php,script,sql.shell',
`taskflow_id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '所属任务流的id',
`parent_id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '上级节点',
`taskflow_is_cron` char(1) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '是否跟随任务流的时间设置?1不跟随(默认),2跟随',
`alarm_email` varchar(256) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '报警邮件',
`executor_timeout` int(11) NULL DEFAULT NULL COMMENT '任务超时时间',
`source_urls` varchar(512) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '脚本文件地址(或远程文件服务器地址)支持多个,逗号分割',
`glue_source` longtext CHARACTER SET utf8 COLLATE utf8_general_ci NULL COMMENT '调度任务源码',
`glue_remark` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '源码备注',
`glue_updatetime` datetime(0) NULL DEFAULT NULL COMMENT '源码更新时间',
`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 '任务节点修改时间',
`removal_mark` char(1) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '删除标志',
`repeat_times` int(13) NOT NULL COMMENT '重复次数 -1永久运行',
`flow_version` 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;
SET FOREIGN_KEY_CHECKS = 1;
/*
Navicat Premium Data Transfer
Source Server : 10.0.10.118
Source Server Type : MySQL
Source Server Version : 50721
Source Host : 10.0.10.118:3306
Source Schema : myth-job
Target Server Type : MySQL
Target Server Version : 50721
File Encoding : 65001
Date: 13/12/2019 10:20:47
*/
SET NAMES utf8mb4;
SET FOREIGN_KEY_CHECKS = 0;
-- ----------------------------
-- Table structure for job_flow_snapshot
-- ----------------------------
DROP TABLE IF EXISTS `job_flow_snapshot`;
CREATE TABLE `job_flow_snapshot` (
`id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '任务流的id',
`flow_name` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '任务流的名字',
`node_count` int(13) NULL DEFAULT NULL COMMENT '任务流中节点的个数',
`flow_version_id` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '当前任务流的版本',
`timeout` bigint(20) NULL DEFAULT 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 COMMENT '任务流的运行标识',
`remaining_node` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '初始值为节点总数,每次一个节点运行成功就将总数-1',
`flow_status` char(1) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '1 未开始 2运行 3暂停 4结束',
`trigger_next_time` bigint(13) NOT NULL DEFAULT 0 COMMENT '下次调度时间',
`flow_run_result` char(1) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '1 成功 2 失败',
`dispatch_ip` 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;
SET FOREIGN_KEY_CHECKS = 1;
/*
Navicat Premium Data Transfer
Source Server : 10.0.10.118
Source Server Type : MySQL
Source Server Version : 50721
Source Host : 10.0.10.118:3306
Source Schema : myth-job
Target Server Type : MySQL
Target Server Version : 50721
File Encoding : 65001
Date: 13/12/2019 10:20:55
*/
SET NAMES utf8mb4;
SET FOREIGN_KEY_CHECKS = 0;
-- ----------------------------
-- Table structure for job_lock
-- ----------------------------
DROP TABLE IF EXISTS `job_lock`;
CREATE TABLE `job_lock` (
`lock_name` varchar(50) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL,
`locl_desc` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL,
PRIMARY KEY (`lock_name`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8 COLLATE = utf8_general_ci ROW_FORMAT = Dynamic;
SET FOREIGN_KEY_CHECKS = 1;
/*
Navicat Premium Data Transfer
Source Server : 10.0.10.118
Source Server Type : MySQL
Source Server Version : 50721
Source Host : 10.0.10.118:3306
Source Schema : myth-job
Target Server Type : MySQL
Target Server Version : 50721
File Encoding : 65001
Date: 13/12/2019 10:21:00
*/
SET NAMES utf8mb4;
SET FOREIGN_KEY_CHECKS = 0;
-- ----------------------------
-- Table structure for job_log
-- ----------------------------
DROP TABLE IF EXISTS `job_log`;
CREATE TABLE `job_log` (
`id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '日志主键',
`job_group` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '执行器主键',
`job_id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '任务节点主键',
`flow_id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '任务流主键',
`flow_version` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '任务流版本',
`run_id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '运行标识',
`executor_handler` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '执行器任务handler',
`executor_params` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '执行器的参数',
`trigger_time` datetime(0) NULL DEFAULT NULL COMMENT '调度时间',
`trigger_code` char(1) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '调度结果 1成功 2失败',
`trigger_msg` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '调度日志',
`handle_time` datetime(0) NULL DEFAULT NULL COMMENT '执行时间',
`handle_code` char(1) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '执行结果1成功 2失败',
`handle_msg` text CHARACTER SET utf8 COLLATE utf8_general_ci NULL COMMENT '执行日志',
`alarm_status` char(1) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '告警状态 0-默认 1-无需警告 2-告警成功 3-告警失败',
PRIMARY KEY (`id`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8 COLLATE = utf8_general_ci ROW_FORMAT = Dynamic;
SET FOREIGN_KEY_CHECKS = 1;
/*
Navicat Premium Data Transfer
Source Server : 10.0.10.118
Source Server Type : MySQL
Source Server Version : 50721
Source Host : 10.0.10.118:3306
Source Schema : myth-job
Target Server Type : MySQL
Target Server Version : 50721
File Encoding : 65001
Date: 13/12/2019 10:21:04
*/
SET NAMES utf8mb4;
SET FOREIGN_KEY_CHECKS = 0;
-- ----------------------------
-- Table structure for job_logglue
-- ----------------------------
DROP TABLE IF EXISTS `job_logglue`;
CREATE TABLE `job_logglue` (
`id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '历史源码主键',
`job_node_id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '节点主键',
`glue_source` longtext CHARACTER SET utf8 COLLATE utf8_general_ci NULL COMMENT '源代码',
`glue_remark` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '备注',
`add_time` datetime(0) NULL DEFAULT NULL COMMENT '添加时间',
PRIMARY KEY (`id`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8 COLLATE = utf8_general_ci ROW_FORMAT = Dynamic;
SET FOREIGN_KEY_CHECKS = 1;
/*
Navicat Premium Data Transfer
Source Server : 10.0.10.118
Source Server Type : MySQL
Source Server Version : 50721
Source Host : 10.0.10.118:3306
Source Schema : myth-job
Target Server Type : MySQL
Target Server Version : 50721
File Encoding : 65001
Date: 13/12/2019 10:21:09
*/
SET NAMES utf8mb4;
SET FOREIGN_KEY_CHECKS = 0;
-- ----------------------------
-- Table structure for job_read_ahead
-- ----------------------------
DROP TABLE IF EXISTS `job_read_ahead`;
CREATE TABLE `job_read_ahead` (
`id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '主键 任务id',
`job_name` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '任务名称',
`executor_handler` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '插件任务的调度key',
`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',
`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 '任务参数',
`job_type` varchar(32) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '任务类型,java,pyton,php,script,sql.shell',
`taskflow_id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '所属任务流的id',
`parent_id` varchar(36) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '上级节点',
`taskflow_is_cron` char(1) CHARACTER SET utf8 COLLATE utf8_general_ci NOT NULL COMMENT '是否跟随任务流的时间设置?1不跟随(默认),2跟随',
`alarm_email` varchar(256) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '报警邮件',
`executor_timeout` int(11) NULL DEFAULT NULL COMMENT '任务超时时间',
`source_urls` varchar(512) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '脚本文件地址(或远程文件服务器地址)支持多个,逗号分割',
`trigger_status` char(1) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '调度状态:0-暂停,1-运行',
`glue_source` longtext CHARACTER SET utf8 COLLATE utf8_general_ci NULL COMMENT '调度任务源码',
`glue_remark` varchar(255) CHARACTER SET utf8 COLLATE utf8_general_ci NULL DEFAULT NULL COMMENT '源码备注',
`glue_updatetime` datetime(0) NULL DEFAULT NULL COMMENT '源码更新时间',
`trigger_next_time` bigint(20) NULL DEFAULT NULL COMMENT '下次调度时间',
`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 COMMENT '运行的唯一标识',
PRIMARY KEY (`id`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8 COLLATE = utf8_general_ci ROW_FORMAT = Dynamic;
SET FOREIGN_KEY_CHECKS = 1;
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment