Commit d2dc8846 by huangfusuper

增加运行日志本地记录,任务日志全异步记录

parent 7d4d61ff
package com.byit.controller;
import com.byit.job.dto.PluginJobRunResultDto;
import com.byit.job.model.MythJobReadAhead;
import com.byit.conf.MythJobAutoConfigure;
import com.byit.job.dto.JobRunResultDto;
import com.byit.job.dto.PluginBeanJobInfo;
import com.byit.job.model.MythJobReadAhead;
import com.byit.job.utils.SourceObj2TargetObjUtil;
import com.byit.job.vo.ReturnResult;
import com.byit.service.JobReadAheadService;
import com.byit.thread.LogCallbackThread;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.*;
......@@ -22,6 +23,7 @@ import java.util.List;
public class JobController {
@Autowired
private JobReadAheadService jobReadAheadService;
@PostMapping(value = "addJob")
public String addJob(@RequestBody PluginBeanJobInfo pluginBeanJobInfo){
MythJobReadAhead mythJobReadAhead = SourceObj2TargetObjUtil.pluginBeanJobInfo2MythJobReadAhead(pluginBeanJobInfo);
......@@ -30,9 +32,8 @@ public class JobController {
}
@PostMapping(value = "callbackRes")
public String callbackRes(@RequestBody PluginJobRunResultDto pluginJobRunResultDto){
System.out.println(pluginJobRunResultDto.getReturnResult().getCode()+"-----"+pluginJobRunResultDto.getReturnResult().getMsg());
return "好的,我知道你执行成功了";
public void callbackRes(@RequestBody JobRunResultDto jobRunResultDto){
MythJobAutoConfigure.LOG_CALLBACK.execute(new LogCallbackThread(jobRunResultDto));
}
@GetMapping(value = "getJobInfo")
......
......@@ -7,3 +7,10 @@ spring:
mybatis:
mapper-locations: /mapper/*.xml
logging:
path: /data/mythjob
file: myth_log_file
<?xml version="1.0" encoding="UTF-8"?>
<configuration>
<appender name="console" class="ch.qos.logback.core.ConsoleAppender">
<!-- <filter class="ch.qos.logback.classic.filter.ThresholdFilter">
<level>ERROR</level>
</filter>-->
<encoder>
<pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{50} - %msg%n</pattern>
</encoder>
</appender>
<appender name="file-debug" class="ch.qos.logback.core.rolling.RollingFileAppender">
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<fileNamePattern>${LOG_PATH}/${LOG_FILE}_debug.%d{yyyy-MM-dd}.log</fileNamePattern>
</rollingPolicy>
<encoder>
<pattern>%d{HH:mm:ss.SSS} %contextName [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
<filter class="ch.qos.logback.classic.filter.LevelFilter">
<level>DEBUG</level>
<onMatch>ACCEPT</onMatch>
<onMismatch>DENY</onMismatch>
</filter>
</appender>
<appender name="file-info" class="ch.qos.logback.core.rolling.RollingFileAppender">
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<fileNamePattern>${LOG_PATH}/${LOG_FILE}_info.%d{yyyy-MM-dd}.log</fileNamePattern>
</rollingPolicy>
<encoder>
<pattern>%d{HH:mm:ss.SSS} %contextName [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
<filter class="ch.qos.logback.classic.filter.LevelFilter">
<level>INFO</level>
<onMatch>ACCEPT</onMatch>
<onMismatch>DENY</onMismatch>
</filter>
</appender>
<appender name="file-error" class="ch.qos.logback.core.rolling.RollingFileAppender">
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<fileNamePattern>${LOG_PATH}/${LOG_FILE}_error.%d{yyyy-MM-dd}.log</fileNamePattern>
</rollingPolicy>
<encoder>
<pattern>%d{HH:mm:ss.SSS} %contextName [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
<filter class="ch.qos.logback.classic.filter.LevelFilter">
<level>ERROR</level>
<onMatch>ACCEPT</onMatch>
<onMismatch>DENY</onMismatch>
</filter>
</appender>
<appender name="file-myth-job" class="ch.qos.logback.core.rolling.RollingFileAppender">
<rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
<fileNamePattern>${LOG_PATH}/${LOG_FILE}_myth_job.%d{yyyy-MM-dd}.log</fileNamePattern>
</rollingPolicy>
<encoder>
<pattern>%d{HH:mm:ss.SSS} %contextName [%thread] %-5level %logger{36} - %msg%n</pattern>
</encoder>
</appender>
<logger name="com.byit.thread" level="debug" additivity="false">
<appender-ref ref="file-myth-job" />
</logger>
<logger name="org.springframework.web" level="DEBUG"/>
<logger name="com.ibatis" level="DEBUG" />
<logger name="com.ibatis.common.jdbc.SimpleDataSource" level="DEBUG" />
<logger name="com.ibatis.common.jdbc.ScriptRunner" level="DEBUG" />
<logger name="com.ibatis.sqlmap.engine.impl.SqlMapClientDelegate" level="DEBUG" />
<logger name="java.sql.Connection" level="DEBUG" />
<logger name="java.sql.Statement" level="DEBUG" />
<logger name="java.sql.PreparedStatement" level="DEBUG" />
<root level="info">
<appender-ref ref="console" />
<appender-ref ref="file-debug" />
<appender-ref ref="file-info" />
<appender-ref ref="file-error" />
</root>
</configuration>
\ No newline at end of file
......@@ -62,8 +62,6 @@
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>
</dependencies>
</project>
\ No newline at end of file
......@@ -3,6 +3,10 @@ package com.byit.conf;
import org.mybatis.spring.annotation.MapperScan;
import org.springframework.context.annotation.*;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
/**
* @program: byit-myth-job->AppConf
* @description: 稳健配置
......@@ -12,4 +16,20 @@ import org.springframework.context.annotation.*;
@Configuration
@MapperScan("com.byit.dao")
public class MythJobAutoConfigure {
/**
* 日志的回调线程池,主要是将执行结果写到库里面
* LinkedBlockingQueue 不指定容量就变成了无界队列
* 判断核心线程数是否已满,核心线程数大小和corePoolSize参数有关,未满则创建线程执行任务
* 若核心线程池已满,判断队列是否满,队列是否满和workQueue参数有关,若未满则加入队列中
* 若队列已满,判断线程池是否已满,线程池是否已满和maximumPoolSize参数有关,若未满创建线程执行任务
* 若线程池已满,则采用拒绝策略处理无法执执行的任务,拒绝策略和handler参数有关
*/
public static final ThreadPoolExecutor LOG_CALLBACK = new ThreadPoolExecutor(
10,
100,
60L,
TimeUnit.SECONDS,
new LinkedBlockingQueue<Runnable>(100),
r ->new Thread(r, "MythJob Thread of Job Run Log Callback Warehouse-" + r.hashCode()));
}
package com.byit.dao;
import com.byit.job.model.MythJobLog;
import org.apache.ibatis.annotations.Param;
import org.springframework.stereotype.Repository;
import java.util.List;
/**
* @program: byit-myth-job->JobLogMapper
* @description: 任务日志的操作
* @author: huangfu
* @date: 2019/12/16 14:35
**/
@Repository
public interface JobLogMapper {
/**
* 根据运行id和所有的父类id 查询没有完成的数据
* 后期查询一个节点的父节点是否完成时只需要根据运行标识和父节点id查询 返回结果为null的情况下成立
* @param runId
* @param parenIds
* @return
*/
List<MythJobLog> findUndoneMythJobLogByParenIdAndInRunId(@Param("runId") String runId, @Param("parenIds") List<String> parenIds);
/**
* 修改一条日志数据
* @param mythJobLog
*/
void updateOneMythJobLog(MythJobLog mythJobLog);
/**
* 添加一条数据
* @param mythJobLog
*/
void addOneMythJobLog(MythJobLog mythJobLog);
}
......@@ -3,7 +3,6 @@ package com.byit.job;
import io.netty.util.HashedWheelTimer;
import io.netty.util.TimerTask;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.TimeUnit;
/**
......@@ -15,22 +14,17 @@ import java.util.concurrent.TimeUnit;
public class WorkRoulette {
/**
* hashedWheelTimer:工作轮盘
* HASHED_WHEEL_TIMER:工作轮盘
* ThreadFactory:创建work线程
* tickDuration:每个刻度的时间
* ticsPerWheel:轮盘一圈大小
* maxPendingTimeouts:最大等待处理超时
*/
private static final HashedWheelTimer hashedWheelTimer = new HashedWheelTimer(new ThreadFactory() {
@Override
public Thread newThread(Runnable r) {
return new Thread(r, "hashedWheelTimer" + r.hashCode());
}
}, 1, TimeUnit.SECONDS, 8, true, 0);
private static final HashedWheelTimer HASHED_WHEEL_TIMER = new HashedWheelTimer(r -> new Thread(r, "HASHED_WHEEL_TIMER" + r.hashCode()), 1, TimeUnit.SECONDS, 8, true, 0);
public static void addJob(TimerTask timerTask) {
hashedWheelTimer.newTimeout(timerTask, TimeUnit.SECONDS.toNanos(20), TimeUnit.NANOSECONDS);
public static void addJob(TimerTask timerTask,long triggerNextTime) {
HASHED_WHEEL_TIMER.newTimeout(timerTask, TimeUnit.MILLISECONDS.toNanos(triggerNextTime-System.currentTimeMillis()), TimeUnit.NANOSECONDS);
}
}
package com.byit.service;
import com.byit.job.dto.JobRunResultDto;
import com.byit.job.model.MythJobLog;
import java.util.List;
/**
* @program: byit-myth-job->JobLogService
* @description: 日志表操作
* @author: huangfu
* @date: 2019/12/16 17:01
**/
public interface JobLogService {
/**
* 根据运行id和所有的父类id 查询没有完成的数据
* 后期查询一个节点的父节点是否完成时只需要根据运行标识和父节点id查询 返回结果为null的情况下成立
* @param runId
* @param parenIds
* @return
*/
List<MythJobLog> findUndoneMythJobLogByParenIdAndInRunId(String runId,List<String> parenIds);
/**
* 修改一条日志数据
* @param jobRunResultDto
*/
void updateOneMythJobLog(JobRunResultDto jobRunResultDto);
/**
* 添加一条数据
* @param mythJobLog
*/
void addOneMythJobLog(MythJobLog mythJobLog);
}
package com.byit.service.impl;
import com.byit.dao.JobLogMapper;
import com.byit.job.dto.JobRunResultDto;
import com.byit.job.model.MythJobLog;
import com.byit.service.JobLogService;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.util.Date;
import java.util.List;
/**
* @program: byit-myth-job->JobLogServiceImpl
* @description: 日志表操作
* @author: huangfu
* @date: 2019/12/16 17:02
**/
@Service
public class JobLogServiceImpl implements JobLogService {
@Autowired
private JobLogMapper jobLogMapper;
@Override
public List<MythJobLog> findUndoneMythJobLogByParenIdAndInRunId(String runId, List<String> parenIds) {
return jobLogMapper.findUndoneMythJobLogByParenIdAndInRunId(runId,parenIds);
}
@Override
public void updateOneMythJobLog(JobRunResultDto jobRunResultDto) {
MythJobLog mythJobLog = new MythJobLog();
mythJobLog.setId(jobRunResultDto.getLogId());
mythJobLog.setHandleTime(new Date());
mythJobLog.setHandleCode(jobRunResultDto.getReturnResult().getCode());
mythJobLog.setHandleMsg(jobRunResultDto.getReturnResult().getMsg());
jobLogMapper.updateOneMythJobLog(mythJobLog);
}
@Override
public void addOneMythJobLog(MythJobLog mythJobLog) {
jobLogMapper.addOneMythJobLog(mythJobLog);
}
}
......@@ -28,6 +28,8 @@ public class JobReadAheadServiceImpl implements JobReadAheadService {
@Override
public void addOneMythJobReadAhead(MythJobReadAhead mythJobReadAhead) {
jobInfoMapper.addOneMythJobReadAhead(mythJobReadAhead);
}
......
......@@ -5,13 +5,19 @@ import com.alibaba.fastjson.JSON;
import com.byit.job.dto.AdminSenPluginDto;
import com.byit.job.exceptions.plugin.PluginException;
import com.byit.job.model.MythJobFlightSchedule;
import com.byit.job.model.MythJobLog;
import com.byit.job.utils.IpUtil;
import com.byit.rpc.remoting.invoker.route.LoadBalance;
import com.byit.rpc.remoting.invoker.route.RpcLoadBalance;
import com.byit.service.impl.JobLogServiceImpl;
import com.byit.util.SpringUtil;
import io.netty.util.Timeout;
import io.netty.util.TimerTask;
import lombok.extern.slf4j.Slf4j;
import java.util.Date;
import java.util.UUID;
/**
* @program: byit-myth-job->JvavBeanJobTask
* @description: 开发者编写的java job
......@@ -34,13 +40,18 @@ public class JavaBeanJobTask implements TimerTask {
String jobHandelName = mythJobFlightSchedule.getExecutorHandler( );
String param = mythJobFlightSchedule.getJobParam( );
String runId = mythJobFlightSchedule.getRunID();
String logId = UUID.randomUUID( ).toString().replace("-","");
AdminSenPluginDto adminSenPluginDto = new AdminSenPluginDto();
adminSenPluginDto.setCallbackUrl("http://127.0.0.1:8080/job/callbackRes");
adminSenPluginDto.setJobHandelName(jobHandelName);
adminSenPluginDto.setJobParam(param);
adminSenPluginDto.setRunId(runId);
adminSenPluginDto.setLogId(logId);
//这里获取的是调度结果
String result = HttpUtil.post(url, JSON.toJSONString(adminSenPluginDto),10*1000);
saveLog(mythJobFlightSchedule,logId,url,result);
log.debug("---------------{}------------",result);
}catch (PluginException ignored){
log.error("-------------通讯异常:{},{}",ignored.getIEnum().getCode(),ignored.getIEnum().getMsg());
......@@ -48,4 +59,30 @@ public class JavaBeanJobTask implements TimerTask {
e.printStackTrace();
}
}
private void saveLog(MythJobFlightSchedule mythJobFlightSchedule,String logId,String url,String result){
MythJobLog mythJobLog = new MythJobLog( );
//设置邮件
mythJobLog.setId(logId);
mythJobLog.setJobGroup(url);
mythJobLog.setJobId(mythJobFlightSchedule.getId());
mythJobLog.setFlowId(mythJobFlightSchedule.getTaskFlowId());
mythJobLog.setRunId(mythJobFlightSchedule.getRunID());
mythJobLog.setExecutorHandler(mythJobFlightSchedule.getExecutorHandler());
mythJobLog.setExecutorParams(mythJobFlightSchedule.getJobParam());
//TODO 有问题
mythJobLog.setTriggerCode("1");
mythJobLog.setTriggerMsg(result);
mythJobLog.setTriggerTime(new Date());
mythJobLog.setJobName(mythJobFlightSchedule.getJobName());
mythJobLog.setAlarmEmail(mythJobFlightSchedule.getAlarmEmail());
JobLogServiceImpl bean = SpringUtil.getBean(JobLogServiceImpl.class);
bean.addOneMythJobLog(mythJobLog);
}
}
......@@ -240,7 +240,7 @@ public class JobScheduleHelper{
if ("BEAN".equals(mythJobFlightScheduleByTriggerNextTime.getJobType())) {
JavaBeanJobTask javaBeanJobTask = new JavaBeanJobTask(mythJobFlightScheduleByTriggerNextTime);
jobFlightScheduleService.deleteMythJobFlightScheduleById(mythJobFlightScheduleByTriggerNextTime.getId());
WorkRoulette.addJob(javaBeanJobTask);
WorkRoulette.addJob(javaBeanJobTask,mythJobFlightScheduleByTriggerNextTime.getTriggerNextTime());
}
});
}else{
......
package com.byit.thread;
import com.byit.job.dto.JobRunResultDto;
import com.byit.service.JobLogService;
import com.byit.service.impl.JobLogServiceImpl;
import com.byit.util.SpringUtil;
import lombok.extern.slf4j.Slf4j;
/**
* @program: byit-myth-job->CallbackThread
* @description: 日志回调线程
* @author: huangfu
* @date: 2019/12/17 12:04
**/
@Slf4j
public class LogCallbackThread implements Runnable {
private JobRunResultDto jobRunResultDto;
public LogCallbackThread(JobRunResultDto jobRunResultDto) {
this.jobRunResultDto = jobRunResultDto;
}
@Override
public void run() {
log.debug("--------------------任务执行完成---------------------");
JobLogService jobLogService = SpringUtil.getBean(JobLogServiceImpl.class);
jobLogService.updateOneMythJobLog(jobRunResultDto);
}
}
......@@ -23,7 +23,6 @@
<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>
......
<?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.JobLogMapper">
<sql id="Base_Column_List">
t.id,
t.job_group,
t.job_id,
t.flow_id,
t.run_id,
t.executor_handler,
t.executor_params,
t.trigger_time,
t.trigger_code,
t.trigger_msg,
t.handle_time,
t.handle_code,
t.handle_msg,
t.alarm_status,
t.job_name,
t.alarm_email
</sql>
<resultMap id="mythJobLog" type="com.byit.job.model.MythJobLog">
<id column="id" property="id"/>
<result column="job_group" property="jobGroup"/>
<result column="job_id" property="jobId"/>
<result column="flow_id" property="flowId"/>
<result column="run_id" property="runId"/>
<result column="executor_handler" property="executorHandler"/>
<result column="executor_params" property="executorParams"/>
<result column="trigger_time" property="triggerTime"/>
<result column="trigger_code" property="triggerCode"/>
<result column="trigger_msg" property="triggerMsg"/>
<result column="handle_time" property="handleTime"/>
<result column="handle_code" property="handleCode"/>
<result column="handle_msg" property="handleMsg"/>
<result column="alarm_status" property="alarmStatus"/>
<result column="job_name" property="jobName"/>
<result column="alarm_email" property="alarmEmail"/>
</resultMap>
<select id="findUndoneMythJobLogByParenIdAndInRunId" resultMap="mythJobLog">
SELECT <include refid="Base_Column_List"/> FROM job_log t
<where>
run_id = #{runId} AND handle_code = '0' AND
job_id IN
<foreach collection="parenIds" item="jobId" index="index" open="(" close=")" separator=",">
#{jobId}
</foreach>
</where>
</select>
<update id="updateOneMythJobLog" parameterType="com.byit.job.model.MythJobLog">
UPDATE job_log
<trim prefix="SET" suffixOverrides=",">
<if test="jobGroup!= null and jobGroup != ''">job_group=#{jobGroup},</if>
<if test="jobId!= null and jobId != ''">job_id=#{jobId},</if>
<if test="flowId!= null and flowId != ''">flow_id=#{flowId},</if>
<if test="runId!= null and runId != ''">run_id=#{runId},</if>
<if test="executorHandler!= null and executorHandler != ''">executor_handler=#{executorHandler},</if>
<if test="executorParams!= null and executorParams != ''">executor_params=#{executorParams},</if>
<if test="triggerTime != null">trigger_time=#{triggerTime},</if>
<if test="triggerCode!= null and triggerCode != ''">trigger_code=#{triggerCode},</if>
<if test="triggerMsg!= null and triggerMsg != ''">trigger_msg=#{triggerMsg},</if>
<if test="handleTime!= null">handle_time=#{handleTime},</if>
<if test="handleCode!= null and handleCode != ''">handle_code=#{handleCode},</if>
<if test="handleMsg!= null and handleMsg != ''">handle_msg=#{handleMsg},</if>
<if test="alarmStatus!= null and alarmStatus != ''">alarm_status=#{alarmStatus},</if>
<if test="jobName!= null and jobName != ''">job_name=#{jobName},</if>
<if test="alarmEmail!= null and alarmEmail != ''">alarm_email=#{alarmEmail},</if>
</trim>
WHERE id = #{id}
</update>
<insert id="addOneMythJobLog" parameterType="com.byit.job.model.MythJobLog">
INSERT INTO job_log (
id,job_group,job_id,flow_id,run_id,executor_handler,
executor_params,trigger_time,trigger_code,trigger_msg,handle_time,
handle_code,handle_msg,alarm_status,job_name,alarm_email
)
values (
#{id},#{jobGroup},#{jobId},#{flowId},#{runId},#{executorHandler},#{executorParams},
#{triggerTime},#{triggerCode},#{triggerMsg},#{handleTime},#{handleCode},#{handleMsg},#{alarmStatus},
#{jobName},#{alarmEmail}
)
</insert>
</mapper>
\ No newline at end of file
......@@ -36,4 +36,8 @@ public class AdminSenPluginDto {
* 检测是否有心跳参数,有心跳参数则为测试参数,且为PENG的话,服务端回复 PONG
*/
private String heartbeat;
/**
* 日志ID
*/
private String logId;
}
......@@ -18,7 +18,7 @@ import java.util.Date;
@NoArgsConstructor
@AllArgsConstructor
@ToString
public class PluginJobRunResultDto {
public class JobRunResultDto {
/**
* 任务的运行标识
*/
......@@ -36,4 +36,6 @@ public class PluginJobRunResultDto {
*/
private Date endTime;
private String logId;
}
......@@ -11,7 +11,7 @@ public interface IEnum {
* 返码值
* @return
*/
int getCode();
String getCode();
/**
* 返回执行信息
......
......@@ -5,23 +5,23 @@ package com.byit.job.enums;
* @author huangfu
*/
public enum JobResultEnum implements IEnum {
SUCCESS(200,"任务执行成功"),
FAIL(500,"任务执行失败"),
FAIL_TIMEOUT(502,"超时错误");
private int code;
SUCCESS("200","任务执行成功"),
FAIL("500","任务执行失败"),
FAIL_TIMEOUT("502","超时错误");
private String code;
private String msg;
JobResultEnum() {
}
JobResultEnum(int code, String msg) {
JobResultEnum(String code, String msg) {
this.code = code;
this.msg = msg;
}
@Override
public int getCode() {
public String getCode() {
return this.code;
}
......
......@@ -8,18 +8,18 @@ import com.byit.job.enums.IEnum;
*/
public enum PluginEnum implements IEnum {
REQUEST_PORT_OR_IP_IS_MISSING("请求ip或者port为null",10000),
NO_SERVICE_AVAILABLE("无可用的服务",11000);
REQUEST_PORT_OR_IP_IS_MISSING("请求ip或者port为null","10000"),
NO_SERVICE_AVAILABLE("无可用的服务","11000");
private String msg;
private int code;
private String code;
PluginEnum(String msg, int code) {
PluginEnum(String msg, String code) {
this.msg = msg;
this.code = code;
}
@Override
public int getCode() {
public String getCode() {
return this.code;
}
......
package com.byit.job.model;
import lombok.*;
import java.util.Date;
/**
* @program: byit-myth-job->FlowVersion
* @description: 工作流版本
* @author: huangfu
* @date: 2019/12/17 16:52
**/
@AllArgsConstructor
@NoArgsConstructor
@Data
@ToString
@EqualsAndHashCode
public class FlowVersion {
private String id;
/**
* 版本名称
*/
private String versionName;
/**
* 本版本的节点数
*/
private Integer nodeCount;
/**
* 当前版本的介绍
*/
private String flowDesc;
/**
* 当前版本的添加是时间
*/
private Date addTime;
/**
* 当前版本的额标识
*/
private String versionSign;
/**
* 工作流id
*/
private String flowId;
/**
* 任务流的超时时间
*/
private Long timeout;
/**
* 此版本的下一次的执行时间
*/
private Long triggerNextTime;
/**
* 重复次数
*/
private Integer repeatCount;
/**
* 剩余次数
*/
private Integer remainingCount;
/**
* 报警邮箱
*/
private String alarmEmail;
}
......@@ -137,4 +137,5 @@ public class MythJobFlightSchedule {
* 运行标识
*/
private String runID;
}
package com.byit.job.model;
import lombok.*;
import java.util.Date;
/**
* @program: byit-myth-job->MythJobFlow
* @description: 任务流
* @author: huangfu
* @date: 2019/12/17 16:47
**/
@AllArgsConstructor
@NoArgsConstructor
@Data
@ToString
@EqualsAndHashCode
public class MythJobFlow {
private String id;
private String flowName;
private Date updateTime;
private Date addTime;
private String removeMark;
}
package com.byit.job.model;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.ToString;
import lombok.*;
import java.util.Date;
......@@ -17,6 +14,7 @@ import java.util.Date;
@AllArgsConstructor
@NoArgsConstructor
@ToString
@EqualsAndHashCode
public class MythJobFlowNodes {
/**
* 任务节点的id
......@@ -90,10 +88,6 @@ public class MythJobFlowNodes {
*/
private String parentId;
/**
* 是否跟随任务流的时间设置?1不跟随(默认),2跟随
*/
private String taskFlowIsCron;
/**
* 报警邮件
*/
private String alarmEmail;
......
package com.byit.job.model;
import lombok.*;
import java.util.Date;
/**
* @program: byit-myth-job->MythJobLog
* @description: 任务日志表实体
* @author: huangfu
* @date: 2019/12/16 14:22
**/
@Data
@AllArgsConstructor
@NoArgsConstructor
@ToString
@EqualsAndHashCode
public class MythJobLog {
private String id;
/**
* 执行器的id
*/
private String jobGroup;
/**
* 任务节点主键
*/
private String jobId;
/**
* 任务流id
*/
private String flowId;
/**
* 运行标识
*/
private String runId;
/**
* 执行器任务handler
*/
private String executorHandler;
/**
* 执行器的参数
*/
private String executorParams;
/**
*调度时间
*/
private Date triggerTime;
/**
* 调度结果 1成功 2失败
*/
private String triggerCode;
/**
* 调度日志
*/
private String triggerMsg;
/**
*执行时间
*/
private Date handleTime;
/**
* 执行结果0未执行完成 1成功 2失败
*/
private String handleCode;
/**
* 执行日志
*/
private String handleMsg;
/**
*告警状态 1- 默认无需警告 2-告警成功 3-告警失败
*/
private String alarmStatus;
/**
* 任务名称
*/
private String jobName;
/**
* 告警邮箱
*/
private String alarmEmail;
}
package com.byit.job.model;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.ToString;
import lombok.*;
import java.util.Date;
......@@ -17,6 +14,7 @@ import java.util.Date;
@AllArgsConstructor
@NoArgsConstructor
@ToString
@EqualsAndHashCode
public class MythJobReadAhead {
/**
* 任务节点的id
......
......@@ -42,7 +42,7 @@ public class SourceObj2TargetObjUtil {
mythJobReadAhead.setTaskFlowIsCron("1");
//TODO 暂时随机分配 未来这个东西是线程自动添加的
mythJobReadAhead.setRunID(UUID.randomUUID( ).toString().replace("-",""));
mythJobReadAhead.setTriggerNextTime(System.currentTimeMillis()+60000);
mythJobReadAhead.setTriggerNextTime(System.currentTimeMillis()+20000);
return mythJobReadAhead;
}
}
......@@ -21,7 +21,7 @@ public class ReturnResult<T> implements Serializable {
public static final ReturnResult SUCCESS = new ReturnResult(null);
public static final ReturnResult FAIL = new ReturnResult(JobResultEnum.FAIL.getCode(), JobResultEnum.FAIL.getMsg());
public static final ReturnResult FAIL_TIMEOUT = new ReturnResult(JobResultEnum.FAIL_TIMEOUT.getCode(),JobResultEnum.FAIL_TIMEOUT.getMsg());
private int code;
private String code;
private String msg;
private T content;
......@@ -30,7 +30,7 @@ public class ReturnResult<T> implements Serializable {
* @param code
* @param msg
*/
public ReturnResult(int code, String msg) {
public ReturnResult(String code, String msg) {
this.code = code;
this.msg = msg;
}
......
......@@ -2,7 +2,7 @@ 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.dto.JobRunResultDto;
import com.byit.job.handler.interfaces.IJobHandler;
import com.byit.job.vo.ReturnResult;
import com.byit.utils.JobUtils;
......@@ -27,20 +27,21 @@ public class RunJobThread implements Runnable {
@Override
public void run() {
//创建回复对象
PluginJobRunResultDto pluginJobRunResultDto = new PluginJobRunResultDto();
pluginJobRunResultDto.setStartTime(new Date());
JobRunResultDto jobRunResultDto = new JobRunResultDto();
jobRunResultDto.setStartTime(new Date());
//设定运行标识
pluginJobRunResultDto.setJobRunId(adminSenPluginDto.getRunId());
jobRunResultDto.setJobRunId(adminSenPluginDto.getRunId());
//运行任务
ReturnResult<String> stringReturnResult = runJob(adminSenPluginDto.getJobHandelName(), adminSenPluginDto.getJobParam());
//设置运行结果
pluginJobRunResultDto.setReturnResult(stringReturnResult);
jobRunResultDto.setReturnResult(stringReturnResult);
//获取回调通知URL
String callbackUrl = adminSenPluginDto.getCallbackUrl();
//设置结束时间
pluginJobRunResultDto.setEndTime(new Date());
cn.hutool.http.HttpUtil.post(callbackUrl, JSON.toJSONString(pluginJobRunResultDto));
log.info("---------服务器端:{}-------------",pluginJobRunResultDto);
jobRunResultDto.setEndTime(new Date());
jobRunResultDto.setLogId(adminSenPluginDto.getLogId());
cn.hutool.http.HttpUtil.post(callbackUrl, JSON.toJSONString(jobRunResultDto));
log.info("---------服务器端:{}-------------", jobRunResultDto);
}
private ReturnResult<String> runJob(String jobHandlerName,String param){
......
......@@ -26,9 +26,11 @@ public class Mains {
String requestPort="8080";
String param="不延迟任务";
String name="addJob";
PluginBeanJobInfo pluginBeanJobInfo = new PluginBeanJobInfo(name,plServerUrl,mythCron,routingStrategy,blockingStrategy,callbackToken,gatewayToken,requestIP,requestPort,param,"","");
String email = "huangfusuper163.com";
String email1 = "huangfusuper@163.com";
PluginBeanJobInfo pluginBeanJobInfo = new PluginBeanJobInfo(name,plServerUrl,mythCron,routingStrategy,blockingStrategy,callbackToken,gatewayToken,requestIP,requestPort,param,email,"");
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,email1,"");
javaBeanJobInfo1.setJobHandelName("addJob1");
javaBeanJobInfo1.setParam("延迟任务");
JobUtils.addJob(javaBeanJobInfo1);
......
......@@ -59,6 +59,7 @@
<maven-source-plugin.version>3.1.0</maven-source-plugin.version>
<maven-javadoc-plugin.version>3.1.1</maven-javadoc-plugin.version>
<maven-gpg-plugin.version>1.6</maven-gpg-plugin.version>
<byit.validation.starter>1.0-SNAPSHOT</byit.validation.starter>
......@@ -161,6 +162,11 @@
<version>${mybatis.spring.boot.starter}</version>
</dependency>
<dependency>
<groupId>com.byit</groupId>
<artifactId>byit-validation-starter</artifactId>
<version>${byit.validation.starter}</version>
</dependency>
</dependencies>
</dependencyManagement>
......
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