job相关功能封装为solon插件

This commit is contained in:
aoshiguchen
2023-03-12 01:12:48 +08:00
parent b0c3c311af
commit 8d6a2860d9
38 changed files with 404 additions and 507 deletions
+39
View File
@@ -0,0 +1,39 @@
<?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">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.noear</groupId>
<artifactId>solon-parent</artifactId>
<version>2.2.1</version>
</parent>
<groupId>fun.asgc</groupId>
<artifactId>job-solon-plugin</artifactId>
<packaging>jar</packaging>
<dependencies>
<dependency>
<groupId>org.noear</groupId>
<artifactId>solon</artifactId>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</dependency>
<!--quartz-->
<dependency>
<groupId>org.quartz-scheduler</groupId>
<artifactId>quartz</artifactId>
<version>2.3.1</version>
</dependency>
<!--hutool-->
<dependency>
<groupId>cn.hutool</groupId>
<artifactId>hutool-core</artifactId>
<version>5.8.15</version>
</dependency>
</dependencies>
</project>
@@ -0,0 +1,34 @@
package fun.asgc.solon.extend.job;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.atomic.AtomicInteger;
/**
*
* @author: aoshiguchen
* @date: 2022/9/4
*/
public class CustomThreadFactory implements ThreadFactory {
private final ThreadGroup group;
private final AtomicInteger threadNumber = new AtomicInteger(1);
private final String namePrefix;
public CustomThreadFactory(String prefix) {
SecurityManager s = System.getSecurityManager();
group = (s != null) ? s.getThreadGroup() :
Thread.currentThread().getThreadGroup();
namePrefix = prefix + "-thread-";
}
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(group, r, namePrefix + threadNumber.getAndIncrement(), 0);
if (t.isDaemon()) {
t.setDaemon(false);
}
if (t.getPriority() != Thread.NORM_PRIORITY) {
t.setPriority(Thread.NORM_PRIORITY);
}
return t;
}
}
@@ -0,0 +1,17 @@
package fun.asgc.solon.extend.job;
/**
*
* @author: aoshiguchen
* @date: 2022/9/4
*/
public interface IJobCallback {
/**
* 执行日志
* @param jobInfo
* @param param
* @param throwable
*/
void executeLog(JobInfo jobInfo, String param, Throwable throwable);
}
@@ -0,0 +1,34 @@
package fun.asgc.solon.extend.job;
/**
*
* @author: aoshiguchen
* @date: 2022/9/4
*/
public interface IJobExecutor {
/**
* 初始化
* @throws JobException
*/
void init() throws Exception;
/**
* 新增job
* @param jobInfo
*/
void add(JobInfo jobInfo);
/**
* 删除job
* @param jobName
*/
void remove(String jobName);
/**
* 触发
* @param jobName
* @param param
*/
void trigger(String jobName, String param);
}
@@ -0,0 +1,16 @@
package fun.asgc.solon.extend.job;
/**
* @author: aoshiguchen
* @date: 2022/9/4
*/
public interface IJobHandler {
/**
* job执行
* @param param
* @throws Exception
*/
void execute(String param) throws Exception;
}
@@ -0,0 +1,17 @@
package fun.asgc.solon.extend.job;
import java.util.List;
/**
*
* @author: aoshiguchen
* @date: 2022/9/4
*/
public interface IJobSource {
/**
* 获取所有job列表
* @return
*/
List<JobInfo> sourceList();
}
@@ -0,0 +1,26 @@
package fun.asgc.solon.extend.job;
import fun.asgc.solon.extend.job.impl.JobExecutor;
import lombok.extern.slf4j.Slf4j;
import org.noear.solon.Solon;
import org.quartz.Job;
import org.quartz.JobExecutionContext;
import org.quartz.JobExecutionException;
/**
*
* @author: aoshiguchen
* @date: 2022/9/4
*/
@Slf4j
public class JobBean implements Job {
@Override
public void execute(JobExecutionContext jobExecutionContext) throws JobExecutionException {
JobExecutor jobExecutor = Solon.context().getBean(JobExecutor.class);
if (null != jobExecutor) {
jobExecutor.execute(jobExecutionContext);
}
}
}
@@ -0,0 +1,23 @@
package fun.asgc.solon.extend.job;
import lombok.Data;
import lombok.experimental.Accessors;
import java.util.Map;
/**
*
* @author: aoshiguchen
* @date: 2022/9/4
*/
@Accessors(chain = true)
@Data
public class JobInfo {
private String id;
private String name;
private String desc;
private String cron;
private String param;
private boolean enable;
private Map<String, Object> extension;
}
@@ -0,0 +1,45 @@
package fun.asgc.solon.extend.job;
import fun.asgc.solon.extend.job.annotation.EnableJob;
import fun.asgc.solon.extend.job.impl.DefaultJobCallback;
import fun.asgc.solon.extend.job.impl.DefaultJobSource;
import fun.asgc.solon.extend.job.impl.JobExecutor;
import org.noear.solon.Solon;
import org.noear.solon.core.AopContext;
import org.noear.solon.core.Plugin;
import org.noear.solon.core.event.AppLoadEndEvent;
/**
* @author: aoshiguchen
* @date: 2023/3/11
*/
public class XPluginImp implements Plugin {
@Override
public void start(AopContext context) throws Throwable {
EnableJob enableJob = Solon.app().source().getAnnotation(EnableJob.class);
if (null == enableJob || !enableJob.value()) {
return;
}
//应用加载完后,再启动任务
Solon.app().onEvent(AppLoadEndEvent.class, e -> {
IJobSource jobSource = context.getBean(IJobSource.class);
IJobCallback jobCallback = context.getBean(IJobCallback.class);
if (null == jobSource) {
jobSource = new DefaultJobSource();
}
if (null == jobCallback) {
jobCallback = new DefaultJobCallback();
}
JobExecutor jobExecutor = new JobExecutor();
jobExecutor.setJobSource(jobSource);
jobExecutor.setJobCallback(jobCallback);
jobExecutor.start();
context.wrapAndPut(JobExecutor.class, jobExecutor);
});
}
}
@@ -0,0 +1,16 @@
package fun.asgc.solon.extend.job.annotation;
import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
/**
* @author: aoshiguchen
* @date: 2023/3/12
*/
@Target({ElementType.TYPE})
@Retention(RetentionPolicy.RUNTIME)
public @interface EnableJob {
boolean value() default true;
}
@@ -0,0 +1,20 @@
package fun.asgc.solon.extend.job.annotation;
import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
/**
*
* @author: aoshiguchen
* @date: 2022/9/4
*/
@Target({ElementType.TYPE})
@Retention(RetentionPolicy.RUNTIME)
public @interface JobHandler {
String name();
String desc() default "";
String cron();
String param() default "";
}
@@ -0,0 +1,23 @@
package fun.asgc.solon.extend.job.impl;
import fun.asgc.solon.extend.job.IJobCallback;
import fun.asgc.solon.extend.job.JobInfo;
import lombok.extern.slf4j.Slf4j;
/**
* @author: aoshiguchen
* @date: 2023/3/12
*/
@Slf4j
public class DefaultJobCallback implements IJobCallback {
@Override
public void executeLog(JobInfo jobInfo, String param, Throwable throwable) {
if (null == throwable) {
log.debug("[Solon Plugin Job] Job执行 id:{} name:{} desc:{} param:{}", jobInfo.getId(), jobInfo.getName(), jobInfo.getDesc(), param);
} else {
log.error("[Solon Plugin Job] Job执行 id:{} name:{} desc:{} param:{}", jobInfo.getId(), jobInfo.getName(), jobInfo.getDesc(), param, throwable);
}
}
}
@@ -0,0 +1,44 @@
package fun.asgc.solon.extend.job.impl;
import cn.hutool.core.collection.CollectionUtil;
import cn.hutool.core.util.StrUtil;
import fun.asgc.solon.extend.job.IJobHandler;
import fun.asgc.solon.extend.job.IJobSource;
import fun.asgc.solon.extend.job.annotation.JobHandler;
import fun.asgc.solon.extend.job.JobInfo;
import org.noear.solon.Solon;
import java.util.List;
/**
*
* @author: aoshiguchen
* @date: 2022/9/4
*/
public class DefaultJobSource implements IJobSource {
@Override
public List<JobInfo> sourceList() {
List<IJobHandler> jobHandlerList = Solon.context().getBeansOfType(IJobHandler.class);
if (CollectionUtil.isEmpty(jobHandlerList)) {
return CollectionUtil.newArrayList();
}
List<JobInfo> jobInfoList = CollectionUtil.newArrayList();
for (IJobHandler jobHandler : jobHandlerList) {
JobHandler handler = jobHandler.getClass().getAnnotation(JobHandler.class);
if (null == handler || StrUtil.isEmpty(handler.name()) || StrUtil.isEmpty(handler.cron())) {
continue;
}
jobInfoList.add(new JobInfo()
.setId(handler.name())
.setName(handler.name())
.setDesc(handler.desc())
.setCron(handler.cron())
.setParam(handler.param())
.setEnable(true)
);
}
return jobInfoList;
}
}
@@ -0,0 +1,175 @@
package fun.asgc.solon.extend.job.impl;
import cn.hutool.core.collection.CollectionUtil;
import cn.hutool.core.util.StrUtil;
import fun.asgc.solon.extend.job.*;
import fun.asgc.solon.extend.job.annotation.JobHandler;
import org.noear.solon.Solon;
import org.quartz.*;
import org.quartz.impl.StdSchedulerFactory;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
/**
* Job执行器
* @author: aoshiguchen
* @date: 2022/9/4
*/
public class JobExecutor implements IJobExecutor {
private IJobSource jobSource;
private ThreadPoolExecutor threadPoolExecutor;
private Map<String, JobInfo> jobInfoMap = new ConcurrentHashMap<>();
private SchedulerFactory schedulerFactory;
private Scheduler scheduler;
private Map<String, IJobHandler> jobHandlerMap = new ConcurrentHashMap<>();
private Set<String> runJobSet = CollectionUtil.newHashSet();
private IJobCallback jobCallback;
private Map<String, TriggerKey> triggerKeyMap = new ConcurrentHashMap<>();
public void start() {
if (null == threadPoolExecutor) {
threadPoolExecutor = new ThreadPoolExecutor(5, 20, 10L, TimeUnit.SECONDS,
new LinkedBlockingQueue<>(), new CustomThreadFactory("SolonJob"));
}
List<IJobHandler> jobHandlerList = Solon.context().getBeansOfType(IJobHandler.class);
if (!CollectionUtil.isEmpty(jobHandlerList)) {
for (IJobHandler item : jobHandlerList) {
JobHandler jobHandler = item.getClass().getAnnotation(JobHandler.class);
if (null == jobHandler) {
continue;
}
jobHandlerMap.put(jobHandler.name(), item);
}
}
try {
this.schedulerFactory = new StdSchedulerFactory();
this.scheduler = schedulerFactory.getScheduler();
this.init();
} catch (Exception e){
throw new RuntimeException("job初始化异常");
}
}
public void setJobSource(IJobSource jobSource) {
this.jobSource = jobSource;
}
public void setThreadPoolExecutor(ThreadPoolExecutor threadPoolExecutor) {
this.threadPoolExecutor = threadPoolExecutor;
}
public void setJobCallback(IJobCallback jobCallback) {
this.jobCallback = jobCallback;
}
@Override
public void init() throws Exception {
List<JobInfo> jobInfoList = jobSource.sourceList();
if (CollectionUtil.isEmpty(jobInfoList)) {
return;
}
for (JobInfo jobInfo : jobInfoList) {
add(jobInfo);
}
}
@Override
public void add(JobInfo jobInfo) {
if (null == jobInfo || StrUtil.isEmpty(jobInfo.getId()) || StrUtil.isEmpty(jobInfo.getName()) ||
StrUtil.isEmpty(jobInfo.getCron())) {
return;
}
synchronized (jobInfo.getId()) {
runJobSet.add(jobInfo.getName());
jobInfoMap.put(jobInfo.getId(), jobInfo);
TriggerKey triggerKey = TriggerKey.triggerKey(jobInfo.getId());
triggerKeyMap.put(jobInfo.getId(), triggerKey);
JobKey jobKey = new JobKey(jobInfo.getName());
CronScheduleBuilder cronScheduleBuilder = CronScheduleBuilder.cronSchedule(jobInfo.getCron()).withMisfireHandlingInstructionDoNothing();
CronTrigger cronTrigger = TriggerBuilder.newTrigger().withIdentity(triggerKey).withSchedule(cronScheduleBuilder).build();
JobDetail jobDetail = JobBuilder.newJob(JobBean.class).withIdentity(jobKey).build();
try {
if (jobInfo.isEnable()) {
scheduler.scheduleJob(jobDetail, cronTrigger);
scheduler.start();
}
} catch (Exception e) {
throw new RuntimeException(String.format("新增job[name=%s]异常", jobInfo.getName()));
}
}
}
@Override
public void remove(String jobId) {
JobInfo jobInfo = jobInfoMap.get(jobId);
if (null == jobInfo) {
return;
}
synchronized (jobId) {
runJobSet.remove(jobInfo.getName());
unscheduleJob(jobId);
}
}
@Override
public void trigger(String jobId, String param) {
doExecute(jobId, param);
}
public void execute(JobExecutionContext context) throws JobExecutionException {
if (null == context || null == context.getTrigger()) {
return;
}
String jobId = context.getTrigger().getKey().getName();
JobInfo jobInfo = jobInfoMap.get(jobId);
if (null == jobInfo) {
unscheduleJob(jobId);
return;
}
doExecute(jobId, jobInfo.getParam());
}
private void unscheduleJob(String jobId) {
TriggerKey triggerKey = triggerKeyMap.get(jobId);
if (null != triggerKey) {
try {
scheduler.unscheduleJob(triggerKey);
} catch (Exception e) {
e.printStackTrace();
}
}
}
private void doExecute(String jobId, String param) {
JobInfo jobInfo = jobInfoMap.get(jobId);
if (null == jobInfo) {
return;
}
IJobHandler jobHandler =jobHandlerMap.get(jobInfo.getName());
if (null == jobHandler) {
return;
}
threadPoolExecutor.submit(() -> {
try {
jobHandler.execute(param);
if (null != jobCallback) {
jobCallback.executeLog(jobInfo, param, null);
}
} catch (Throwable e) {
jobCallback.executeLog(jobInfo, param, e);
}
});
}
}
@@ -0,0 +1 @@
package fun.asgc.solon.extend.job;
@@ -0,0 +1,2 @@
solon.plugin=fun.asgc.solon.extend.job.XPluginImp
solon.plugin.priority=2
@@ -0,0 +1 @@
package fun.asgc.solon.extend.orika;