ARTICLE DETAIL

资讯详情

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

自研分布式任务调度系统ax:时间轮、分片与故障转移实践

自研分布式任务调度系统ax:时间轮、分片与故障转移实践 后端干久了你会发现“定时任务”四个字能撑起半个分布式系统的复杂度。我主要负责交易链路那块报警群里被问得最多的就是“怎么又有任务没跑”“谁把同一张表跑重了”。前前后后折腾过好几套方案之后我整理出一套适合自己团队的调度设计内部代号就是“ax”我们一般叫它 ax调度。它不是什么黑科技核心就解决三件事任务什么时候触发、交给哪台机器执行、失败之后怎么办。这篇文章把 ax调度的设计思路、实现要点和落地过程完整拆一遍适合正在用 Quartz 救火、或者准备自研调度系统的后端同学参考。1. 为什么我不直接用现成的调度框架而是自己搞了套 ax调度每次聊任务调度第一反应肯定是“为什么不用现成的”。市面上 Quartz、ElasticJob、XXL-JOB 都成熟得不行直接拿来用不香吗我一开始也是这样想的但真正在业务里跑一段时间后发现“能用”和“好用”之间有很长一段路。1.1 Quartz 的分布式困境Quartz 在单机时代是绝对王者但一旦上集群痛点就非常明显。它的分布式方案是靠数据库行锁实现的调度线程定时去 triggers 表里抢记录抢到锁的节点才有资格执行触发逻辑。这套机制在小规模场景没问题可当触发器表到了几万条、数据库连接一旦紧张锁竞争就会让调度延迟从秒级恶化到分钟级。我在压测环境里实测过仅仅 50 个并发调度线程、单表 5 万条 trigger 记录就有接近 3% 的触发动作延迟超过 30 秒。对于财务对账、库存同步这类对时效敏感的任务30 秒意味着业务已经出了可感知的误差。更头疼的是双机房部署时的网络抖动。节点 A 刚拿到锁准备触发节点 B 因为网络分区误判自己失联转而竞争同一把锁等网络恢复两个节点同时认为自己是 Leader结果就是同一个任务被触发两次。Quartz 本身没有很好的幂等保护业务侧只能自己加分布式锁这个成本往往被很多人低估。1.2 大而全调度平台的改造成本后来我也认真调研过 XXL-JOB 这类功能齐全的调度平台。它们确实很强有可视化控制台、有权限体系、有动态任务管理拿来即用。但我们的团队规模并不大核心服务不到 20 个真要引入一套完整的平台学习成本和维护成本会压在本来就紧张的研发资源上。控制台、权限、多租户、归档报表这些功能对我来说属于“偶尔用到、平时添乱”的范畴。而且这类平台的扩展点虽然多但真要改核心调度逻辑需要熟读大量源码对团队里每个接手的人来说都是负担。我需要的不是大而全而是小而可靠调度器专注时间计算和任务分发执行器专注跑任务和回传状态中间通信尽可能薄。这就促成了 ax调度的雏形。1.3 ax调度的设计取舍ax调度的架构就三个角色调度中心ax-server、执行器ax-worker、元数据库。调度中心只干两件事维护任务元数据、计算触发时间并下发执行指令。执行器只干三件事注册自身地址、接收指令执行任务、上报日志和结果。通信走 HTTP 端点请求携带全局唯一的 batchId执行批次号从触发到执行结束全链路可以通过 batchId 把日志串起来。这套拆分的直接收益是执行器可以做到语言无关。调度中心下发的是通用指令任何语言只要实现一个 HTTP 回调接口就能接入。我在团队里就有两个 Python 写的任务进程注册方式跟 Java 服务一样只是回调接口略有区别。这个特性在微服务多语言环境下特别实用。2. 核心细节解析与实操要点很多调度框架的原理文档写得抽象落地时全靠自己猜。下面把 ax调度里我认为最核心的几个模块拆开讲偏实现向尽量说清楚“为什么这么设计”。2.1 调度引擎的“多层时间轮 延迟队列”混合模型时间轮是解决大量定时任务扫描问题的经典方案核心思想是把时间分成一个个 slot用数组存待触发的任务指针每 tick 跳一个 slot命中哪个槽就触发哪个槽上的任务链表。相比每秒钟全表扫一遍数据库时间轮在插入和触发上都能做到 O(1) 级复杂度。ax调度里用的是两层时间轮第一层 512 个 slot每个 slot 对应 200ms所以能覆盖 102.4 秒第二层 60 个 slot每个 slot 对应 1 分钟用于存放更远的触发任务。第一层指针转完一圈会把第二层中最近一分钟的任务降级到第一层里。这种设计避免了单层时间轮内存占用过大的问题。但纯时间轮有个缺陷任务量稀疏的时候指针空转浪费 CPU。所以在 ax调度里增加了一个容量约 1000 的 DelayQueue用来缓存最近 5 分钟内需要触发的任务。调度线程从 DelayQueue 里阻塞取最近到期的任务到期后再根据任务精确时间放入时间轮。这样即使任务数量很少线程也能休眠等待而不是空转。这里有一个很容易踩的坑调度中心的时间精度不能完全依赖系统 tick。调度进程需要对 NTP 时间同步敏感一旦宿主机时钟漂移超过 500ms调度触发的准确性就无从谈起。ax调度专门加了一个“时钟偏差检测”模块每隔 30 秒跟元数据库时间做对比偏差过大直接在当前任务日志里打 WARN。2.2 任务分片、路由与故障转移策略任务触发之后调度中心要决定把这批任务发给谁。ax调度支持三种基本路由策略轮询、一致性哈希、分片广播。轮询适合任务短时间内多、但单个执行成本低的场景比如批量发送通知自动均匀分散到各个 worker。一致性哈希适合需要把同一业务维度积累到同一节点的场景比如按用户 ID 做缓存预热保证同一个用户总是落在同一台执行器上提高本地缓存命中率。分片广播则是我用得最多的策略把一个大任务切分成 M 个分片每个执行器拿到分片号shardingId和总分片数shardingTotal自行处理属于自己的一部分数据。分片任务最大的坑是数据边界不清晰。如果分片逻辑是“订单号取模后等于当前分片号”那新增 worker 会导致分片总数变化老任务可能漏数据。所以我在 ax调度里规定执行器启动后分片总数是「当前在线 worker 数 × 每个 worker 允许并发数」动态扩容时并不立即改变已经下发任务的 shardingTotal而是等下一轮触发再生效。虽然短期可能存在资源不均但换来了数据一致性。故障转移的处理我设计成“先标记后剔除”。执行器每 10 秒发一次心跳调度中心连续 3 次没收到心跳就把它标记为“可能失联”此时不立即摘除而是再等 15 秒尝试一次主动探测。如果探测失败才从在线列表摘除并把尚未完成执行的批次重新路由到其他节点。这种两段式方案能有效避免因为网络瞬断导致的大规模任务迁移。2.3 触发规则、线程池和重试参数怎么定调度系统的参数配置如果不解释清楚就是黑魔法。先说触发规则ax调度支持 cron 表达式和固定间隔两种。cron 表达式解析后先转成下一次触发时间点再通过时间轮排程。值得注意的是cron 表达式解析要显式指定时区默认用 Asia/Shanghai不要用服务器系统时区否则服务器时区被改一下所有任务秒变“时差任务”。执行器收到指令后会丢进自己的业务线程池。线程池参数我一般这样建议核心线程数 CPU 核心数 × 2 1最大线程数 核心线程数 × 2有界队列容量 1024拒绝策略用 CallerRunsPolicy 的改进版不阻塞本次任务执行而是立即把失败原因写进 execution_log 表同时触发告警。这里要特别说明调度场景最怕的是“队列积压但任务还在排队”的假象。如果你把队列设成无界遇到下游接口变慢任务会一直在队列里堆着执行时间远晚于调度时间排查时根本看不出是哪一环拖的。有界队列配合拒绝策略至少能让问题立刻暴露出来而不是让延迟“温水煮青蛙”。重试策略我经历过从“无脑重试”到“谨慎重试”的转变。默认重试 3 次重试间隔指数退避1 秒、2 秒、4 秒。超过 3 次就不重试转人工/告警。为什么不用更大的重试次数因为任务失败的重试是有成本的如果目标系统已经过载重试只会加剧雪崩如果任务是不可幂等的比如发短信重试还会造成重复发送。所有重试逻辑必须配合幂等令牌使用这个令牌就是 batchId。2.4 两套核心表结构设计调度中心能稳定运行元数据表结构非常关键。ax调度只保留两张核心表任务表 job_info 和 执行日志表 execution_log。job_info 主要字段包括任务主键、应用名 app_name、处理器标识 handler、cron 表达式、路由策略、超时时间、重试次数、是否开启并发执行、最近触发时间、下次触发时间、状态。索引重点落在 app_name 和 next_trigger_time 上调度扫描时就靠这两个字段定位任务。execution_log 主要字段包括日志主键、任务主键、batch_id、触发时间、实际执行时间、worker 地址、执行状态、耗时、错误信息、日志快照。索引建议落在 job_id batch_id 和 trigger_time 上。batch_id 一定要有唯一索引这是防止调度中心重复下发的最后一道防线即使调度侧逻辑出错插入同一 batch_id 也会失败从而中止重复执行。表设计要说一个血泪经验不要把任务详情和执行日志混在一张表。最开始我图省事把 handler 参数直接拼在任务表的字段里结果日志查询稍微变多任务表就出现锁竞争调度扫描也被拖慢。拆成两张表后调度链路的查询路径非常干净日志表的膨胀也不会影响调度性能。3. 实操过程与核心环节实现讲完原理进入正文操作。以一个订单服务为例从零搭建 ax调度环境、注册执行器、跑通第一个调度任务。3.1 十分钟搭起调度中心环境调度中心我打包成了一个独立可执行 Jar依赖一个 MySQL 库。启动前先初始化 SQL 脚本创建 ax_dispatch 数据库脚本里建好 job_info 和 execution_log 两张表以及基础索引。mysql -u root -p -e CREATE DATABASE ax_dispatch DEFAULT CHARACTER SET utf8mb4; mysql -u root -p ax_dispatch sql/init.sql然后启动调度中心wget https://mirrors.example.com/ax-server-1.0.0.tar.gz tar xzf ax-server-1.0.0.tar.gz cd ax-server sh bin/start.sh --server.port8081 --spring.datasource.urljdbc:mysql://127.0.0.1:3306/ax_dispatch启动后访问 http://127.0.0.1:8081/actuator/health 返回 UP 说明调度中心起来了。调度中心自身不存储执行逻辑只负责计算、下发行令和维护 worker 列表所以它的重启对已注册执行器是无感的正在执行的任务不会被中断。3.2 注册执行器并跑通第一个任务执行器端接入很简单。我是 Maven 项目加依赖dependency groupIdcom.ax/groupId artifactIdax-executor-spring-boot-starter/artifactId version1.0.0/version /dependency配置文件 application.yml 里加ax: executor: app-name: order-service address: 127.0.0.1:9099 registry-url: http://127.0.0.1:8081/ax/registry注意 address 是执行器对外暴露的回调地址在容器环境里要配置成宿主 IP 和映射端口不能配成 pod 内部 IP否则调度中心无法回调。这个我在后面容器化章节还会再提。写一个最简单的任务Component public class HealthCheckJob implements SimpleJob { Override public JobResult execute(JobContext ctx) { log.info(health check start, batchId{}, ctx.getBatchId()); return JobResult.success(); } }执行器启动后日志中会出现“register to scheduler success”的字样。接着在调度中心控制台新建任务应用名 order-service、处理器标识 healthCheckJob、cron 表达式每 5 分钟一次保存后即可看到下一次触发时间。到点后执行日志里会出现 batchId、起止耗时和状态码。到这里第一个 ax调度任务就跑通了。3.3 分片任务批量订单数据处理的完整示例光跑通 Hello World 没意思看一个真实场景每天晚上要同步前一天的订单数据到数仓订单量有几百万单台机器跑要几个小时必须分片。任务实现采用分片模式Component public class OrderSyncShardingJob implements ShardingJob { Override public JobResult execute(JobContext ctx) { int shardId ctx.getShardingId(); int shardTotal ctx.getShardingTotal(); int batchSize 5000; long minOrderId ctx.getJobParam().getLong(minOrderId); long maxOrderId ctx.getJobParam().getLong(maxOrderId); // 每个分片只处理属于自己范围内的订单 long rangeLen (maxOrderId - minOrderId) / shardTotal; long startId minOrderId rangeLen * shardId; long endId (shardId shardTotal - 1) ? maxOrderId : startId rangeLen; ListOrderDO orders orderMapper.selectRange(startId, endId, batchSize); for (OrderDO order : orders) { syncToWarehouse(order, ctx.getBatchId()); } return JobResult.success(); } }控制台配置路由策略为“分片广播”假设当前 4 台 worker任务触发后每台机器拿到的 shardingId 分别是 0、1、2、3shardingTotal 4。由于上面代码里每个分片只处理订单 ID 区间的一部分合计起来正好覆盖全量不会重复也不会漏单。实际跑起来我遇到过一个边界问题如果订单 ID 不是连续分布中间有空洞按 ID 区间分片会导致分片 1 可能比分片 2 多处理几十万条。后来我改成先统计 min/max再在任务参数里携带全量 ID 集合的布隆过滤器每个分片查询时本地过滤一遍空洞。虽然多了一点内存消耗但负载均衡效果好了很多。3.4 配置监控告警提前发现问题分布式调度最怕“没跑”和“跑重”而这两类问题靠日志事后查是非常痛苦的。ax调度通过暴露 Prometheus 指标让你提前感知风险# HELP ax_dispatch_task_trigger_total 触发任务总数 # TYPE ax_dispatch_task_trigger_total counter ax_dispatch_task_trigger_total{apporder-service} 1280 # HELP ax_dispatch_task_failure_total 失败任务总数 # TYPE ax_dispatch_task_failure_total counter ax_dispatch_task_failure_total{apporder-service} 3在 Prometheus 里配置抓取调度中心的 /actuator/prometheus 端点然后配两条核心告警规则一是调度延迟任务实际执行时间减去触发时间超过 15 秒触发 Warning二是失败率5 分钟内失败率超过 10% 触发 Critical。我自己的习惯是给每个重要任务额外加一条“SLA 未触发告警”——如果一个任务在预定窗口内完全没有触发记录说明调度链路本身可能挂了这比失败告警能更早暴露问题。4. 常见问题与排查技巧实录不管设计多完整落地过程总会有各种幺蛾子。这里把我在 ax调度使用中踩过的坑集中列一下附上排查思路方便遇到类似问题的人直接定位。4.1 任务总是延迟几十秒才触发现象cron 设的是整点触发但 execution_log 里记录的实际执行时间比触发时间晚了 30 秒甚至 1 分钟。排查第一步看调度中心的线程池是否打满。ax调度的下发动作会经过一个独立的 IO 线程池如果业务任务里有人直接写了阻塞式 HTTP 调用且超时设成 60 秒这个线程池会被占满后续任务全部排队。解决办法有两个层面任务代码里禁止在调度执行线程中调用不可控外部接口改成异步提交把同步下发改成“批量聚合下发”同一秒触发的任务合并成一条请求发给执行器降低 IO 次数。我在 4.1 节中实测发现聚合下发能减少约 40% 的下发延迟波动。另一个隐蔽原因是线程池配置把核心线程数设得太小默认 2 个线程扛不住瞬时 20 个任务同时触发的尖峰。调度线程池的队列要有界但线程数不能太保守建议按“单机任务触发峰值 × 1.5”来设计。4.2 同一个任务被重复执行重复执行是调度系统最严重的故障之一。我遇到过一次典型的重复场景调度中心在任务执行超时后判定失败立即重试路由但第一次执行的线程并没有被真正终止数据库查询卡住了而已导致同一个批次两个 worker 同时在跑。这是“超时误判 不可中断操作”共同作用的结果。对应的防御手段我总结成三个层级第一层数据库约束。execution_log 表为 job_id batch_id 建立唯一索引重复下发会直接插入失败第二层执行器内部并发控制。每个任务在执行前先对 job_id 加分布式锁锁的 key 是 job_id batch_id用 Redis 的 SETNX 实现过期时间设为任务超时时间的两倍第三层严格禁止并发。在任务配置里强制开启“并发运行阻止”同一任务即使被触发多次也只会有一个执行中的实例。三层都做了也不能说 100% 安全但至少能挡住绝大多数误操作。4.3 长任务导致的任务堆积有些任务是跑 3 小时才算完的批处理任务如果 cron 是 2 小时一次第二次触发时第一次还没结束。任务会越积越多每个任务还都在抢数据库连接最终拖垮整个服务。ax调度里我建议对这类任务单独设置“最大执行时间”超时后触发中断信号。注意中断不是简单的 Thread.stop而是标记线程中断状态并把执行结果写成“超时失败”让业务代码在数据库操作处感知到异常退出。如果中断后任务还是没停比如卡在了不可中断的锁等待则配合线程池的 discardPolicy 强制丢弃该批次的新触发请求。更治本的手段是分片 增量。把全量同步改成上一天增量用 offset 记录进度每次任务只处理新增数据这样单次执行时间能压到 10 分钟以内彻底摆脱“任务跑不完”。4.4 调度中心与执行器时钟偏差分布式的世界里“时间”是个伪命题。我遇到过一个问题执行日志里显示任务实际执行时间比触发时间早 2 秒导致链路追踪的时间线倒挂排查问题的时候非常误导。原因是执行器的宿主机时钟快了 2 秒它记录的是本地时间而调度中心记录的是自身时钟时间。解决方案是统一以调度中心的时间作为事实来源。触发时间、计划时间都用调度中心生成的时间戳执行器只记录 CPU 耗时、内存等相对指标不再记录本地时间用于链路排障。需要真实时间做业务判断时从调度中心下发的报文头里取而不是本地 System.currentTimeMillis()。同时所有机器必须强制配置 NTP 同步偏差超过 500ms 时监控告警。4.5 调度问题排查速查表我把上面这些经验浓缩成一张排查表贴在团队内部文档里遇到问题直接对照现象可能原因排查方向任务触发延迟大调度中心线程池阻塞查调度中心活跃线程数、队列堆积数同一任务执行多次超时误判、路由重复下发查 execution_log 中 job_id batch_id 是否唯一任务长时间排队执行器线程池满或任务不可中断查执行器活跃线程数、等待队列长度执行时间倒挂执行器本地时钟漂移统一走调度中心时间戳校准 NTP日志数据缺失日志上报链路被阻塞查执行器到调度中心日志端点的网络状态任务分片数据不均分片区间数据分布不均改用一致性哈希或实现消费进度记录这张表最大的价值是帮团队节省了“猜”的时间。遇到调度问题先看表中对应的最可能原因再通过命令查指标基本能在 5 分钟内定位到问题层。5. 容器化部署与容量规划现在部署基本都往 K8s 走了调度系统也要适配容器环境。这一章把我自己在容器化落地过程中的要点写清楚。5.1 K8s 下怎么部署 ax调度调度中心无状态可以做成 Deployment 跑两个副本通过 Service 对外暴露端口。但底层依赖的 MySQL 不建议也容器化至少生产环境用云数据库会更稳。调度中心多副本部署时副本之间会通过数据库锁做任务触发互斥所以不用担心两个副本同时触发同一个任务但如果用了数据库锁连接池大小要给足否则高并发下锁等待会拖慢调度。执行器部署稍微复杂点每个 Pod 就是一个执行器实例注册地址必须配成宿主机 IP 加 NodePort 映射出来的端口。如果直接用 Pod IP调度中心从集群外访问不到。我用的是 StatefulSet 部署执行器每个 Pod 有个稳定网络标识配合 headless service 做回调。Pod 生命周期结束时preStop 钩子里调用执行器提供的 unregister 接口把这个节点从调度中心在线列表里摘除避免已经下线的执行器继续被路由。lifecycle: preStop: exec: command: [sh, -c, curl -X POST http://127.0.0.1:9099/ax/executor/unregister; sleep 5]这个 5 秒的 sleep 是为了等调度中心完成摘除操作再真正终止容器否则正在执行的任务会被强制杀掉。5.2 多租户与权限隔离团队大了之后不同业务线都要用调度平台直接混在一个命名空间里很不安全。ax调度在应用层面支持简单的多租户隔离每个任务绑定一个 namespace执行器注册时声明自己属于哪个 namespace。调度中心只把任务路由到相同 namespace 下的执行器。控制台权限上按角色分为管理员、开发者、只读三种。管理员能管理所有 namespace开发者只能在自己的 namespace 里创建和修改任务只读只能查看日志和监控。这种隔离粒度对中小团队足够用了没必要一开始就上完整的 RBAC 权限系统不然又是维护负担。等需要更细粒度的权限时再基于 OpenID Connect 做一层对接也不迟。5.3 生产环境需要多大规格的机器容量规划不能靠感觉。我自己跑过一轮压测作为参考调度中心使用 8C16G 的虚机、MySQL 使用 4C8G 的云数据库单调度中心能支撑 5000 个任务峰值触发频率约 120 次/秒端到端调度延迟 P99 在 80ms 以内。这个数据是在任务逻辑都比较简单的情况下测的如果你的执行器处理很重瓶颈一般不在调度中心而在执行器线程池。建议的规划比例是一个执行器实例最多注册 200 个任务如果超过 200 个就拆成多个执行器应用。一个调度中心最多支撑 20 个执行器应用在线超过之后建议做多调度中心分片按业务线拆成多个独立集群。这样出问题时影响面能控制在一个集群内排查也方便。我个人在实际操作中的体会是调度系统最怕的不是负载高而是你把所有鸡蛋放在一个篮子里。宁可初期拆成两个小集群也不要等到一个集群里几百个任务互相影响时再迁移。另外最后分享一个小技巧给每个任务设置一个“SLA 窗口”如果任务在这个窗口内没有被触发就立即告警这比观察失败率更能提前发现调度链路的问题。把调度当第一公民来看待而不是最后一环系统稳定性的提升是肉眼可见的。
返回列表