Commit 22ef2245 by huangfusuper

rpc自动注册以及回调通知修改

parent a8ef595a
package com.byit;
import com.byit.annotations.EnablePluginClient;
import com.byit.rpc.remoting.provider.annotation.RpcService;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
......@@ -12,6 +13,7 @@ import org.springframework.boot.autoconfigure.SpringBootApplication;
**/
@SpringBootApplication
@RpcService(http_type = true)
@EnablePluginClient
public class AdminApplication {
public static void main(String[] args) {
SpringApplication.run(AdminApplication.class,args);
......
package com.byit.api;
import com.byit.conf.MythJobAutoConfigure;
import com.byit.dto.executor.JobRunResultDto;
import com.byit.packet.response.PluginRpcResponsePacket;
import com.byit.thread.JavaTaskCallbackThread;
import com.byit.thread.LogCallbackThread;
import io.swagger.annotations.Api;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
/**
* @author huangfu
*/
@Api(tags = "回调api")
@RestController
@RequestMapping("api/callback")
public class ApiCallbackController {
@PostMapping(value = "callbackRes")
public void callbackRes(@RequestBody PluginRpcResponsePacket pluginRpcResponsePacket){
MythJobAutoConfigure.LOG_CALLBACK.execute(new JavaTaskCallbackThread(pluginRpcResponsePacket));
}
}
......@@ -1226,10 +1226,10 @@ public class ApiFlowServiceImpl implements ApiFlowService {
flow.setStartUp(FlowPropertyEnum.IS_START.getCode());
//如果是周期调度修改下次执行时间
if (FlowPropertyEnum.SCHEDULE_MODE.getCode().equals(flow.getExecType())){
flow.setTriggerNextTime(new CronExpression(flow.getFlowCron()).getNextValidTimeAfter(new Date()).getTime());
}
}else {
flow.setStartUp(FlowPropertyEnum.NO_START.getCode());
flow.setTriggerNextTime(new CronExpression(flow.getFlowCron()).getNextValidTimeAfter(new Date()).getTime());
}
flowMapper.updateByIdSelective(flow);
List<Node> nodeList = nodeMapper.findVirtualByFlowId(flow.getFlowId());
......
......@@ -47,3 +47,11 @@ file:
myth-job:
filestystem: FASTDFS
snapshoot-date: 60 #快照的保存时间 单位天
myth:
plugin:
env: plugin_test
biz: byit-myth-job
port: 8971
register:
url: http://localhost:8080/myth-register
\ No newline at end of file
......@@ -80,6 +80,13 @@
<groupId>redis.clients</groupId>
<artifactId>jedis</artifactId>
</dependency>
<dependency>
<groupId>myth-job</groupId>
<artifactId>plugin-spring-boot-starter</artifactId>
<version>1.0-SNAPSHOT</version>
</dependency>
</dependencies>
<build>
......
package com.byit.conf;
import cn.hutool.http.HttpRequest;
import com.alibaba.fastjson.JSON;
import com.byit.dto.web.ResponseResult;
import com.byit.packet.response.PluginRpcResponsePacket;
import com.byit.param.ResultCallback;
import io.netty.channel.ChannelHandlerContext;
import static com.alibaba.fastjson.serializer.SerializerFeature.WriteClassName;
/**
* @author huangfu
*/
public class RpcResultHttpCallback implements ResultCallback {
private String gatewayIpAndPort;
public RpcResultHttpCallback(String gatewayIpAndPort) {
this.gatewayIpAndPort = "http://"+gatewayIpAndPort+"/myth-job-admin/api/callback/callbackRes";
}
@Override
public void resultCallback(ChannelHandlerContext ctx, PluginRpcResponsePacket pluginRpcResponsePacket) {
Object result = pluginRpcResponsePacket.getResult();
String responseStr = JSON.toJSONString(pluginRpcResponsePacket, WriteClassName);
HttpRequest post = HttpRequest.post(gatewayIpAndPort);
post.header("token", "Token");
post.body("param="+responseStr).execute();
System.out.println("-------消息回复成功-----------");
}
}
......@@ -10,7 +10,8 @@ public enum ScheduleTypeEnum {
NORMAL(1, "正常跑批"),
REPEAT(2, "重跑"),
REPAIR(3, "补批"),
REAL(4, "插件端立即运行")
REAL(4, "插件端立即运行"),
JAVA_SYNC(5, "java同步调度")
;
private Integer code;
......
package com.byit.mapper;
import com.byit.dto.plugin.JavaTask;
import org.apache.ibatis.annotations.Param;
import org.springframework.stereotype.Repository;
import java.util.List;
@Repository
public interface JavaTaskMapper {
int deleteById(Integer id);
/**
* 当前时间+设定时间要执行的类
* @param triggerTime
* @return
*/
List<JavaTask> findAllByTriggerTimeLessThanEqual(@Param("triggerTime") Long triggerTime);
int insertSelective(JavaTask record);
JavaTask getById(Integer id);
......
package com.byit.service;
import com.byit.dto.plugin.JavaTask;
import java.util.List;
/**
* java任务操作
* @author huangfu
*/
public interface JavaTaskService {
/**
* 查询当前时间+若干秒内该执行的任务节点
* @param triggerTime
* @return
*/
List<JavaTask> findAllByTriggerTimeLessThanEqual(Long triggerTime);
/**
* 修改javaTask
* @param record
*/
void updateById(JavaTask record);
}
package com.byit.service;
import com.byit.packet.request.PluginRpcRequestPacket;
import com.byit.packet.response.PluginRpcResponsePacket;
/**
* @author huangfu
*/
public interface RunJavaService {
/**
* 执行远程java的委托方法
* @param rpcRequestPacket
* @return
*/
PluginRpcResponsePacket runJava(PluginRpcRequestPacket rpcRequestPacket);
}
package com.byit.service.impl;
import com.byit.dto.plugin.JavaTask;
import com.byit.mapper.JavaTaskMapper;
import com.byit.service.JavaTaskService;
import org.springframework.stereotype.Service;
import java.util.List;
/**
* @author huangfu
*/
@Service
public class JavaTaskServiceImpl implements JavaTaskService {
private final JavaTaskMapper javaTaskMapper;
public JavaTaskServiceImpl(JavaTaskMapper javaTaskMapper) {
this.javaTaskMapper = javaTaskMapper;
}
@Override
public List<JavaTask> findAllByTriggerTimeLessThanEqual(Long triggerTime) {
return javaTaskMapper.findAllByTriggerTimeLessThanEqual(triggerTime);
}
@Override
public void updateById(JavaTask record) {
javaTaskMapper.updateByIdSelective(record);
}
}
......@@ -73,6 +73,7 @@ public class JobTaskRunLogServiceImpl implements JobTaskRunLogService {
}
@Override
@Transactional(rollbackFor = Exception.class)
public int updateJobTaskRunLog(JobTaskRunLog jobTaskRunLog) {
return jobTaskRunLogMapper.updateJobTaskRunLog(jobTaskRunLog);
}
......
package com.byit.service.impl;
import com.byit.call.TaskServer;
import com.byit.packet.request.PluginRpcRequestPacket;
import com.byit.packet.response.PluginRpcResponsePacket;
import com.byit.service.RunJavaService;
import com.byit.task.annotations.TaskClient;
import org.springframework.stereotype.Service;
/**
* 直接运行
* @author huangfu
*/
@Service
public class RunJavaServiceImpl implements RunJavaService {
@TaskClient
private TaskServer taskServer;
@Override
public PluginRpcResponsePacket runJava(PluginRpcRequestPacket rpcRequestPacket) {
return taskServer.call(rpcRequestPacket);
}
}
......@@ -88,6 +88,13 @@ public class RunNodeServiceImpl implements RunNodeServer, ApplicationEventPublis
});
jobTaskService.saveJobTasks(jobTasks);
log.info("-------------【开始修改工作流{}的下次运行时间,以及各种状态】---------------",flow);
try {
flow.setTriggerNextTime(new CronExpression(flow.getFlowCron()).getNextValidTimeAfter(new Date()).getTime());
} catch (ParseException e) {
flow.setTriggerNextTime(999999999999999999L);
e.printStackTrace();
}
updateFlow(flow);
applicationEventPublisher.publishEvent(new FlowScanEndEvent(this,flow.getFlowId()));
log.info("-------saveRunRecAndTaskAndUpdate end-----------【运行结束】-----------------");
......
package com.byit.service.mapservice;
import com.byit.dto.plugin.JavaTask;
/**
* @author huangfu
*/
public interface JavaTaskAndLogService {
/**
* 修改javaTask和报错任务实例的组合操作
* @param javaTask
*/
void updateJavaTaskAndSaveLog(JavaTask javaTask);
}
package com.byit.service.mapservice.impl;
import com.byit.dto.plugin.JavaTask;
import com.byit.job.WorkRoulette;
import com.byit.job.utils.CronExpression;
import com.byit.model.RunRecording;
import com.byit.service.JavaTaskService;
import com.byit.service.mapservice.JavaTaskAndLogService;
import com.byit.task.JavaTaskJobTask;
import org.springframework.stereotype.Service;
import java.text.ParseException;
import java.util.Date;
/**
* @author Administrator
*/
@Service
public class JavaTaskAndLogServiceServiceImpl implements JavaTaskAndLogService {
private final JavaTaskService javaTaskService;
public JavaTaskAndLogServiceServiceImpl(JavaTaskService javaTaskService) {
this.javaTaskService = javaTaskService;
}
@Override
public void updateJavaTaskAndSaveLog(JavaTask javaTask) {
try {
javaTask.setTriggerTime(new CronExpression(javaTask.getCron()).getNextValidTimeAfter(new Date()).getTime());
} catch (ParseException e) {
javaTask.setTriggerTime(999999999999999999L);
e.printStackTrace();
}
javaTaskService.updateById(javaTask);
//保存到调度轮
JavaTaskJobTask javaTaskJobTask = new JavaTaskJobTask(javaTask);
WorkRoulette.addJob(javaTaskJobTask,javaTask.getTriggerTime());
}
}
package com.byit.task;
import com.byit.conf.RpcResultHttpCallback;
import com.byit.dto.plugin.JavaTask;
import com.byit.dto.plugin.JobTaskRunLog;
import com.byit.enums.ScheduleTypeEnum;
import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.packet.request.PluginRpcRequestPacket;
import com.byit.packet.response.PluginRpcResponsePacket;
import com.byit.param.defaultparam.DefaultResultCallback;
import com.byit.service.impl.JobTaskRunLogServiceImpl;
import com.byit.service.impl.RunJavaServiceImpl;
import com.byit.util.SpringUtil;
import io.netty.util.Timeout;
import io.netty.util.TimerTask;
import java.util.Date;
/**
* java节点的执行器
* @author huangfu
*/
public class JavaTaskJobTask implements TimerTask {
private final JavaTask javaTask;
public JavaTaskJobTask(JavaTask javaTask) {
this.javaTask = javaTask;
}
@Override
public void run(Timeout timeout) throws Exception {
RunJavaServiceImpl service = SpringUtil.getBean(RunJavaServiceImpl.class);
PluginRpcRequestPacket request = new PluginRpcRequestPacket();
//保存到日志
Integer logId = saveLog(javaTask);
request.setExtension(logId+"");
request.setCallbackUrl("http://127.0.0.1:8081/myth-job-admin/api/callback/callbackRes");
request.setParam(javaTask.getParam());
request.setJobName(javaTask.getTaskName());
PluginRpcResponsePacket pluginRpcResponsePacket = service.runJava(request);
JobTaskRunLogWithBLOBs log = new JobTaskRunLogWithBLOBs();
log.setLogId(logId);
if(pluginRpcResponsePacket.isStatus()){
log.setTriggerCode("1");
}else{
log.setTriggerCode("2");
}
log.setTriggerMsg(pluginRpcResponsePacket.getMsg());
JobTaskRunLogServiceImpl jobTaskRunLogService = SpringUtil.getBean(JobTaskRunLogServiceImpl.class);
jobTaskRunLogService.updateJobTaskRunLogWithBLOBs(log);
}
private Integer saveLog(JavaTask javaTask){
JobTaskRunLogWithBLOBs log = new JobTaskRunLogWithBLOBs();
log.setScheduleType(ScheduleTypeEnum.JAVA_SYNC.getCode());
log.setNodeId(javaTask.getId());
log.setNodeName(javaTask.getJobName());
log.setHandlerName(javaTask.getTaskName());
log.setRunParams(javaTask.getParam());
log.setTriggerTime(new Date());
log.setJobType("JAVA_SYNC");
JobTaskRunLogServiceImpl jobTaskRunLogService = SpringUtil.getBean(JobTaskRunLogServiceImpl.class);
int i = jobTaskRunLogService.saveJobTaskRunLog(log);
return log.getLogId();
}
}
package com.byit.thread;
import com.byit.dto.executor.JobRunResultDto;
import com.byit.dto.web.ReturnResult;
import com.byit.enums.RunRecordingEnum;
import com.byit.model.JobTaskRunLogWithBLOBs;
import com.byit.packet.response.PluginRpcResponsePacket;
import com.byit.service.impl.JobTaskRunLogServiceImpl;
import com.byit.util.SpringUtil;
import java.util.Map;
/**
* @author huangfu
*/
public class JavaTaskCallbackThread implements Runnable {
private PluginRpcResponsePacket pluginRpcResponsePacket;
public JavaTaskCallbackThread(PluginRpcResponsePacket pluginRpcResponsePacket) {
this.pluginRpcResponsePacket = pluginRpcResponsePacket;
}
@Override
public void run() {
String logIdStr = pluginRpcResponsePacket.getExtension();
int logId = Integer.parseInt(logIdStr);
JobTaskRunLogServiceImpl mythJobTaskRunLogService = SpringUtil.getBean(JobTaskRunLogServiceImpl.class);
Map<String,String> result = (Map<String,String>)pluginRpcResponsePacket.getResult();
JobTaskRunLogWithBLOBs jobTaskRunLog = new JobTaskRunLogWithBLOBs();
jobTaskRunLog.setRunMsg(result.get("msg"));
jobTaskRunLog.setRunCode(result.get("code"));
jobTaskRunLog.setLogId(logId);
mythJobTaskRunLogService.updateJobTaskRunLogWithBLOBs(jobTaskRunLog);
}
}
package com.byit.thread.helper;
import com.byit.thread.BaseThreadRunHelper;
import javax.sql.DataSource;
/**
* TODO 邮箱的告警功能待商讨
* @author huangfu
*/
public class JavaTaskEmailThreadRunHelper extends BaseThreadRunHelper {
private static final String LOCK_NAME = "java_task_email_lock";
private final DataSource dataSource;
public JavaTaskEmailThreadRunHelper(DataSource dataSource) {
this.dataSource = dataSource;
}
@Override
public Long start() {
return null;
}
@Override
public DataSource getDataSource() {
return dataSource;
}
@Override
public String getLockName() {
return LOCK_NAME;
}
}
package com.byit.thread.helper;
import cn.hutool.core.collection.CollectionUtil;
import com.byit.dto.plugin.JavaTask;
import com.byit.service.JavaTaskService;
import com.byit.service.mapservice.JavaTaskAndLogService;
import com.byit.thread.BaseThreadRunHelper;
import org.springframework.stereotype.Component;
import javax.sql.DataSource;
import java.util.List;
/**
* java节点扫描
* @author huangfu
*/
@Component
public class JavaTaskThreadRunHelper extends BaseThreadRunHelper {
private static final String LOCK_NAME = "java_task_lock";
private final DataSource dataSource;
private final JavaTaskService javaTaskService;
private final JavaTaskAndLogService javaTaskAndLogService;
private static final long PRE_READ_MS = 7000;
public JavaTaskThreadRunHelper(DataSource dataSource, JavaTaskService javaTaskService, JavaTaskAndLogService javaTaskAndLogService) {
this.dataSource = dataSource;
this.javaTaskService = javaTaskService;
this.javaTaskAndLogService = javaTaskAndLogService;
}
@Override
public Long start() {
long nowTime = System.currentTimeMillis();
List<JavaTask> javaTasks = javaTaskService.findAllByTriggerTimeLessThanEqual(nowTime + PRE_READ_MS);
if(CollectionUtil.isNotEmpty(javaTasks)){
javaTasks.forEach(javaTask -> {
//查看该任务的剩余执行次数
Integer repeatCount = javaTask.getRepeatCount();
if(repeatCount > 0){
javaTask.setRepeatCount(javaTask.getRepeatCount()-1);
}
javaTaskAndLogService.updateJavaTaskAndSaveLog(javaTask);
});
}
return PRE_READ_MS;
}
@Override
public DataSource getDataSource() {
return this.dataSource;
}
@Override
public String getLockName() {
return LOCK_NAME;
}
}
......@@ -28,6 +28,14 @@
from java_task
where id = #{id,jdbcType=INTEGER}
</select>
<select id="findAllByTriggerTimeLessThanEqual" resultMap="BaseResultMap">
select
<include refid="Base_Column_List" />
from java_task
where trigger_time <![CDATA[ <= ]]> #{triggerTime,jdbcType=BIGINT} and remaining_count != 0
</select>
<delete id="deleteById" parameterType="java.lang.Integer">
<!-- generated @mbg.generated date: 2020-04-14 -->
delete from java_task
......@@ -157,7 +165,7 @@
</set>
where id = #{id,jdbcType=INTEGER}
</update>
<update id="updateByJobName"parameterType="com.byit.dto.plugin.JavaTask">
<update id="updateByJobName" parameterType="com.byit.dto.plugin.JavaTask">
<!-- generated @mbg.generated date: 2020-04-14 -->
update java_task
<set>
......
......@@ -29,7 +29,8 @@ public class PluginRpcRequestPacket extends BasePacketModel {
private String param;
private ResultCallback resultCallback;
private String callbackUrl;
@Override
public Command getCommand() {
......
......@@ -15,6 +15,7 @@ public class NettyPluginClient extends PluginClient {
@Override
public void send(String address, PluginRpcRequestPacket pluginRpcRequestPacket) throws Exception {
System.out.println("--------------我发送信息了");
PluginConnectClient.send(pluginRpcRequestPacket,address,pluginConnectClientImpl,pluginClientInitialization);
}
}
......@@ -58,6 +58,7 @@ public class PluginClientInitialization {
return Proxy.newProxyInstance(Thread.currentThread().getContextClassLoader(),
new Class[]{this.iFace},
(proxy, method, args) ->{
PluginRpcRequestPacket pluginRpcRequestPacket = (PluginRpcRequestPacket)args[0];
if(pluginRpcRequestPacket == null){
throw new RuntimeException("服务参数为null:【com.byit.packet.request.PluginRpcRequestPacket】!");
......
package com.byit.callback;
import cn.hutool.http.HttpRequest;
import cn.hutool.http.HttpUtil;
import com.alibaba.fastjson.JSON;
import com.byit.dto.plugin.JavaTask;
......@@ -31,8 +32,11 @@ public class DefaultRemainingOperationsCallBack implements RemainingOperationsCa
if(StringUtils.isNotBlank(publishUrl)) {
JavaTask javaTask = annotationParse(annotation);
String javaTaskJson = JSON.toJSONString(javaTask,WriteClassName);
String post = HttpUtil.post(publishUrl, javaTaskJson);
ResponseResult responseResult = JSON.parseObject(post, ResponseResult.class);
HttpRequest post = HttpRequest.post(publishUrl);
post.header("token", "Token");
String body = post.body("param="+javaTaskJson).execute().body();
System.out.println(body);
ResponseResult responseResult = JSON.parseObject(body, ResponseResult.class);
String msg = responseResult.getMsg();
System.err.println("------"+nodeClass.getName()+":"+msg+"-------------");
}
......
package com.byit.server.netty.handler;
import cn.hutool.http.HttpRequest;
import com.alibaba.fastjson.JSON;
import com.byit.dto.web.ReturnResult;
import com.byit.enums.ResponseTyEnum;
import com.byit.factory.PluginServerFactory;
import com.byit.packet.request.PluginRpcRequestPacket;
import com.byit.packet.response.PluginRpcResponsePacket;
import com.byit.param.PluginBeat;
import com.byit.param.ResultCallback;
import com.byit.task.handler.interfaces.IJobHandler;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
......@@ -65,11 +65,18 @@ public class NettyPluginServerHandler extends SimpleChannelInboundHandler<Plugin
rpcResponsePacket.setType(ResponseTyEnum.RESPONSE.getType());
rpcResponsePacket.setRequestId(msg.getRequestId());
//ctx.channel().writeAndFlush(rpcResponsePacket);
ResultCallback resultCallback = msg.getResultCallback();
if (resultCallback != null) {
System.out.println(resultCallback);
resultCallback.resultCallback(ctx,rpcResponsePacket);
}
// ResultCallback resultCallback = msg.getResultCallback();
// if (resultCallback != null) {
// System.out.println(resultCallback);
// resultCallback.resultCallback(ctx,rpcResponsePacket);
// }
String responseStr = JSON.toJSONString(rpcResponsePacket);
HttpRequest post = HttpRequest.post(msg.getCallbackUrl());
post.header("token", "Token");
post.header("contentType","application/json");
post.body(responseStr).execute();
System.out.println("-------消息回复成功-----------");
});
transferPluginRpcResponse.setRequestId(msg.getRequestId());
......
......@@ -9,10 +9,11 @@ import org.springframework.stereotype.Component;
* @author huangfu
*/
@Component
@TaskHandler(taskName = "sentEmailServer")
@TaskHandler(cron = "0 0/1 * * * ?",taskName = "sentEmailServer",autoPublish = true,publishUrl = "http://127.0.0.1:8081/myth-job-admin/api/node/autoAddJavaTask")
public class SentEmailServer extends BaseJobHandler {
@Override
public ReturnResult<String> execute(String param) throws Exception {
System.out.println("-------------SentEmailServer-被调度执行------------");
return ReturnResult.SUCCESS;
}
}
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