- 文档
- 教程
- 后端
【免费下载链接】CodeGuide
:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总,旨在为大家提供一个清晰详细的学习教程,侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助,请给予支持(关注、点赞、分享)!
本文以《SpringBoot 中间件设计和开发》专栏第 15 章《分布式任务调度》为核心,完整继承其"前言 + 需求背景"的论述骨架,并依托本仓库同主题完整源码文档,深入讲解 DcsSchedule 分布式任务中间件的设计思路、使用方式与核心实现。读者学完本篇文章后,将理解"为什么单机 Schedule 撑不住业务体量、分布式任务系统解决什么问题",并掌握一套以 Zookeeper 为配置中心、可统一启停、支持宕机灾备的 SpringBoot 分布式任务调度中间件的完整落地方案。
一、前言:CRUD 程序员会不会越来越便宜?
CRUD程序员会不会越来越便宜?——这是每一个身处业务开发一线的 Java 工程师都该认真思考的问题。
CRUD,是程序员的自嘲,指自己经常开发增删改查或者接口包装的简单逻辑代码。但恰恰是这部分"简单逻辑"的代码,几乎占据了现阶段互联网公司里最消耗研发人员的部分:任务的业务需求实现中,大量存在重复的、简单的、单一的功能和逻辑开发,而这些无论是业务功能还是技术组件都没有被单独抽离出来。于是每次开发需求都要重新折腾一遍,最终导致研发、测试到交付一整条线的人员投入,重复造轮子、重复做验证。
对个人来说,开发 CRUD 几乎没有技术成长。开发 CRUD 只是程序员成长过程中的一个阶段,随着个人能力的提升以及跳槽,必然会走向更核心的开发。站在公司技术部门的层面,也都希望投入更少的人实现更高的交付能力,所以组件化、物料化以及低代码编排会越来越抢占 CRUD 的市场。
这正是本仓库(CodeGuide)《SpringBoot 中间件设计和开发》小册(参见 第 2 章 小册学习介绍&源码授权 与 小册上线说明)选择"分布式任务调度"作为其中一章的原因——它既是业务系统里最高频的基础设施之一,也是锻炼"注册中心、任务调度、控制台"三方联动设计能力的最佳实践场景。
二、需求背景:单机 Schedule 为什么撑不住业务体量
在互联网开发的业务场景中,常常有一块功能或者一个独立的服务,专门用于处理定时任务。例如:
- 扫描库表待结算日息;
- 扫描待开始活动状态;
- 扫描用户会员过期时间;
- 处理一些异常流程的补偿动作。
等等诸如此类的功能。
一般最开始的时候,一台单机的任务计算能力就可以支撑起业务体量。但随着业务规模的逐步增加,系统的承载量也随之加大,此时一个单机的任务系统就很难再支撑起整个业务体量的任务扫描工作了。最典型的例子就是:每天 0 点到 3 点需要扫描贷款日息,由于单机任务处理能力有限,会发现已经到了第二天的 0 点,第一天的数据还没有处理完。
所以,这个时候我们需要一个分布式的任务系统:可以把任务作业分散到各个服务处理实例节点上去,加强整个服务的运算承载能力。
而 SpringBoot 自带的@Scheduled定时任务,简单易用(见下方代码),在开发中如果需要做一些定时或指定时刻循环执行的逻辑时,基本都会使用它:
@SpringBootApplication @EnableScheduling public class Application { public static void main(String[] args) { SpringApplication.run(Application.class, args); } @Scheduled(cron = "0/3 * * * * *") public void demoTask() { //... } }但是,如果任务是比较大型的,比如定时跑批 T+1 结算、商品秒杀前状态变更、刷新数据预热到缓存等等,这些定时任务都有相同的特点:作业量大、实时性强、可用率高。而这时候如果只是单纯使用Schedule就显得不足以控制。
于是,分布式 DcsSchedule 任务的产品需求就出来了(本文完整实现内容可对照仓库文档 《开发基于SpringBoot的分布式任务中间件DcsSchedule》):
- 多机器部署任务:把任务分散到多个服务实例上执行;
- 统一控制中心启停:通过控制台统一管理所有任务的启动与关闭;
- 宕机灾备,自动启动执行:节点宕机可被感知,保证任务不中断;
- 实时检测任务执行信息:部署数量、任务总量、成功次数、失败次数、执行耗时等。
下面这张图就是 DcsSchedule 控制台的首页监控界面,直观展示了"部署总数、服务总数、实例统计、任务总数"等实时运行数据:
而任务列表页则可以按服务筛选查看每一个任务的 IP、服务ID、对象名称、方法名称、任务描述、cron 计划与运行状态,并提供「启动」「关闭」两个操作按钮:
三、整体设计:注册中心 + 任务 + 控制台三方联动
在动手实现前,先明确这个中间件的技术选型与架构思路。开发一款基于 SpringBoot 的分布式任务中间件,需要具备以下知识工具:
- 读取 Yml 自定义配置:中间件的连接地址、服务标识等都需要从配置文件读取;
- 使用 Zookeeper 作为配置中心:这样如果有机器宕机了,就可以通过临时节点监听感知到;
- 利用 Spring 的
ApplicationContextAware、BeanPostProcessor、ApplicationListener:完成服务启动、注解扫描、节点挂载; - 分布式任务统一控制台:通过 Zookeeper 的接口功能做数据展示和启停操作。
整体上,DcsSchedule 由三部分组成:schedule-spring-boot-starter(接入业务工程的任务中间件)、Zookeeper(注册中心/配置中心)、itstack-middleware-control(统一控制台)。控制台本身并不复杂,只是使用中间件提供的 ZK 功能接口做展示和操作。
四、中间件使用:三步接入 SpringBoot 工程
1. 环境准备
JDK 1.8;
SpringBoot 2.x;
配置中心 Zookeeper(完整版文档中调试环境使用 3.4.14,本仓库小册 第 2 章 给出的开发环境为 Zookeeper 3.6.0)。准备好 Zookeeper 服务后:
- 下载解压,在 bin 同级路径创建
data、logs文件夹; - 修改
conf/zoo.cfg,配置数据与日志目录,例如:
dataDir=D:\Program Files\apache-zookeeper-3.4.14\data dataLogDir=D:\Program Files\apache-zookeeper-3.4.14\logs- 下载解压,在 bin 同级路径创建
打包部署控制平台
itstack-middleware-control,部署后访问http://localhost:7397即可打开控制台。
2. 配置 POM 依赖
<dependency> <groupId>org.itstack.middleware</groupId> <artifactId>schedule-spring-boot-starter</artifactId> <version>1.0.0-RELEASE</version> </dependency>3. 引入 @EnableDcsScheduling 开启分布式任务
与 SpringBoot 的@EnableScheduling非常像,@EnableDcsScheduling是中间件的统一入口注解,尽可能降低使用难度:
@SpringBootApplication @EnableDcsScheduling public class HelloWorldApplication { public static void main(String[] args) { SpringApplication.run(HelloWorldApplication.class, args); } }4. 在任务方法上添加 @DcsScheduled 注解
这个注解也和 SpringBoot 的@Scheduled很像,但多了desc描述和启停初始化控制:
cron:执行计划;desc:任务描述;autoStartup:默认启动状态(true表示启动后自动执行,false表示需在控制台手动启动)。
如果任务需要参数,可以通过引入 Service 去调用获取等方式。
@Component("demoTaskThree") public class DemoTaskThree { @DcsScheduled(cron = "0 0 9,13 * * *", desc = "03定时任务执行测试:taskMethod01", autoStartup = false) public void taskMethod01() { System.out.println("03定时任务执行测试:taskMethod01"); } @DcsScheduled(cron = "0 0/30 8-10 * * *", desc = "03定时任务执行测试:taskMethod02", autoStartup = false) public void taskMethod02() { System.out.println("03定时任务执行测试:taskMethod02"); } }5. 启动验证
- 启动 SpringBoot 工程即可,
autoStartup = true的任务会自动启动(任务以多线程并行方式执行); - 启动控制平台
itstack-middleware-control,访问http://localhost:7397/,即可在首页看到部署总数、服务总数、实例统计、任务总数等实时数据,并在任务列表中对任务执行「启动」「关闭」验证。
五、中间件开发:核心源码拆解
1. 工程模型
以 SpringBoot 为基础开发中间件,工程结构如下(org.itstack.middleware.schedule包):
schedule-spring-boot-starter └── src ├── main │ ├── java │ │ └── org.itstack.middleware.schedule │ │ ├── annotation │ │ │ ├── DcsScheduled.java │ │ │ └── EnableDcsScheduling.java │ │ ├── annotation │ │ │ └── InstructStatus.java │ │ ├── config │ │ │ ├── DcsSchedulingConfiguration.java │ │ │ ├── StarterAutoConfig.java │ │ │ └── StarterServiceProperties.java │ │ ├── domain │ │ │ ├── DataCollect.java │ │ │ ├── DcsScheduleInfo.java │ │ │ ├── DcsServerNode.java │ │ │ ├── ExecOrder.java │ │ │ └── Instruct.java │ │ ├── export │ │ │ └── DcsScheduleResource.java │ │ ├── service │ │ │ ├── HeartbeatService.java │ │ │ └── ZkCuratorServer.java │ │ ├── task │ │ │ ├── TaskScheduler.java │ │ │ ├── ScheduledTask.java │ │ │ ├── SchedulingConfig.java │ │ │ └── SchedulingRunnable.java │ │ ├── util │ │ │ └── StrUtil.java │ │ └── DoJoinPoint.java │ └── resources │ └── META_INF │ └── spring.factories └── test └── java └── org.itstack.demo.test └── ApiTest.java2. 自定义注解:@EnableDcsScheduling
注解上的一堆元注解,都是为了开始启动执行中间件:
@Target(ElementType.TYPE):标识需要放到类上执行;@Retention(RetentionPolicy.RUNTIME):注释将由编译器记录在类文件中,并且在运行时由 VM 保留,因此可以被反射读取;@Import:引入入口资源,在程序启动时会执行到自己定义的类中,方便初始化配置/服务、启动任务、挂载节点;@ComponentScan:告诉程序扫描位置(org.itstack.middleware.*,这也是自定义切面能否被扫描到的关键,否则自定义切面会失效)。
@Target({ElementType.TYPE}) @Retention(RetentionPolicy.RUNTIME) @Import({DcsSchedulingConfiguration.class}) @ImportAutoConfiguration({SchedulingConfig.class, CronTaskRegister.class, DoJoinPoint.class}) @ComponentScan("org.itstack.middleware.*") public @interface EnableDcsScheduling { }3. 扫描注解、初始化配置/服务、启动任务、挂载节点
注解已经写到方法上了,怎么拿到呢?核心在DcsSchedulingConfiguration:
- 通过实现
BeanPostProcessor.postProcessAfterInitialization,在每个 Bean 实例化的时候进行扫描; - 这里会遇到一个有趣的问题:一个方法会得到两次,因为有一个 CGLIB 代理出来的类,几乎一模一样。通过
method.getDeclaredAnnotations()判断"生命注解批注有没有",即可区分出真实方法; - 扫描下来的任务信息汇总到
Map中,等 Spring 初始化完成后,再执行中间件内容(太早执行会喧宾夺主,Spring 也不允许)。
@Override public Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException { Class<?> targetClass = AopProxyUtils.ultimateTargetClass(bean); if (this.nonAnnotatedClasses.contains(targetClass)) return bean; Method[] methods = ReflectionUtils.getAllDeclaredMethods(bean.getClass()); if (methods == null) return bean; for (Method method : methods) { DcsScheduled dcsScheduled = AnnotationUtils.findAnnotation(method, DcsScheduled.class); if (null == dcsScheduled || 0 == method.getDeclaredAnnotations().length) continue; List<ExecOrder> execOrderList = Constants.execOrderMap.computeIfAbsent(beanName, k -> new ArrayList<>()); ExecOrder execOrder = new ExecOrder(); execOrder.setBean(bean); execOrder.setBeanName(beanName); execOrder.setMethodName(method.getName()); execOrder.setDesc(dcsScheduled.desc()); execOrder.setCron(dcsScheduled.cron()); execOrder.setAutoStartup(dcsScheduled.autoStartup()); execOrderList.add(execOrder); this.nonAnnotatedClasses.add(targetClass); } return bean; }接下来初始化服务连接 Zookeeper 配置中心,连接后将创建节点并添加监听——这个监听主要负责分布式消息通知,收到通知后负责控制任务启停。这里包含了循环创建节点以及批量节点删除的逻辑:
private void init_server(ApplicationContext applicationContext) { try { //获取zk连接 CuratorFramework client = ZkCuratorServer.getClient(Constants.Global.zkAddress); //节点组装 path_root_server = StrUtil.joinStr(path_root, LINE, "server", LINE, schedulerServerId); path_root_server_ip = StrUtil.joinStr(path_root_server, LINE, "ip", LINE, Constants.Global.ip); //创建节点&递归删除本服务IP下的旧内容 ZkCuratorServer.deletingChildrenIfNeeded(client, path_root_server_ip); ZkCuratorServer.createNode(client, path_root_server_ip); ZkCuratorServer.setData(client, path_root_server, schedulerServerName); //添加节点&监听 ZkCuratorServer.createNodeSimple(client, Constants.Global.path_root_exec); ZkCuratorServer.addTreeCacheListener(applicationContext, client, Constants.Global.path_root_exec); } catch (Exception e) { logger.error("itstack middleware schedule init server error!", e); throw new RuntimeException(e); } }启动标记为true的 Schedule 任务。@Scheduled默认是单线程执行的,这里扩展为多线程并行执行:
private void init_task(ApplicationContext applicationContext) { CronTaskRegister cronTaskRegistrar = applicationContext.getBean("itstack-middlware-schedule-cronTaskRegister", CronTaskRegister.class); Set<String> beanNames = Constants.execOrderMap.keySet(); for (String beanName : beanNames) { List<ExecOrder> execOrderList = Constants.execOrderMap.get(beanName); for (ExecOrder execOrder : execOrderList) { if (!execOrder.getAutoStartup()) continue; SchedulingRunnable task = new SchedulingRunnable(execOrder.getBean(), execOrder.getBeanName(), execOrder.getMethodName()); cronTaskRegistrar.addCronTask(task, execOrder.getCron()); } } }最后把任务节点挂载到 Zookeeper。按照不同的场景,有些内容挂载到虚拟节点(创建的是永久节点,虚拟值通过子节点数据追加)。节点路径结构path_root_server_ip_clazz_method为:根目录 / 服务 / IP / 类 / 方法:
private void init_node() throws Exception { Set<String> beanNames = Constants.execOrderMap.keySet(); for (String beanName : beanNames) { List<ExecOrder> execOrderList = Constants.execOrderMap.get(beanName); for (ExecOrder execOrder : execOrderList) { String path_root_server_ip_clazz = StrUtil.joinStr(path_root_server_ip, LINE, "clazz", LINE, execOrder.getBeanName()); String path_root_server_ip_clazz_method = StrUtil.joinStr(path_root_server_ip_clazz, LINE, "method", LINE, execOrder.getMethodName()); String path_root_server_ip_clazz_method_status = StrUtil.joinStr(path_root_server_ip_clazz, LINE, "method", LINE, execOrder.getMethodName(), "/status"); //添加节点 ZkCuratorServer.createNodeSimple(client, path_root_server_ip_clazz); ZkCuratorServer.createNodeSimple(client, path_root_server_ip_clazz_method); ZkCuratorServer.createNodeSimple(client, path_root_server_ip_clazz_method_status); //添加节点数据[临时] ZkCuratorServer.appendPersistentData(client, path_root_server_ip_clazz_method + "/value", JSON.toJSONString(execOrder)); //添加节点数据[永久] ZkCuratorServer.setData(client, path_root_server_ip_clazz_method_status, execOrder.getAutoStartup() ? "1" : "0"); } } }4. Zookeeper 控制服务:监听下发启停指令
ZkCuratorServer提供一个 ZK 的方法集合,其中最重要的方法是添加监听。Zookeeper 的特性是:对这个路径添加监听后,当节点内容发生变化时会收到通知,宕机同样可以感知到——这正是后面开发灾备能力的核心触发点。
public static void addTreeCacheListener(final ApplicationContext applicationContext, final CuratorFramework client, String path) throws Exception { TreeCache treeCache = new TreeCache(client, path); treeCache.start(); treeCache.getListenable().addListener((curatorFramework, event) -> { //... switch (event.getType()) { case NODE_ADDED: case NODE_UPDATED: if (Constants.Global.ip.equals(instruct.getIp()) && Constants.Global.schedulerServerId.equals(instruct.getSchedulerServerId())) { //执行命令 Integer status = instruct.getStatus(); switch (status) { case 0: //停止任务 cronTaskRegistrar.removeCronTask(instruct.getBeanName() + "_" + instruct.getMethodName()); setData(client, path_root_server_ip_clazz_method_status, "0"); logger.info("itstack middleware schedule task stop {} {}", instruct.getBeanName(), instruct.getMethodName()); break; case 1: //启动任务 cronTaskRegistrar.addCronTask(new SchedulingRunnable(scheduleBean, instruct.getBeanName(), instruct.getMethodName()), instruct.getCron()); setData(client, path_root_server_ip_clazz_method_status, "1"); logger.info("itstack middleware schedule task start {} {}", instruct.getBeanName(), instruct.getMethodName()); break; case 2: //刷新任务 cronTaskRegistrar.removeCronTask(instruct.getBeanName() + "_" + instruct.getMethodName()); cronTaskRegistrar.addCronTask(new SchedulingRunnable(scheduleBean, instruct.getBeanName(), instruct.getMethodName()), instruct.getCron()); setData(client, path_root_server_ip_clazz_method_status, "1"); logger.info("itstack middleware schedule task refresh {} {}", instruct.getBeanName(), instruct.getMethodName()); break; } } break; case NODE_REMOVED: break; default: break; } }); }从代码可以看到,指令状态语义为:0= 停止任务、1= 启动任务、2= 刷新任务(先移除再按新 cron 注册)。每一次指令下发后,都会同步回写节点的status数据,保证控制台展示的状态与任务实际运行状态一致。
5. 并行任务注册:打破 @Scheduled 单线程限制
由于默认的 SpringBoot@Scheduled是单线程的,这里做了改造以支持多线程并行执行,包括添加任务和删除任务(即执行future.cancel(true)):
public void addCronTask(SchedulingRunnable task, String cronExpression) { if (null != Constants.scheduledTasks.get(task.taskId())) { removeCronTask(task.taskId()); } CronTask cronTask = new CronTask(task, cronExpression); Constants.scheduledTasks.put(task.taskId(), scheduleCronTask(cronTask)); } public void removeCronTask(String taskId) { ScheduledTask scheduledTask = Constants.scheduledTasks.remove(taskId); if (scheduledTask == null) return; scheduledTask.cancel(); }6. 待扩展的自定义 AOP
最开始配置的@ComponentScan("org.itstack.middleware.*"),主要就是服务于这里的自定义注解,否则是扫描不到的(即自定义切面失效的效果)。目前该切面并未扩展业务功能,基本只打印方法执行耗时;后续任务执行耗时监听等能力,可以基于这个切入点继续完善:
@Pointcut("@annotation(org.itstack.middleware.schedule.annotation.DcsScheduled)") public void aopPoint() { } @Around("aopPoint()") public Object doRouter(ProceedingJoinPoint jp) throws Throwable { long begin = System.currentTimeMillis(); Method method = getMethod(jp); try { return jp.proceed(); } finally { long end = System.currentTimeMillis(); logger.info("\nitstack middleware schedule method:{}.{} take time(m):{}", jp.getTarget().getClass().getSimpleName(), method.getName(), (end - begin)); } }六、Jar 包发布:让中间件可被 Maven 中央仓库引用
中间件开发完成后,还需要将 Jar 包发布到 Maven 中央仓库,这样使用者才能通过 POM 依赖直接引入。完整的发布流程(GPG 签名密钥生成、Sonatype 工单申请、Staging 仓库 Release、中央仓库同步)在仓库文档 《发布Jar包到Maven中央仓库,为开发开源中间件做准备》 中有详细步骤,其要点包括:
- 准备 GPG 密钥:Maven 中央仓库要求构件用 PGP 密钥签名以验证真实性,需要下载 GPG 工具生成 OpenPGP 密钥对并上传公钥到密钥服务器;
- Sonatype 工单申请:在工单系统创建 New Project 工单,填写 GroupId(与域名绑定,需通过 DNS TXT 记录验证域名归属)、Project URL、SCM 地址,等待人工审核;
- Staging 仓库发布:上传的 Jar 包先进入 Sonatype 的 Staging 仓库(如
orgitstackmiddleware-1000),Release 后即同步到 Maven 中央仓库,随后可在镜像仓库与阿里云仓库搜索到org.itstack.middleware:schedule-spring-boot-starter。
版本记录如下:
| 序号 | 版本 | 发布日期 | 备注 |
|---|---|---|---|
| 1 | 1.0.0-RELEASE | 2019-12-07 | 基本功能实现:任务接入、分布式启停 |
| 2 | 2019-12-07 | 上传测试版本 |
七、总结
从单机@Scheduled到分布式 DcsSchedule,本质上是把"任务扫描"从业务代码中抽离为独立的基础设施能力。本文以《分布式任务调度》章节的前言与需求背景为骨架,完整串联了 DcsSchedule 的落地全过程:
- 为什么需要:业务量增长后,单机定时任务无法在窗口期内完成扫描,需要把任务作业分散到多个服务实例节点;
- 怎么用:引入
schedule-spring-boot-starter,加@EnableDcsScheduling开启中间件,在任务方法上加@DcsScheduled(cron, desc, autoStartup),配合控制台统一启停与实时监控; - 怎么实现:
BeanPostProcessor扫描注解 → 汇总任务信息 → 连接 Zookeeper 创建/挂载节点 → TreeCache 监听下发启停指令 → 自定义CronTaskRegister多线程并行调度,最终通过 AOP 记录任务执行耗时,为后续监控能力留好扩展点。
在此基础上还可以继续深挖的点包括:分布式任务控制台itstack-middleware-control的具体实现(它只是使用中间件的 ZK 功能接口做展示和操作)、宕机灾备的完整触发链路、以及更多中间件设计与实现的源码对照(可继续阅读 第 3 章 服务治理·统一白名单控制、第 13 章 数据库路由组件 等同系列章节)。中间件开发是一件非常有意思的事情,不同于业务开发,它更像是对框架源码、数据结构、算法理论的最佳实践,也是程序员突破 CRUD 瓶颈、走向核心开发的必经之路。
- 文档
- 教程
- 后端
【免费下载链接】CodeGuide
:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总,旨在为大家提供一个清晰详细的学习教程,侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助,请给予支持(关注、点赞、分享)!
相关推荐
基于 SpringBoot 的分布式任务中间件 DcsSchedule:从 @Scheduled 到多机任务统一管控的完整设计与实现
基于 SpringBoot 的分布式任务中间件 DcsSchedule:从 @Scheduled 到多机任务统一管控的完整设计与实现 本文以 CodeGuide
文档教程后端从 CRUD 到中间件:基于 SpringBoot 的分布式任务调度中间件 DcsSchedule 设计与实现
从 CRUD 到中间件:基于 SpringBoot 的分布式任务调度中间件 DcsSchedule 设计与实现 导读 :单机定时任务(Spring 原生 @Sc
文档教程后端react-text-loop 未来展望:动画库的发展趋势与技术演进
react text loop 未来展望:动画库的发展趋势与技术演进 react text loop 作为一款轻量级的 React 文字动画库,通过简洁的 AP
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考