ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

CodeGuide 分布式任务调度实战:从 @Scheduled 单机定时任务,到 Zookeeper 驱动的 DcsSchedule 中间件

CodeGuide 分布式任务调度实战:从 @Scheduled 单机定时任务,到 Zookeeper 驱动的 DcsSchedule 中间件 文档教程后端【免费下载链接】CodeGuide:books: 本代码库是作者小傅哥多年从事一线互联网 Java 开发的学习历程技术汇总旨在为大家提供一个清晰详细的学习教程侧重点更倾向编写Java核心内容。如果本仓库能为您提供帮助请给予支持(关注、点赞、分享)项目地址https://gitcode.com/gh_mirrors/code/CodeGuide点击查看免费下载本文以《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() { //... } }但是如果任务是比较大型的比如定时跑批 T1 结算、商品秒杀前状态变更、刷新数据预热到缓存等等这些定时任务都有相同的特点作业量大、实时性强、可用率高。而这时候如果只是单纯使用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.8SpringBoot 2.x配置中心 Zookeeper完整版文档中调试环境使用 3.4.14本仓库小册 第 2 章 给出的开发环境为 Zookeeper 3.6.0。准备好 Zookeeper 服务后下载解压在 bin 同级路径创建data、logs文件夹修改conf/zoo.cfg配置数据与日志目录例如dataDirD:\Program Files\apache-zookeeper-3.4.14\data dataLogDirD:\Program Files\apache-zookeeper-3.4.14\logs打包部署控制平台itstack-middleware-control部署后访问http://localhost:7397即可打开控制台。2. 配置 POM 依赖dependency groupIdorg.itstack.middleware/groupId artifactIdschedule-spring-boot-starter/artifactId version1.0.0-RELEASE/version /dependency3. 引入 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; ListExecOrder 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); SetString beanNames Constants.execOrderMap.keySet(); for (String beanName : beanNames) { ListExecOrder 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 { SetString beanNames Constants.execOrderMap.keySet(); for (String beanName : beanNames) { ListExecOrder 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 单线程限制由于默认的 SpringBootScheduled是单线程的这里做了改造以支持多线程并行执行包括添加任务和删除任务即执行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-1000Release 后即同步到 Maven 中央仓库随后可在镜像仓库与阿里云仓库搜索到org.itstack.middleware:schedule-spring-boot-starter。版本记录如下序号版本发布日期备注11.0.0-RELEASE2019-12-07基本功能实现任务接入、分布式启停21.0.1-RELEASE2019-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核心内容。如果本仓库能为您提供帮助请给予支持(关注、点赞、分享)项目地址https://gitcode.com/gh_mirrors/code/CodeGuide点击查看免费下载相关推荐基于 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上一篇基于 Rube MCP 自动化 Fingertip 操作awesome-codex-skills 中 fingertip-automation 技能实战指南下一篇JuiceFS Hadoop Java SDK 完全指南在 Hadoop 生态中平滑接入 JuiceFS创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表