ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

分布式任务调度框架设计实践:从单机Cron到分片调度

分布式任务调度框架设计实践:从单机Cron到分片调度 说句实话我第一次听到ax调度这个名字的时候以为是某个内部项目的代号后来才知道是团队里沉淀下来的一套分布式任务调度组件。前前后后踩了不少坑也重构过两轮今天把这套东西从设计思路到落地细节都整理出来。如果你是那种还在用单机Cron、或者正准备引入调度框架的开发者这篇文章应该能帮你省掉很多弯路。1. 场景回顾多实例部署之后定时任务只剩下一地鸡毛1.1 最初的项目规模和任务底数先说背景。我们当时的业务线有8个微服务定时任务加起来40多个覆盖了订单超时关单、每日数据对账、凌晨报表生成、缓存预热、外呼任务扫描这些场景。一开始所有任务都是直接用Spring的Scheduled注解写在业务代码里Cron表达式写死在配置文件中每个服务自己跑自己的。表面上看起来没什么问题一个服务里面几个定时任务到点触发跑完拉倒。但等我们把服务从单实例扩展到多实例部署之后问题就全冒出来了。这也是很多团队都会经历的阶段代码还是那套代码拓扑变了以前隐藏的假设一下子全暴露了。1.2 单机Cron模式暴露出的四个痛点第一个痛点是重复执行。同样一段Scheduled代码部署了3个实例到点之后3个实例各跑一遍。你说加分布式锁加是可以加但锁只能保证同时只有一个实例在执行不能保证每次触发都落在同一个实例上而且锁的粒度、过期时间、异常释放这些问题都得自己处理写过的人都知道有多恶心。更关键的是加锁只是治标任务的状态、结果、触发历史还是一团黑。第二个痛点是负载不均。所有实例上的定时任务都是同时触发的比如凌晨1点整系统里十几个任务同时启动每个任务都要查库、算数据数据库连接池瞬间被打满主库CPU直接飙到80%以上。更无语的是其实只需要一个实例执行的报表任务3个实例都在抢锁抢不到锁的实例就干瞪眼白占线程资源。第三个痛点是任务不可观测。某个任务昨天跑了没、跑了多长时间、花了多少数据量完全没有记录。每次出问题都是感觉没跑然后人去日志里翻翻半天也不一定找得到。第四个痛点是运维成本。每新增一个定时任务就要改代码、发版、改配置。Cron表达式写错了只能等下一天才能发现根本没有试运行的概念。综合这些情况我们需要的是一个统一调度、支持动态注册、能把任务分片并行、还带完整执行日志的组件。这就是ax调度的由来。2. AX调度器的整体设计调度中心与执行器各管各事2.1 一次完整调度的生命周期ax调度总体上分为两部分ax-server调度中心和ax-worker执行器SDK。调度中心负责Cron表达式解析、任务触发、路由策略、调度记录存储。它不执行业务代码只负责到什么时间点应该让谁去做什么事。执行器SDK是嵌入到业务服务里的一个依赖负责和调度中心通信、维护心跳、接收任务指令、拉起本地线程池执行任务、上报执行结果。那么一次完整调度是什么样子的我用大白话描述一下ax-server在任务配置的Cron时间点到达后生成一条调度记录然后根据该任务配置的路由策略从在线执行器列表里面挑一个或几个实例执行器内部有一个长轮询线程每隔几秒去调度中心拉一次待运行的任务指令拉到之后就更新状态为运行中并提交到本地线程池真正执行业务代码跑完之后执行器把成功或失败的结果上报给调度中心调度中心把耗时、报错信息都写进调度日志表。这里有一个设计决策很关键为什么选拉取模式而不是推送模式因为业务服务可能部署在各种复杂的网络环境里比如多机房、网关后面、容器化平台中机器不见得能暴露一个对外端口让调度中心主动连进来。拉取模式只需要业务服务能访问到调度中心的地址就够了入站方向完全不用开端口接入成本低很多。2.2 核心数据表与任务状态机ax-server存储侧主要有三张表ax_job_info任务配置表。字段包括任务名称、Cron表达式、路由策略、分片总数、阻塞处理策略、超时时间、启停状态、负责人。ax_executor_group执行器注册表。每个接入ax-worker的服务有一个appName同一appName下的多个实例会周期性上报心跳表里记录实例IP、端口、最后心跳时间、在线状态。ax_job_log调度日志表。记录每次触发的调度时间、分片序号、落在哪个执行器实例、耗时、执行结果、错误信息。任务状态可以用一条流转线来概括待调度 - 待运行 - 运行中 - 成功/失败/超时。调度记录刚生成的时候是待调度到了触发点之后进入队列并变为待运行执行器拉走之后变为运行中最后根据执行结果落到终态。这里要强调一个容易被忽视的点任务配置ax_job_info和调度记录ax_job_log是分离的。任务配置是静态的调度记录是每触发一次就生成一条的动态数据。这种设计的好处是方便追溯每次执行情况也方便做失败重试——重试的是某一条调度记录而不是改配置。2.3 为什么采用分片而不是简单的分发在规划ax调度能力的时候我们就确定了一个原则只解决调度本身的问题不掺和具体业务逻辑。但怎么把一个任务在多台机器上并行跑这是调度层面必须回答的问题。最朴素的做法是分发一个任务只选一台机器执行。这能解决重复执行、全局管控的问题但解决不了单实例处理能力有限的问题。比如每天凌晨对20万条订单做对账一台机器算要跑两个小时这个时间成本会直接影响后续依赖任务的开始时间。所以我们在设计上直接支持了分片把一个任务在逻辑上切成N个分片shard每个分片内容独立可以由不同实例并行处理。分片总数和分片序号对业务代码是可见的业务方根据这两个参数计算自己该处理哪一部分数据。3. 接入AX调度的实操记录从引入依赖到跑通第一个任务3.1 调度中心的部署方式ax-server本身是一个Spring Boot应用我们通过Docker Compose部署了3个节点做高可用前面挂一个负载均衡数据库用的MySQL注册数据量不大所以没有额外引入Redis。接入的第一步其实很简单从内部制品库拉取ax-server的镜像配置好数据源地址启动之后打开管理页面看到执行器列表是空的说明调度中心已经就绪。管理页面主要做两件事维护任务配置、查看调度日志。这里提醒一句调度中心的数据库连接池一定要单独评估。因为每触发一次任务就要写一条调度日志任务量大之后写库这个动作会成为瓶颈我们后来把连接池从默认的10调到了30详见第6章的压测数据。3.2 业务服务接入执行器端的配置业务服务接入ax-worker更像加一个普通中间件依赖。如果是Maven项目pom里引入ax-worker的依赖然后在yml里加上配置ax: worker: admin-addresses: http://ax-admin.internal:8080 app-name: biz-order enable: true三个配置的含义admin-addresses是调度中心的地址多个地址用逗号分隔。执行器启动后会周期性向这个地址发起注册和心跳请求。app-name是执行器分组的名称。同一个业务服务部署了10个实例这10个实例的app-name应该是一样的调度中心才会把它们归属到同一个执行器分组里。enable是开关如果临时不想让这个服务参与调度可以动态置为false不用下线进程。接入之后服务启动日志里会看到一条ax worker register success的输出同时调度中心的执行器列表里会多出这个实例的IP和端口记录。如果列表里一直没出来90%是网络不通——先telnet一下调度中心地址的端口再看本地防火墙别一上来就查代码。3.3 一个对账任务从注册到成功执行的全过程接着我们注册了第一个任务每日订单对账。业务侧只需要写一个继承AxJobHandler的类Component Slf4j public class OrderCheckHandler extends AxJobHandler { Override public AxJobResult execute(AxJobContext context) { String shardTotal context.getShardTotal(); String shardIndex context.getShardIndex(); log.info(order check task start, shardIndex{}, shardTotal{}, shardIndex, shardTotal); ListLong orderIds orderMapper.selectNeedCheckIds(); for (Long orderId : orderIds) { // 按分片参数过滤每个分片只处理属于自己的那部分数据 if (orderId % Long.parseLong(shardTotal) ! Long.parseLong(shardIndex)) { continue; } orderCheckService.check(orderId); } return AxJobResult.success(); } }这段代码里最关键的就是那两行分片判断用主键ID对分片总数取模结果等于当前分片序号时才处理。这样不管任务被切成多少片、在哪台机器上跑只要所有分片都执行完毕全局数据一定被完整覆盖而且不会重复。然后在调度中心的管理页面新建一个任务任务名称每日订单对账Cron表达式0 0 1 * * ? 每天凌晨1点路由策略分片路由分片总数4阻塞处理策略丢弃后续调度超时时间1800秒配置好后先点一次手动触发到调度日志里看到4条记录全部成功耗时数据也正常再启用Cron定时触发。这里强烈建议新任务都先走一遍手动触发验证否则Cron写错或者业务代码有坑要等第二天才能发现。4. 分片执行与动态扩缩容最容易被忽略的幂等细节4.1 分片参数与业务代码的配合方式分片一共分为两步调度中心切分业务代码消化。调度中心负责把任务切成分片总数份每一份生成一条调度记录业务代码拿到当前分片的序号之后自己去决定这段数据怎么处理。前面订单对账的示例里用的是对ID取模的方式这是最简单、最容易理解的做法适合主键ID均匀分布的表。但如果数据表的ID不是连续的或者业务上更适合按日期、按商家维度切分也可以把分片参数直接暴露给任务代码让业务方自己决定怎么切。比如有个积分补发任务需要给一批用户补发积分。我们可以把用户ID当作切分维度同样用取模的方式来决定每个分片处理哪些用户。再比如有任务想按日期范围处理那可以在分片逻辑里用日期加序号做位移每个分片处理一周的数据。分片没有固定的公式重要的是保证两点所有分片处理的集合是全集任意两个分片处理的集合没有交集。4.2 节点变化时任务归属怎么重算ax调度会把分片总数作为一个固定配置但它和在线实例数不绑定。比如任务配置分片总数是10在线执行器只有3个那么10个分片会被分配到3台实例上——有的机器分到3片有的分到4片。如果某个执行器实例挂了调度中心的心跳检测会在约90秒后把它标记为离线它手上分到的那些分片会被重新计算分配给其他还活着的实例。这带来的好处是分片自身是稳定的分片序号不会因为实例数变化而改变。业务代码中取模的边界完全由分片总数决定只要这个值不变同一类数据始终落在同一个分片序号上。这一点对幂等非常关键——你不能因为机器少了一台分片序号也跟着跳否则一批数据可能会从分片3跳到分片5状态记录就乱套了。当然如果线上实例数经常变化比如弹性扩缩容很频繁我建议把分片总数设得大一些比如20甚至50。分片多了之后单台实例故障导致的数据转移范围会更小任务的总体验收时间也不会被拖得太明显。4.3 为什么分片之后仍然必须做幂等校验一个常见的误解是我分片了数据就不会重复处理了。这话只对了一半。分片解决的是多台机器并行跑同一份任务时的重叠问题但解决不了任务重试和重复触发带来的重复问题。举个例子某个分片执行到一半执行器所在节点发布了新版本服务正在重启调度中心因为超时把这个分片标记为失败。任务配置了失败重试于是重试请求又过来了。这时候最原始的取模过滤是拦不住同一批数据的——因为逻辑上它就是同一个分片该处理的数据依然还是那些。所以分片任务必须有幂等防线。常用的思路有三种第一种是数据库唯一约束比如对账结果表里以订单ID建立唯一键重复插入直接报错后捕获掉第二种是用Redis做一个分布式防重标记处理前先SETNX一下已存在就跳过第三种是业务状态跳过比如订单已经是已对账状态就直接不处理。三种方式可以组合使用核心原则是哪怕同一个分片被物理执行了两次业务数据也只应该被有效处理一次。5. 线上深度排查三个真实故障的处理全过程5.1 执行器明明在线任务却全部路由到同一台机器这个故障发生在一次扩容之后。我们给某个执行器分组扩到了6台实例任务的路由策略配的是轮询——按道理每次触发应该轮流落在不同的机器上。但实际运行了一段时间后执行器列表里6个实例都在线某个高频任务却几乎只在一台机器上跑。排查链路是这样一步步走的。我先去调度中心的任务配置页面看了一眼路由策略没错配的确实是轮询。接着去看执行器列表里6个实例的注册时间发现其中5个是昨天扩容后新注册的还有1个是几天前的旧实例。然后我对比了旧实例对应的服务进程发现它实际已经重启过一次但注册信息里的端口还是旧的。问题就在这服务重启后主机名没变、IP没变但注册端口变了旧端口的心跳虽然还在续实际对应的进程却已经换了。当时ax-server对执行器实例的唯一标识是IP端口appName服务重启后如果端口是随机分配的就会在注册表里产生一条新记录同时旧记录因为心跳还在没有被清除于是轮询算法在6条记录里轮转其中有几条指向的其实是同一台机器的不同端口。根因清楚之后修复方案是两步第一步把执行器实例的唯一标识改为实例ID每次服务启动时生成一个全局唯一ID上报避免用端口做身份第二步清理掉那些端口已失效但心跳仍在的僵尸注册记录并且把执行器的端口从随机改为固定端口方便排障。这个坑给我们的教训是任何基于IP端口做实例身份的系统在容器化和弹性扩缩容场景下都会出问题。实例身份必须是一个稳定的、生命周期等于进程生命周期的ID而不是网络层面的属性。5.2 长耗时任务被重复触发账务数据差点算错第二个故障更严重。一个数据重跑任务单次执行需要30分钟到1个小时Cron配置是每5分钟触发一次。结果这个任务在跑的过程中调度中心每5分钟又触发一次新的调度记录执行器拉回去之后发现上一个任务还在跑又把这个新指令也提交到线程池了——于是同一批数据被多个线程反复处理。这个问题的直接原因是任务配置里阻塞处理策略选错了。ax调度提供了三种阻塞策略丢弃后续调度、覆盖之前调度、并行执行。当时这个任务被配置成了并行执行语义是不管上一个有没有跑完新的触发直接开跑这是一个为短任务设计的选项用在长任务上就是事故。根因找到之后我们做了两处修正。第一处把该任务的阻塞策略改为丢弃后续调度也就是说如果上一个还没跑完后续的Cron触发直接丢弃避免堆积。第二处在任务配置里设置合理的超时时间超时后调度中心会强制将该调度记录标记为超时失败这样即使任务异常卡住也不会无限占住资源。这里要提醒的是超时时间和阻塞策略要一起看光设置超时但允许并行执行依然会出现重复。最稳妥的长任务配置组合是丢弃后续调度 明确的超时时间 业务侧幂等。这张组合拳打下来长任务基本就不会出大问题。5.3 调度日志疯狂膨胀数据库磁盘告警第三个问题来得比较隐蔽。某天监控告警说数据库磁盘使用率超过了85%上去一看最大的表不是业务订单表而是ax_job_log调度日志表。这个表60天内的数据堆积到了25GB而且还在以每天600MB的速度增长。原因有两层。第一层是日志清理任务没有配置默认的清理逻辑虽然是有的但需要创建一个定时任务去触发它我们当时没有建这个任务相当于自动清理一直没生效。第二层是调度记录里的错误信息字段存了完整堆栈一个异常堆栈动不动就是几千个字符几百次失败任务轻轻松松就能把存储撑爆。修复分三步走。第一步马上给日志表按日期做分区按月归档旧分区磁盘压力立刻缓解第二步启用调度中心的日志清理任务每天晚上清理90天前的数据第三步限制保存日志的详细程度错误堆栈只保留前100个字符完整堆栈放对象存储里按需拉取。这个故障给我们的启示是调度组件看起来是个辅助系统但它产生数据的速率可能是业务系统的数倍。所有中间件在规划存储的时候都要考虑数据生命周期管理别等磁盘告警了才想起来要清理。6. 压力测试与参数调优中的一些实测数据6.1 压测环境与测试方法ax调度整体稳定运行了一段时间后我们做了一次针对性的压测目的是回答两个问题调度中心能支撑多大的触发频率哪些参数会先成为瓶颈。压测环境是3台ax-server4核8G1台MySQL8核16G10个执行器实例。测试方式是通过管理接口批量创建1000个Cron任务触发频率统一为每分钟1次记录调度中心的触发延迟从Cron到点时间到调度记录写入完成和执行失败率。6.2 压测数据与瓶颈分析压测结果整理成了表格任务规模平均触发延迟P99触发延迟执行失败率主要瓶颈500个/分钟23ms89ms0%无明显瓶颈2000个/分钟58ms232ms0.3%数据库写入开始出现锁等待5000个/分钟130ms621ms2.1%调度中心线程池排队日志表写入变慢从数据能明显看出2000个/分钟以内是一个比较舒适的范围超过之后数据库写日志开始成为主要瓶颈。执行器端反而没出现太大压力10个实例分摊下来每个实例每分钟只需要处理几百个任务指令。6.3 根据压测结果做的参数调整压测暴露的最大问题是日志写入。我们把调度日志的写入从同步改成了异步批量先攒一批调度记录再一次性批量插入MySQL。这个改动把2000个/分钟场景下的P99延迟从232ms降到了80ms左右效果非常明显。另外调度中心的调度线程池从默认的16调到了32但这个调大作用有限在5000个/分钟时线程本身不是瓶颈内存里的任务队列才是。我个人建议的参数如下调度线程数不超过CPU核数的2倍别盲目调大调度中心的秒级触发并不需要太多线程线程多了反而增加上下文切换开销。执行器端的任务线程池要按任务类型做隔离比如快速任务的线程池配20个长任务的线程池单独配5个防止长任务把公共线程池占满拖垮短任务。调度日志表必须提前规划分区和归档策略不要在数据量上来之后再补救。执行器心跳周期设置为30秒离线判定阈值设为90秒。心跳太频繁会白白消耗调度中心的CPU太慢会让故障转移迟钝。以我在一线维护调度系统的经验来说调度中心本身的架构并不复杂真正难的是参数调优和边界问题的处理——比如实例身份识别、分片稳定性、超时与阻塞策略的组合、日志存储的生命周期管理。把这些细节都处理到位之后ax调度才算真正变成一个让人省心的基础组件。后续如果想进一步扩展可以考虑往工作流编排的方向走让多个任务之间支持依赖关系和条件分支那会是另一个层面的话题了。
返回列表