Commit a7b4ef9d by huangfusuper

添加本地java任务的执行插件,添加测试案例!添加本低java任务的执行task

parent 1297cb2f
...@@ -25,6 +25,12 @@ ...@@ -25,6 +25,12 @@
<groupId>org.projectlombok</groupId> <groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId> <artifactId>lombok</artifactId>
</dependency> </dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.62</version>
</dependency>
</dependencies> </dependencies>
</project> </project>
\ No newline at end of file
package com.byit;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
/**
* @program: byit-myth-job->AdminApplication
* @description: 管理启动类
* @author: huangfu
* @date: 2019/11/20 11:44
**/
@SpringBootApplication
public class AdminApplication {
public static void main(String[] args) {
SpringApplication.run(AdminApplication.class,args);
}
}
package com.byit.controller;
import com.byit.job.WorkRoulette;
import com.byit.job.model.JavaBeanJobInfo;
import org.springframework.web.bind.annotation.*;
/**
* @program: byit-myth-job->JobHandel
* @description: 测试添加任务
* @author: huangfu
* @date: 2019/11/20 11:57
**/
@RestController
@RequestMapping("job")
public class JobController {
@PostMapping(value = "addJob")
public String addJob(@RequestBody JavaBeanJobInfo javaBeanJobInfo){
WorkRoulette.addJob(javaBeanJobInfo );
return "SUCCESS";
}
}
...@@ -38,6 +38,12 @@ ...@@ -38,6 +38,12 @@
<groupId>myth-job</groupId> <groupId>myth-job</groupId>
<artifactId>myth-core-common</artifactId> <artifactId>myth-core-common</artifactId>
</dependency> </dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.62</version>
</dependency>
</dependencies> </dependencies>
</project> </project>
\ No newline at end of file
package com.byit.job; package com.byit.job;
import com.byit.job.handler.interfaces.IJobHandler; import com.byit.job.model.JavaBeanJobInfo;
import com.byit.job.instancecache.JavaBeanInstanceCache; import com.byit.task.JavaBeanJobTask;
import com.byit.task.JobTask;
import io.netty.util.HashedWheelTimer; import io.netty.util.HashedWheelTimer;
import java.util.TimerTask;
import java.util.concurrent.ThreadFactory; import java.util.concurrent.ThreadFactory;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
...@@ -31,13 +29,9 @@ public class WorkRoulette { ...@@ -31,13 +29,9 @@ public class WorkRoulette {
},1, TimeUnit.SECONDS,8,true,0); },1, TimeUnit.SECONDS,8,true,0);
public static void addJob(){ public static void addJob(JavaBeanJobInfo javaBeanJobInfo){
IJobHandler testJob = JavaBeanInstanceCache.getJobHandler("testJob"); JavaBeanJobTask javaBeanJobTask = new JavaBeanJobTask(javaBeanJobInfo);
JobTask jobTask = new JobTask(testJob, "测试参数"); hashedWheelTimer.newTimeout(javaBeanJobTask,TimeUnit.SECONDS.toNanos(20),TimeUnit.NANOSECONDS);
hashedWheelTimer.newTimeout(jobTask,10000000000L,TimeUnit.NANOSECONDS);
} }
public static void main(String[] args) {
System.out.println(TimeUnit.SECONDS.toNanos(10));
}
} }
package com.byit.task;
import cn.hutool.http.HttpUtil;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.byit.job.model.JavaBeanJobInfo;
import io.netty.util.Timeout;
import io.netty.util.TimerTask;
/**
* @program: byit-myth-job->JvavBeanJobTask
* @description: 开发者编写的java job
* @author: huangfu
* @date: 2019/11/20 15:08
**/
public class JavaBeanJobTask implements TimerTask {
private JavaBeanJobInfo javaBeanJobInfo;
public JavaBeanJobTask(JavaBeanJobInfo javaBeanJobInfo) {
this.javaBeanJobInfo = javaBeanJobInfo;
}
@Override
public void run(Timeout timeout) throws Exception {
String url = javaBeanJobInfo.getUrl( );
String jobHandelName = javaBeanJobInfo.getJobHandelName( );
String param = javaBeanJobInfo.getParam( );
JSONObject jsonObject = new JSONObject();
jsonObject.put("jobHandelName",jobHandelName);
jsonObject.put("param",param);
String result = HttpUtil.post(url, JSON.toJSONString(jsonObject),10*1000);
}
}
...@@ -12,5 +12,32 @@ ...@@ -12,5 +12,32 @@
<groupId>myth-job</groupId> <groupId>myth-job</groupId>
<artifactId>myth-core-common</artifactId> <artifactId>myth-core-common</artifactId>
<dependencies>
<!-- slf4j -->
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
<version>${slf4j-api.version}</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-log4j12</artifactId>
<version>${slf4j-api.version}</version>
<scope>test</scope>
</dependency>
<!-- https://mvnrepository.com/artifact/org.apache.httpcomponents/httpclient -->
<dependency>
<groupId>org.apache.httpcomponents</groupId>
<artifactId>httpclient</artifactId>
<version>4.5.10</version>
</dependency>
<dependency>
<groupId>cn.hutool</groupId>
<artifactId>hutool-all</artifactId>
<version>4.5.11</version>
</dependency>
</dependencies>
</project> </project>
\ No newline at end of file
package com.byit.job.annotation;
import java.lang.annotation.*;
/**
* 主要应用在本地java任务上
*/
@Target({ElementType.TYPE})
@Retention(RetentionPolicy.RUNTIME)
@Documented
public @interface JobHandler {
String value();
}
package com.byit.utils;
import com.byit.annotations.JobHandler;
import com.byit.job.handler.interfaces.IJobHandler;
import java.io.File;
import java.util.ArrayList;
import java.util.List;
import java.util.stream.Collectors;
/**
* @program: byit-myth-job->ScanRootPackage
* @description: 扫描项目
* @author: huangfu
* @date: 2019/11/19 17:31
**/
public class ScanRootPackage {
public ScanRootPackage() {
addJobCache();
}
/**
* 任务基类
*/
private static final Class<IJobHandler> JOB_ROOT_CLASS;
/**
* 任务注解类
*/
private static final Class<JobHandler> JOB_HANDLER_CLASS;
static {
JOB_ROOT_CLASS = IJobHandler.class;
JOB_HANDLER_CLASS = JobHandler.class;
}
/**
* 获取项目的根路径
* @return
*/
private static String getProjectRootPath(){
return new File(ScanRootPackage.class.getResource("/").getPath()).getPath();
}
/**
* 获取任务类型的CLASS对象
* 根据是否是{@link IJobHandler}判断
* @return
*/
private static List<Class> getClassObjective() {
List<String> classFilePath = getClassFilePath( );
String projectRootPath = getProjectRootPath( );
//获取全限定名
List<Class> collect = classFilePath.stream( )
.map(filePth -> {
try {
//获取类全限定名
String replace = filePth.replace(projectRootPath + File.separator, "")
.replace(File.separator, ".")
.replace(".class", "");
return Class.forName(replace);
} catch (ClassNotFoundException e) {
e.printStackTrace( );
}
return null;
}).collect(Collectors.toList( ));
/**
* 符合条件的类
*/
List<Class> meetTheCriteria = collect.stream( ).filter(objectiveClass ->{
if (JOB_ROOT_CLASS.isAssignableFrom(objectiveClass)) {
if(objectiveClass.isAnnotationPresent(JOB_HANDLER_CLASS)){
return true;
}
}
return false;
}).collect(Collectors.toList( ));
return meetTheCriteria;
}
/**
* 向缓存添加任务信息
*/
private static void addJobCache(){
getClassObjective().forEach(jobClass ->{
JobHandler jobHandler = (JobHandler)jobClass.getAnnotation(JobHandler.class);
JobUtils.jobCache.put(jobHandler.value(),jobClass);
});
}
/**
* 获取项目根路径下所有的类文件
* @return
*/
private static List<String> getClassFilePath(){
ArrayList<String> classFileList = new ArrayList<>(10);
scan(getProjectRootPath(), classFileList);
return classFileList;
}
/**
* 递归扫描
* @param rootPath
* @param list
*/
private static void scan(String rootPath,List<String> list){
File rootFile = new File(rootPath);
File[] files = rootFile.listFiles( );
for (File file : files) {
if(file.isFile() && file.getName().endsWith(".class")){
list.add(file.getPath());
}else if(file.isDirectory()){
scan(file.getAbsolutePath(),list);
}
}
}
}
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<artifactId>byit-myth-executor</artifactId>
<groupId>myth-job</groupId>
<version>1.0-SNAPSHOT</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<groupId>myth-job</groupId>
<artifactId>myth-exector</artifactId>
</project>
\ No newline at end of file
...@@ -28,7 +28,7 @@ ...@@ -28,7 +28,7 @@
<dependency> <dependency>
<groupId>myth-job</groupId> <groupId>myth-job</groupId>
<artifactId>myth-exector</artifactId> <artifactId>myth-exector-plugin</artifactId>
<version>${project.version}</version> <version>${project.version}</version>
</dependency> </dependency>
......
package com.byit.conf;
import com.byit.job.executor.SpringExecutor;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* @program: byit-myth-job->ExcecutorConfig
* @description: TODO
* @author: huangfu
* @date: 2019/11/15 17:22
**/
@Configuration
public class ExcecutorConfig {
@Bean
public SpringExecutor springExecutor(){
return new SpringExecutor();
}
}
package com.byit.conf;
import com.byit.rpc.registry.impl.RegistryServiceRegistry;
import com.byit.rpc.remoting.provider.impl.RpcSpringProviderFactory;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
/**
* @program: byit-myth-job->RpcConfig
* @description: Rpc通讯的配置类
* @author: huangfu
* @date: 2019/11/14 15:18
**/
@Configuration
@Slf4j
public class ExecutorRpcConfig {
@Value("${myth-rpc.remoting.port}")
private int port;
@Value("${myth-rpc.registry.address}")
private String address;
@Value("${myth-rpc.registry.env}")
private String env;
@Value("${myth-rpc.registry.biz}")
private String biz;
@Bean
public RpcSpringProviderFactory rpcSpringProviderFactory(){
RpcSpringProviderFactory providerFactory = new RpcSpringProviderFactory();
providerFactory.setPort(port);
providerFactory.setServiceRegistryClass(RegistryServiceRegistry.class);
providerFactory.setServiceRegistryParam(new HashMap<String,String>(){{
put(RegistryServiceRegistry.REGISTRY_ADDRESS, address);
put(RegistryServiceRegistry.BIZ, biz);
put(RegistryServiceRegistry.ENV, env);
}});
log.info(">>>>>>>>>>>>>>> 执行器服务向注册中心注册成功");
return providerFactory;
}
}
package com.byit.executor;
import com.byit.api.ExecutorServer;
import com.byit.job.handler.interfaces.IJobHandler;
import com.byit.job.model.ReturnResult;
import com.byit.rpc.remoting.provider.annotation.RpcService;
import org.springframework.stereotype.Component;
import org.springframework.stereotype.Service;
/**
* @program: byit-myth-job->JavaBeanExcutor
* @description: 解释执行器
* @author: huangfu
* @date: 2019/11/14 11:16
**/
@RpcService
@Service
public class ExecutorServerImpl implements ExecutorServer {
@Override
public ReturnResult<String> execute(IJobHandler iJobHandler,String param) throws Exception {
return iJobHandler.execute(param);
}
}
package com.byit.service;
import com.byit.job.annotation.JobHandler;
import com.byit.job.handler.BaseJobHandler;
import com.byit.job.handler.interfaces.IJobHandler;
import com.byit.job.model.ReturnResult;
import org.springframework.stereotype.Component;
/**
* @program: byit-myth-job->DemoHandel
* @description: huangfu
* @author: huangfu
* @date: 2019/11/15 17:52
**/
@JobHandler("addJob")
@Component
public class DemoHandel extends BaseJobHandler {
@Override
public ReturnResult<String> execute(String param) throws Exception {
System.out.println( "----------------"+param+"----------------------");
return ReturnResult.SUCCESS;
}
@Override
public void init(String param) {
}
@Override
public void destroy(String param) {
}
}
...@@ -14,7 +14,7 @@ ...@@ -14,7 +14,7 @@
<modules> <modules>
<module>myth-executor-api</module> <module>myth-executor-api</module>
<module>myth-executor-server</module> <module>myth-executor-server</module>
<module>myth-exector</module> <module>myth-exector-plugin</module>
</modules> </modules>
......
...@@ -31,6 +31,18 @@ ...@@ -31,6 +31,18 @@
<artifactId>byit-demo-api</artifactId> <artifactId>byit-demo-api</artifactId>
<version>${project.version}</version> <version>${project.version}</version>
</dependency> </dependency>
<dependency>
<groupId>myth-job</groupId>
<artifactId>myth-exector-plugin</artifactId>
<version>1.0-SNAPSHOT</version>
</dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.62</version>
</dependency>
</dependencies> </dependencies>
<build> <build>
......
package com.byit.job;
import com.byit.annotations.JobHandler;
import com.byit.job.handler.BaseJobHandler;
import com.byit.job.model.ReturnResult;
/**
* @program: byit-myth-job->DemoJob
* @description: TODO
* @author: huangfu
* @date: 2019/11/20 12:40
**/
@JobHandler("addJob")
public class DemoJob extends BaseJobHandler {
@Override
public ReturnResult<String> execute(String s) throws Exception {
System.out.println("------------------------------" );
return ReturnResult.SUCCESS;
}
}
package com.byit.job;
import com.alibaba.fastjson.JSON;
import com.byit.job.model.JavaBeanJobInfo;
import com.byit.utils.JobUtils;
import java.io.IOException;
/**
* @program: byit-myth-job->Mains
* @description: TODO
* @author: huangfu
* @date: 2019/11/20 12:18
**/
public class Mains {
public static void main(String[] args) throws IOException {
String plServerUrl = "http://localhost:9999";
String mythCron = "时间";
String routingStrategy = "路由策略";
String blockingStrategy = "阻塞策略";
String callbackToken = "dsasadsad";
String gatewayToken="asdsadsa";
String requestUrl = "http://localhost:8080/job/addJob";
String param="sadsadsa";
String name="addJob";
JavaBeanJobInfo javaBeanJobInfo = new JavaBeanJobInfo(plServerUrl,requestUrl
,param,routingStrategy,blockingStrategy,callbackToken,gatewayToken);
javaBeanJobInfo.setJobHandelName(name);
javaBeanJobInfo.setMythCron(mythCron);
JobUtils.addJob(javaBeanJobInfo);
JobUtils.jobServerStart();
System.out.println(JobUtils.jobCache );
}
}
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