最近的一个小项目是做一个简单的数据仓库,需要将其他数据库的数据抽取出来,并通过而出抽取成页面需要的数据,以空间换时间的方式,让后端报表查询更快。 因为在抽取的过程中,有一定的先后顺序,需要做一个任务调度器,某一优先级的会先执行,然后会进入下一个优先级的队列任务中。 先定义了一个Map的集合,key是优先级,value是任务的集合,某一个优先级内的任务是并发执行的,而不同优先级是串行执行的,前一个优先级执行完之后,后面的才会执行。
ConcurrentHashMap<Integer/* 优先级. */, List<BaseTask>/* 任务集合. */> tasks = new ConcurrentHashMap<>();
这个调度管理有一个演进的过程,我先说第一个,这个是比较好理解的。 第一个版本: 首先对tasks集合中的key进行一个排序,我定义的是数字越小就有限执行,则进行遍历key值,并取出某个优先级的任务队列,执行任务队列的任务。任务的执行交给线程池去执行,在遍历内部,需要不断的检查这个队列中的任务是否都执行了,没有则一直等待否则进入到下个队列,任务执行的时候可能会抛出异常,但是不管任务是否异常,都将任务状态设置已执行。 下面是其核心代码:
1public void run() { 2 //对key值进行排序 3 Enumeration<Integer> keys = tasks.keys(); 4 List<Integer> prioritys = new ArrayList<>(); 5 while (keys.hasMoreElements()) { 6 prioritys.add(keys.nextElement()); 7 } 8 Collections.sort(prioritys);//升序 9 //对key进行遍历,执行某个某个优先级的任务队列 10 for (Integer priority : prioritys) { 11 List<BaseTask> taskList = tasks.get(priority); 12 if (taskList.isEmpty()) { 13 continue; 14 } 15 logger.info("execute priority {} task ", taskList.get(0).priority); 16 for (BaseTask task : taskList) { 17 executor.execute(() -> { 18 try { 19 task.doTask(); 20 } catch (Exception e) { 21 e.printStackTrace(); 22 } 23 });//线程中执行任务 24 } 25 while (true) {//等待所有线程都执行完成之后执行下一个任务队列 26 boolean finish = true; 27 for (BaseTask t : taskList) { 28 if (!t.finish) { 29 finish = false; 30 } 31 } 32 if (finish) {//当前任务都执行完毕 33 break; 34 } 35 Misc.sleep(1000);//Thread.sleep(1000) 36 } 37 Misc.sleep(1000); 38 } 39 }
关键代码很好理解,在任务执行之前,需要对所有任务都初始化,初始化的时候给出每个任务的优先级和任务名称,任务抽象类如下:
1public abstract class BaseTask { 2 public String taskName;//任务名称 3 public Integer priority; //优先级 4 public boolean finish; //任务完成? 5 /** 6 * 执行的任务 7 */ 8 public abstract void doTask(Date date) throws Exception;
第一个版本的思路很简单。 第二个版本稍微有一点点复杂。这里主要介绍该版本的内容,后续将代码的链接附上。 程序是由SpringBoot搭建起来的,定时器是Spring内置的轻量级的Quartz,使用Aop方式拦截异常,使用注解的方式在任务初始化时设置任务的初始变量。使用EventBus解耦程序,其中程序简单实现邮件发送功能(该功能还需要自己配置参数),以上这些至少需要简单的了解一下。 程序的思路:在整个队列执行过程中会有多个管道,某个队列上的管道任务执行完成,可以直接进行到下一个队列中执行,也设置了等待某一个队列上的所有任务都执行完成才执行当前任务。在某个队列任务中会标识某些任务是一队的,其他的为另一队,当这一队任务执行完成,就可以到下一个队列中去,不需要等待另一队。 这里会先初始化每个队列的每个队的条件,这个条件就是每个队的任务数,执行完成减1,当为0时,就进入下一个队列中。 分四个步骤进行完成: 1.bean的初始化 2.条件的设置 3.任务的执行 4.任务异常和任务执行完成之后通知检查是否执行下一个队列的任务
1.bean的初始化
1.创建注解类
1@Retention(RetentionPolicy.RUNTIME) 2@Target(ElementType.TYPE) 3@Documented 4public @interface TaskAnnotation { 5 6 int priority() default 0;//优先级 7 String taskName() default "";//任务名称 8 TaskQueueEnum[] queueName() default {};//队列名称 9}
2.实现BeanPostProcessor,该接口是中有两个方法postProcessBeforeInitialization和postProcessAfterInitialization,分别是bean初始化之前和bean初始化之后做的事情。
1 Annotation[] annotations = bean.getClass().getAnnotations();//获取类上的注解 2 if (ArrayUtils.isEmpty(annotations)) {//注解为空时直接返回(不能返回空,否则bean不会被加载) 3 return bean; 4 } 5 for (Annotation annotation : annotations) { 6 if (annotation.annotationType().equals(TaskAnnotation.class)) { 7 TaskAnnotation taskAnnotation = (TaskAnnotation) annotation;//强转 8 try { 9 Field[] fields = target.getClass().getFields();//需要通过反射将值进行修改,下面的操作仅仅是对象的引用 10 if (!ArrayUtils.isEmpty(fields)) { 11 for (Field f : fields) { 12 f.setAccessible(true); 13 if (f.getName().equals("priority")) { 14 f.set(target, taskAnnotation.priority()); 15 } 16 } 17 } 18 } 19 }
上面需要注意的一点是需要通过反射的机制给bean设置值,不能直接调用bean的方式set值,否则bean的值是空的。 上面的代码通过实现BeanPostProcessor后置处理器,处理任务上的注解,完成对任务的初始化的。
2.条件的初始化
创建条件类,提供初始化的方法。
1public abstract class BaseTask { 2 3 public int nextPriority;//子级节点的优先级 4 5 public String taskName;//任务名称 6 7 public Integer priority; //优先级 8 9 public String queueName;//队列名称 10 11 public boolean finish; //任务完成? 12 13 public boolean allExecute; 14 15 /** 16 * 执行的任务 17 */ 18 public abstract void doTask(Date date) throws Exception; 19 20 //任务完成之后,通过eventBus发送通知,是否需要执行下一个队列 21 public void notifyExecuteTaskMsg(EventBus eventBus, Date date) { 22 EventNotifyExecuteTaskMsg msg = new EventNotifyExecuteTaskMsg(); 23 msg.setDate(date); 24 msg.setNextPriority(nextPriority); 25 msg.setQueueName(queueName); 26 msg.setPriority(priority); 27 msg.setTaskName(taskName); 28 eventBus.post(msg); 29 } 30} 31 32public class TaskExecuteCondition { 33 34 private ConcurrentHashMap<String, AtomicInteger> executeMap = new ConcurrentHashMap<>(); 35 36 /** 37 * 初始化,每个队列进行分组,每个组的任务数量放入map集合中. 38 */ 39 public void init(ConcurrentHashMap<Integer, List<BaseTask>> tasks) { 40 Enumeration<Integer> keys = tasks.keys(); 41 List<Integer> prioritys = new ArrayList<>(); 42 while (keys.hasMoreElements()) { 43 prioritys.add(keys.nextElement()); 44 } 45 Collections.sort(prioritys);//升序 46 for (Integer priority : prioritys) { 47 List<BaseTask> list = tasks.get(priority); 48 if (list.isEmpty()) { 49 continue; 50 } 51 //对每个队列进行分组 52 Map<String, List<BaseTask>> collect = list.stream() 53 .collect(Collectors.groupingBy(x -> x.queueName, Collectors.toList())); 54 for (Entry<String, List<BaseTask>> entry : collect.entrySet()) { 55 for (BaseTask task : entry.getValue()) { 56 addCondition(task.priority, task.queueName); 57 } 58 } 59 } 60 } 61 62 /** 63 * 执行任务完成,条件减1 64 */ 65 public boolean executeTask(Integer priority, String queueName) { 66 String name = this.getQueue(priority, queueName); 67 AtomicInteger count = executeMap.get(name); 68 int sum = count.decrementAndGet(); 69 if (sum == 0) { 70 return true; 71 } 72 return false; 73 } 74 75 /** 76 * 对个某个队列的条件 77 */ 78 public int getCondition(Integer priority, String queueName) { 79 String name = this.getQueue(priority, queueName); 80 return executeMap.get(name).get(); 81 } 82 83 private void addCondition(Integer priority, String queueName) { 84 String name = this.getQueue(priority, queueName); 85 AtomicInteger count = executeMap.get(name); 86 if (count == null) { 87 count = new AtomicInteger(0); 88 executeMap.put(name, count); 89 } 90 count.incrementAndGet(); 91 } 92 93 private void addCondition(Integer priority, String queueName, int sum) { 94 String name = this.getQueue(priority, queueName); 95 AtomicInteger count = executeMap.get(name); 96 if (count == null) { 97 count = new AtomicInteger(sum); 98 executeMap.put(name, count); 99 } else { 100 count.set(sum); 101 } 102 } 103 104 105 private String getQueue(Integer priority, String queueName) { 106 return priority + queueName; 107 } 108 109 /** 110 * 清除队列 111 */ 112 public void clear() { 113 this.executeMap.clear(); 114 } 115}
3.任务的执行
任务执行类提供run方法,执行第一个队列,并提供获取下一个队列优先级方法,执行某个队列某个组的方法。
1public class ScheduleTask { 2 private static final Logger logger = LoggerFactory.getLogger(ScheduleTask.class); 3 4 public ConcurrentHashMap<Integer/* 优先级. */, List<BaseTask>/* 任务集合. */> tasks = new ConcurrentHashMap<>(); 5 6 @Autowired 7 private ThreadPoolTaskExecutor executor;//线程池 8 //任务会先执行第一队列的任务. 9 public void run(Date date) { 10 Enumeration<Integer> keys = tasks.keys(); 11 List<Integer> prioritys = new ArrayList<>(); 12 while (keys.hasMoreElements()) { 13 prioritys.add(keys.nextElement()); 14 } 15 Collections.sort(prioritys);//升序 16 Integer priority = prioritys.get(0); 17 executeTask(priority, date);//执行第一行的任务. 18 } 19 //获取下一个队列的优先级 20 public Integer nextPriority(Integer priority) { 21 Enumeration<Integer> keys = tasks.keys(); 22 List<Integer> prioritys = new ArrayList<>(); 23 while (keys.hasMoreElements()) { 24 prioritys.add(keys.nextElement()); 25 } 26 Collections.sort(prioritys);//升序 27 for (Integer pri : prioritys) { 28 if (priority < pri) { 29 return pri; 30 } 31 } 32 return null;//没有下一个队列 33 } 34 35 public void executeTask(Integer priority) { 36 List<BaseTask> list = tasks.get(priority); 37 if (list.isEmpty()) { 38 return; 39 } 40 for (BaseTask task : list) { 41 execute(task); 42 } 43 } 44 //执行某个队列的某个组 45 public void executeTask(Integer priority, String queueName) { 46 List<BaseTask> list = this.tasks.get(priority); 47 list = list.stream().filter(task -> queueName.equals(task.queueName)) 48 .collect(Collectors.toList()); 49 if (list.isEmpty()) { 50 return; 51 } 52 for (BaseTask task : list) { 53 execute(task); 54 } 55 } 56 57 public void execute(BaseTask task) { 58 executor.execute(() -> { 59 try { 60 task.doTask(date);// 61 } catch (Exception e) {//异常处理已经Aop拦截处理 62 } 63 });//线程中执行任务 64 } 65 66 /** 67 * 增加任务 68 */ 69 public void addTask(BaseTask task) { 70 List<BaseTask> baseTasks = tasks.get(task.priority); 71 if (baseTasks == null) { 72 baseTasks = new ArrayList<>(); 73 List<BaseTask> putIfAbsent = tasks.putIfAbsent(task.priority, baseTasks); 74 if (putIfAbsent != null) { 75 baseTasks = putIfAbsent; 76 } 77 } 78 baseTasks.add(task); 79 } 80 81 /** 82 * 将任务结束标识重新设置 83 */ 84 public void finishTask() { 85 tasks.forEach((key, value) -> { 86 for (BaseTask task : value) { 87 task.finish = false; 88 } 89 }); 90 } 91}
4.任务异常和任务执行完成之后通知检查是否执行下一个队列的任务
1public class EventNotifyExecuteTaskListener { 2 private static final Logger logger = LoggerFactory .getLogger(EventNotifyExecuteTaskListener.class); 3 @Autowired 4 private ScheduleTask scheduleTask; 5 6 @Autowired 7 private TaskExecuteCondition condition; 8 9 @Subscribe 10 public void executeTask(EventNotifyExecuteTaskMsg msg) { 11 //当前队列的某组内容是否都执行完成 12 boolean success = condition.executeTask(msg.getPriority(), msg.getQueueName()); 13 if (success) { 14 Integer nextPriority = scheduleTask.nextPriority(msg.getPriority()); 15 if (nextPriority != null) { 16 scheduleTask.executeTask(nextPriority, msg.getQueueName(), msg.getDate());//执行下一个队列 17 } else {//执行完成,重置任务标识 18 scheduleTask.finishTask(); 19 logger.info("CoreTask end!"); 20 } 21 } 22 } 23}
整个思路介绍到这里,那么接下来是整个项目中出现的一些问题 1.BeanPostProcessor与Aop一起使用时,postProcessAfterInitialization调用之后获取的bean分为不同的了,一个是jdk原生实体对象,一种是Aop注解下的类会被cglib代理,生成带有后缀的对象,如果通过这个对象时反射获取类的注解,字段和方法,就获取不到,在代码中,需要将其转化一下,将cgLib代理之后的类转化为不带后缀的对象。 2.postProcessAfterInitialization的参数bean不能直接设置值,就是如下:
1 TaskAnnotation taskAnnotation = (TaskAnnotation) annotation;//强转 2 BaseTask baseTask = (BaseTask) bean;//强转 3 baseTask.priority = taskAnnotation.priority();
在使用对象时,其中对象的字段时为空的,并需要通过反射的方式去设置字段的值。 上面仅仅只是个人的想法,如果有更好的方式,或者有某些地方可以进行改进的,我们可以共同探讨一下。
链接地址:https://github.com/wangice/task-scheduler 程序中使用了一个公共包:https://github.com/wangice/misc