☰
自研分布式定时任务框架:负载均衡与OpenAPI异步调度实现
2026/10/9 3:49:27 网站建设 项目流程

分布式定时任务框架、负载均衡、OpenAPI异步调用,这三个词叠在一起,很多人第一反应就是“直接用 XXL-Job 不就行了”。但如果你遇到的是这样的场景:外部业务方要通过开放接口动态提交定时任务,平台负责把任务分发给不同的执行节点,并且对外只能走异步语义,不能长时间占用 HTTP 连接,那你大概率会发现现成框架反而成了瓶颈。我就在这种背景下从零写了一套轻量的分布式定时任务框架,里面包含调度中心、执行器、Redis 分布式锁、等开销负载均衡策略,以及一整套 OpenAPI 异步创建任务和查询状态的实现。这篇文章把完整的思路、核心代码和踩坑记录整理出来,适合已经用过 Spring Boot、Redis、MySQL,想自己动手搞明白分布式调度底层逻辑的开发者。

1. 为什么要自己写一个分布式定时任务框架

1.1 业务里被单机调度“卡”住的场景

公司的老系统里,定时任务一直是写在一个单体应用里的。最初只有几个报表任务,用 Spring 的@Scheduled加一个线程池就能跑。但是随着业务方变多,任务数量从十几个涨到几百个,这时候单机调度的问题就非常明显了:任务多了线程池排队,一个任务执行时间过长会把后面的任务全部拖死;服务重启或者宕机,所有定时任务跟着一起没了;线上部署多台实例时,同一个@Scheduled又会在每个实例上各执行一遍,报表重复生成、短信重复发送。

后来尝试过引入成熟框架,但发现两个问题。第一,这些框架本身是面向“运维配置任务”设计的,界面和调度能力很强,可我们的业务方希望直接通过 API 创建定时任务,要求提交后立刻返回,而任务到底执行成不成功,由平台通过回调或者查询接口再告诉他们。这套“异步 OpenAPI”语义,在现成框架里大部分是缺少的,需要大量二次开发。第二,我们并不需要极端的调度能力,更核心的是要把任务均匀地分发到不同执行节点,以及保证同一个任务在集群里不会同时跑多遍。与其被框架的约定绑住,不如自己写一个贴合业务的小框架。

1.2 自研框架的三大核心诉求

动手之前,我把需求收敛成三件事。

第一,任务调度必须是分布式的。多台调度实例同时在线,但同一时刻一个任务只能被一个调度器领走。实现上不能靠“随机抢”,要用分布式锁保证调度动作只有一个执行者。

第二,任务要能被负载均衡地分发到多个执行器。执行器节点注册上来之后,调度器要根据每个节点当前的“开销”而不是简单轮询来选择目标。这里我实现了一个等开销负载均衡策略,综合活跃任务数、历史响应耗时、CPU 空闲率来算权重,让每个节点的实际负载尽量接近。

第三,外部调用必须异步化。OpenAPI 入口接收到创建任务请求后,立刻落库并返回任务 ID,真正执行放在调度链路里。调用方用任务 ID 查询状态,或者等平台执行完成后回调通知。异步的好处很直接——调用方不用为一个可能执行几分钟的任务一直阻塞连接,调度平台也不会因为大量慢任务被打垮。

2. 整体架构和关键设计:调度中心、执行器、分布式锁怎么配合

2.1 模块划分和一次完整调度流程

整个框架从物理上分成两类进程:调度中心(Dispatcher)和执行器(Worker)。调度中心可以部署多台,执行器也可以部署多台,调度中心通过心跳判断执行器是否存活。存储层是 MySQL 加 Redis,MySQL 持久化任务定义和每次执行日志,Redis 用来做分布式锁、执行器心跳和临时状态缓存。

一次完整流程是这样走的:

  1. 业务方调用 OpenAPI 创建任务,接口写入job_info表,返回jobId。
  2. 调度中心每隔一秒去 MySQL 扫描next_trigger_time <= now且状态为等待中的任务。
  3. 调度中心先抢 Redis 锁,抢到锁的实例才有资格执行本次扫描和分发,避免多实例重复调度。
  4. 调度中心根据负载均衡策略从存活的 worker 列表里选一个节点。
  5. 调度中心把任务标记为“已分发”,并调用 worker 的/run接口。
  6. worker 接到任务后丢进本地线程池执行,执行结束后把结果回报给调度中心,写入job_log。
  7. 调用方通过 OpenAPI 查询任务状态,或者接收执行完成后的回调。

这套设计里,最关键的就是第 4 步和第 5 步中间的“状态流转”。如果不小心,两个调度实例会同时发现同一个到期任务,同时去调用 worker,任务就重复执行了。所以我在分发之前加了两道保险:Redis 分布式锁负责“调度动作互斥”,数据库乐观锁负责“任务状态原子变更”。

2.2 Redis 分布式锁:调度防重和任务分发的基石

分布式锁在分布式定时任务里扮演的角色,先想清楚一个问题:为什么不能只用数据库的唯一索引?因为调度中心只是“发现任务”和“分发任务”,任务真正的执行在 worker 侧,数据库层面只能保证任务不被重复领取,但不能保证两个调度中心不会同时把同一个任务发给两个不同的 worker。所以调度动作本身必须互斥。

我的第一版实现很简单:

@Scheduled(fixedDelay = 1000) public void scheduleDispatch() { // 先尝试加锁,如果加不上说明有别的调度实例在干活 Boolean locked = redisTemplate.opsForValue() .setIfAbsent("lock:dispatch", localIp, Duration.ofSeconds(10)); if (Boolean.FALSE.equals(locked)) { return; } try { doDispatch(); } finally { redisTemplate.delete("lock:dispatch"); } }

setIfAbsent对应 Redis 的SET NX,设置过期时间对应PX。这样即使某个调度实例在干活中途宕机,锁也会在 10 秒后自动释放,其他实例可以接管。

但这样写有一个隐患:如果doDispatch()执行时间超过 10 秒,锁会提前过期,另一个调度实例就会冲进来,两个实例同时分发。我后来在锁里加了一个“续期线程”,每 3 秒检查一次锁是否还在自己手里,是的话就把过期时间重置为 10 秒。这个逻辑和 Redisson 的看门狗机制是一致的。生产上可以直接用 Redisson 的RLock,但为了把原理讲透,我建议自己实现一遍。

2.3 负载均衡策略:从轮询到等开销负载均衡

执行器节点多了以后,负载均不均匀直接决定任务能不能及时跑完。最朴素的轮询策略在任务耗时差不多的情况下是够用的,可一旦出现某个节点正在跑一个跑半个小时的报表任务,另一个节点全都空闲,轮询还是会把新任务继续塞给繁忙节点,导致整体吞吐直线下降。

所以我实现了“等开销负载均衡”。这个“等开销”是我自己的叫法,意思是经过负载均衡之后,每个节点的综合开销水平应当尽量相等,而不是简单地任务数量一样多。实现上给每个 worker 算一个综合负载值:

  • 当前活跃任务数activeCount,这个值从 worker 心跳里带上来的。
  • 最近 5 分钟任务的平均耗时avgCostTime。
  • 节点上报的 CPU 空闲率cpuIdlePercent。

得到负载指数:load = activeCount * 10 + avgCostTime / 1000 + (100 - cpuIdlePercent) / 10。调度时选load最小的节点。这样做的好处是,一个节点就算当前任务数少,但如果它上报的 CPU 已经跑得很满,也不会被优先选中。

这里有个细节:worker 上报指标是有延迟的,调度中心拿到的可能是几秒前的状态。所以等开销负载均衡在任务并发不高的时候非常稳定,但瞬时大量任务涌入时还是可能倾斜。为了弥补,我额外加了并发控制,worker 本地线程池的队列长度一旦超过阈值,会主动拒绝任务并要求调度中心换一个节点。

3. 从零到一实现:调度中心、执行器、OpenAPI异步调用

3.1 存储设计与落库 SQL

一套分布式定时任务框架,最需要持久化的是三块内容:任务定义、执行日志、worker 节点信息。下面是核心表结构。

CREATE TABLE `job_info` ( `id` bigint NOT NULL AUTO_INCREMENT, `job_name` varchar(128) NOT NULL, `job_handler` varchar(255) NOT NULL COMMENT '任务处理器的bean名称', `params` text COMMENT '业务参数', `cron_expr` varchar(64) DEFAULT NULL COMMENT 'cron表达式,为空表示一次性任务', `next_trigger_time` bigint DEFAULT NULL COMMENT '下一次触发时间戳(ms)', `status` tinyint NOT NULL DEFAULT '0' COMMENT '0等待中 1运行中 2成功 3失败 4已暂停', `owner_worker_id` varchar(64) DEFAULT NULL COMMENT '当前负责执行的worker', `version` int NOT NULL DEFAULT '0' COMMENT '乐观锁版本号', `callback_url` varchar(512) DEFAULT NULL COMMENT '执行完成后的回调地址', `create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (`id`), KEY `idx_status_next_time` (`status`, `next_trigger_time`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; CREATE TABLE `job_log` ( `id` bigint NOT NULL AUTO_INCREMENT, `job_id` bigint NOT NULL, `worker_id` varchar(64) NOT NULL, `status` tinyint NOT NULL COMMENT '执行状态: 0执行中 1成功 2失败', `message` text, `start_time` bigint NOT NULL, `end_time` bigint DEFAULT NULL, PRIMARY KEY (`id`), KEY `idx_job_id` (`job_id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; CREATE TABLE `worker_node` ( `id` varchar(64) NOT NULL COMMENT 'worker唯一标识', `host` varchar(64) NOT NULL, `port` int NOT NULL, `active_count` int DEFAULT '0', `cpu_idle_percent` int DEFAULT '100', `last_heartbeat_time` bigint NOT NULL, PRIMARY KEY (`id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

job_info里的version字段很关键。调度中心分发任务时,不是直接更新状态,而是执行类似UPDATE job_info SET status=1, version=version+1 WHERE id=? AND status=0 AND version=?的语句。如果影响行数为 0,说明任务已经被别的调度器领走了,当前调度器只能放弃。这个乐观锁配合 Redis 锁,是双层防护。

3.2 调度主循环与分布式锁落地

调度中心的主循环我用 Spring Boot 的@Scheduled来触发,每秒跑一次。每次只处理一批最早到期的任务,避免一次扫描把数据库打满。

核心流程:

public void doDispatch() { // 扫描未来10秒内需要执行的任务,避免一条条绑时间点 long now = System.currentTimeMillis(); List<JobInfo> dueJobs = jobInfoMapper.findJobsByStatusAndTime(0, now, now + 10000); for (JobInfo job : dueJobs) { distributeJob(job); } // 把已经到期的cron任务计算下次触发时间 List<JobInfo> cronJobs = jobInfoMapper.findCronJobs(); for (JobInfo job : cronJobs) { if (job.getNextTriggerTime() <= now) { updateNextTriggerTimeAndDispatch(job); } } }

distributeJob里要做的事情比较多。

public boolean distributeJob(JobInfo job) { // 乐观锁占用任务状态 int rows = jobInfoMapper.tryClaim(job.getId(), localIp, job.getVersion()); if (rows == 0) { return false; } // 选择负载最低的 worker WorkerNode worker = loadBalancer.choose(workerRegistry.getAliveWorkers()); if (worker == null) { // 没有可用worker,把任务重新置回等待中,留待下次调度 jobInfoMapper.resetStatus(job.getId()); return false; } try { restTemplate.postForEntity("http://" + worker.getHost() + ":" + worker.getPort() + "/run", new RunRequest(job.getId(), job.getJobHandler(), job.getParams()), String.class); return true; } catch (Exception e) { // 调用失败,释放任务状态,下一次调度会再试 jobInfoMapper.resetStatus(job.getId()); return false; } }

tryClaim的 SQL 是这样的:

UPDATE job_info SET status = 1, owner_worker_id = #{workerId}, version = version + 1 WHERE id = #{id} AND status = 0 AND version = #{version}

这里有个细节我一开始忽略了:调用 worker 的/run是一个网络请求,请求发出去了不代表 worker 真的开始执行。假如 worker 收到了请求但还没来得及更新自己的状态,然后宕机了,调度中心这边可能认为分发失败并把任务重新置回等待中。这就有可能出现两次执行。要彻底解决这个问题,还得引入任务执行状态确认机制:worker 收到任务后先落一条“已接收”日志,调度中心收到确认后不再处理。不过这会增加复杂度,在简单场景里,我选择了“最多超时重试一次”的策略,配合 worker 端的 Redis 幂等标识来兜底,后文会详细说明。

3.3 执行器实现与心跳上报

执行器本质上是一个独立的 Spring Boot 服务,暴露两个 HTTP 接口:/run接收任务,/heartbeat上报状态。

/run接口的核心逻辑:

@PostMapping("/run") public ResponseEntity<String> run(@RequestBody RunRequest request) { // 幂等校验,同一jobId同一时间戳只允许执行一次 String idempotentKey = request.getJobId() + ":" + request.getTriggerTimestamp(); Boolean isNew = stringRedisTemplate.opsForValue() .setIfAbsent("idempotent:" + idempotentKey, "1", Duration.ofHours(1)); if (!Boolean.TRUE.equals(isNew)) { return ResponseEntity.ok("duplicated"); } // 交给本地线程池异步执行 taskExecutor.execute(() -> { JobLog log = new JobLog(); log.setJobId(request.getJobId()); log.setWorkerId(workerId); log.setStatus(0); log.setStartTime(System.currentTimeMillis()); jobLogMapper.insert(log); try { // 通过Spring上下文获取对应的Handler JobHandler handler = context.getBean(request.getJobHandler(), JobHandler.class); handler.handle(request.getParams()); jobLogMapper.updateStatus(log.getId(), 1); // 执行成功,回调通知 notifyCallback(request, true, null); } catch (Exception e) { jobLogMapper.updateStatus(log.getId(), 2, e.getMessage()); notifyCallback(request, false, e.getMessage()); } }); return ResponseEntity.ok("accepted"); }

心跳上报是另一个定时任务,每 5 秒一次。上报内容包含当前活跃任务数、CPU 空闲率、最近任务平均耗时。调度中心收到心跳后更新worker_node表,同时判断如果超过 15 秒没收到心跳,就把这个 worker 标记为离线,调度任务时自动跳过。

3.4 负载均衡算法代码实现

等开销负载均衡的核心类我写成这样:

@Component public class LoadBalancer { private final WorkerRegistry registry; public WorkerNode choose() { List<WorkerNode> alive = registry.getAliveWorkers(); if (alive.isEmpty()) { return null; } return alive.stream() .min(Comparator.comparingDouble(this::calcLoad)) .orElse(null); } private double calcLoad(WorkerNode worker) { // 活跃任务数权重最高 double activeScore = worker.getActiveCount() * 10; // 平均耗时,单位毫秒,换算成加权分 double avgCostScore = worker.getAvgCostTime() / 1000.0; // CPU空闲率越低,得分越高 double cpuScore = (100 - worker.getCpuIdlePercent()) / 10.0; return activeScore + avgCostScore + cpuScore; } }

在实际测试里,这个策略起效的关键在于 activeCount 的准确性。worker 执行任务前要把计数器加一,任务结束后减一,同时把本次耗时记录到滑动窗口。如果少维护了计数器,负载均衡就退化成随机分发。

除了等开销策略,我还保留了轮询和随机策略。通过配置中心切换策略,方便对比效果。如果你不想自己维护节点心跳,也可以用 Redis 的ZSet按心跳时间戳排序来管理存活节点,调度前过滤掉过期节点,省一张数据库表,代价是查询不够灵活。

3.5 OpenAPI异步创建任务与状态查询

OpenAPI 的设计核心是“异步”。调用方创建一个任务,正常情况下不应该等待任务执行完成,而是立刻得到一个任务 ID。然后调用方循环查询状态,或者由平台在执行完成时回调调用方。

我的 OpenAPI 接口长这样:

@RestController @RequestMapping("/openapi") public class OpenApiController { @PostMapping("/jobs") public ResponseEntity<CreateJobResponse> createJob(@RequestBody CreateJobRequest request) { // 幂等处理,同一个clientRequestId不重复创建 String idempotentKey = "openapi:job:" + request.getClientRequestId(); Boolean isNew = stringRedisTemplate.opsForValue() .setIfAbsent(idempotentKey, "1", Duration.ofDays(1)); if (!Boolean.TRUE.equals(isNew)) { JobInfo existing = jobInfoMapper.findByClientRequestId(request.getClientRequestId()); return ResponseEntity.ok(new CreateJobResponse(existing.getId(), existing.getStatus())); } // 先落库,再返回 JobInfo job = new JobInfo(); job.setJobName(request.getJobName()); job.setJobHandler(request.getJobHandler()); job.setParams(request.getParams()); job.setCronExpr(request.getCronExpr()); job.setNextTriggerTime(calculateNextTriggerTime(request)); job.setCallbackUrl(request.getCallbackUrl()); job.setStatus(0); jobInfoMapper.insert(job); return ResponseEntity.ok(new CreateJobResponse(job.getId(), job.getStatus())); } @GetMapping("/jobs/{id}") public ResponseEntity<JobStatusResponse> getJobStatus(@PathVariable("id") Long id) { JobInfo job = jobInfoMapper.findById(id); return ResponseEntity.ok(new JobStatusResponse(job.getId(), job.getStatus(), job.getParams())); } }

创建接口的时序很重要:必须先落库,再返回一个“已接收”的响应。如果先返回,调用方立刻查状态,会查不到这条记录。这里可以做成事务,保证插入和返回之间没有中间状态。

有些业务方希望任务完成后平台主动通知他们,这时候可以用callbackUrl。在 worker 执行完任务后,由调度中心回调外部系统的接口,回调体里带上任务 ID、执行状态、错误信息。要注意的是,回调接口必须是幂等的,重试时不能造成业务重复。我在回调请求里带了一个executionId(每次执行生成一个 UUID),回调接收方可以根据这个字段去重。

3.6 异步任务如何做到可追踪

很多人在自研时忽略一个点:异步任务虽然返回快,但调用方怎么知道任务到底跑没跑?如果只提供一个“查状态”的接口,调用方就得不停轮询,体验很差。我采用的是“状态流转+变更记录”的方式:job_info表里的状态从 0 到 1 到 2 或 3,每一步都在job_log里写一条记录。调用方查询时,不仅可以拿到任务当前状态,还能拿到每一步的执行时间和错误信息。

对于长时间不回调的任务,我在调度中心加了一个“超时扫描”。比如某个一次性任务,分发出去后 10 分钟还没从status=1变成 2 或 3,就视为执行异常,把状态重置为 0 并重新调度。这个机制既保护了外部调用方体验,也避免任务因为 worker 异常卡死。

4. 上线以来踩过的坑:重复调度、锁超时、任务丢失

4.1 重复执行问题排查与幂等设计

第一版上线后,最严重的问题就是任务偶尔重复执行。不是每一条都重复,而是集中在任务执行时间比较长、调度中心刚好发生集群扩容的时候。排查了很久发现,问题出在 Redis 锁上:调度实例 A 抢到锁,正在逐个分发任务,某一次分发时调用 worker 超时,A 在异常处理时把锁释放了;这时候调度实例 B 抢到锁,把 A 已经分发过的任务又分发了一遍。worker 端收到两个一模一样的执行请求。

光靠 Redis 锁根本防不住这种“跨调度周期”的重复。所以我在 worker 端加了幂等,对同一任务同一次触发只执行一次,用jobId + 触发时间戳作为 Redis 幂等键。但这里也有坑,如果 worker 在执行完任务之后、更新状态之前重启了,Redis 里的幂等键还在,任务就不会被重新执行,造成任务丢失。因此我又在调度中心加了一个超时重置逻辑,超过一定时间没收到执行结果,删除幂等键并重新分发。

最终总结下来,双重幂等是最稳的:调度中心分发前用数据库乐观锁保证状态流转,worker 执行前用 Redis 幂等键保证单次执行不被重复。两者缺一不可。

4.2 worker 宕机后的任务拉起

有一次我做了个破坏性测试:正在执行一个耗时 5 分钟的任务,直接 kill 掉 worker。任务卡在status=1状态,调度中心没有感知,因为心跳数据是每 5 秒上报一次,worker 没了之后调度中心要最多等 15 秒才能判断它离线。

离线判断可以做,但卡在运行中的任务怎么恢复?我在调度中心加了一个“孤儿任务扫描”:

SELECT * FROM job_info WHERE status = 1 AND owner_worker_id IS NOT NULL AND NOT EXISTS ( SELECT 1 FROM worker_node WHERE worker_node.id = job_info.owner_worker_id AND worker_node.last_heartbeat_time > #{now - 15000} )

把这些任务找出来后,把状态重置为 0,并清掉 worker 端的幂等键。这里有一个非常容易被忽略的细节:重置任务前,要确认旧 worker 确实已经死透了。如果旧 worker 只是网络抖动,任务还在它上面执行,你把状态重置了,另一个 worker 也会执行同一个任务,两边同时操作业务数据,后果不堪设想。因此我坚持用“心跳超时 + TCP 探测”双重判断,只有 TCP 连接也连不上,才认为 worker 已经离线。

4.3 Redis 锁超时和续约问题

前面提到过,分布式锁如果过期时间设短了,任务还没分完锁就自动释放,会造成并发调度;如果设长了,调度实例宕机后其他实例要等很久才能接管。我最初设了 10 秒,但是想象一下这个场景:数据库里有几千个到期任务,调度中心在逐条分发,每一条都要走一次网络请求,10 秒根本不够。于是我用“续期线程”定时给 Redis 锁续期,相当于实现了一个简易看门狗。

续期线程实现:

ScheduledExecutorService renewThread = Executors.newSingleThreadScheduledExecutor(r -> { Thread t = new Thread(r, "lock-renew"); t.setDaemon(true); return t; }); renewThread.scheduleAtFixedRate(() -> { String currentValue = redisTemplate.opsForValue().get("lock:dispatch"); if (localIp.equals(currentValue)) { redisTemplate.expire("lock:dispatch", Duration.ofSeconds(10)); } }, 0, 3, TimeUnit.SECONDS);

但注意,如果调度实例长时间持有锁不释放,其他实例会一直空转。所以我在finally块里释放锁之前,还会判断锁的 value 是否还是自己的 IP,防止把别人的锁误删。这个细节建议直接参照 Redisson 的源码实现,自己造的轮子稳定性确实不如成熟库。

4.4 OpenAPI 异步接口的事务与幂等坑

OpenAPI 异步接口最容易出的问题有两个。

一个是事务问题。我一开始创建任务和计算下次触发时间是在同一个事务里,但因为业务方通过 OpenAPI 创建任务的频率很高,数据库连接经常不够用,偶尔会出现任务插入后事务没提交,线程池已经调度到了这条数据的情况。这会导致调度器拿到的是“半条”数据,比如jobHandler字段还是 null。解决方法是把“创建任务”和“任务调度”彻底解耦,插入完成后直接提交事务,调度线程下一次扫描时自然能看到完整数据。

另一个是重复创建。调用方可能因为网络重试,把同一个创建请求发两次。我在接口入口用clientRequestId做幂等:第一次请求把clientRequestId存入 Redis,后续相同请求直接返回已有任务。这里要注意 Redis 的setIfAbsent和业务落库必须保证最终数据一致。如果客户端发的是同一个请求但参数不同,幂等键只能拦截完全相同的情况,参数不同会直接报错。

5. 压测评估和后续扩展方向

5.1 小规模压测实测结果

整个框架在一个三节点的测试环境里跑了一轮,一个调度中心,两个 worker,配置都是 4C8G。准备 1000 个一次性任务,全部通过 OpenAPI 批量创建,然后观察调度链路。实测结果是完整的任务吞吐在每分钟 2500 个左右,超过这个量之后 worker 线程池开始排队,接口响应变慢但不会丢任务。单任务从创建到执行的延迟基本在 1 秒到 2 秒之间,因为调度主循环是每秒扫一次,所以这个延迟是可预期的。

压测中我发现job_info表的状态索引特别重要。如果没有idx_status_next_time,扫描到期任务会导致全表扫描,任务量超过万级以后调度延迟会直线上升。加了这个组合索引后,扫描基本稳定在 10 毫秒内。

5.2 后续可以怎么扩展

这套框架目前最大的短板是缺少停机维护时的优雅处理。worker 在接受新任务前,可以先把自己的状态置为“维护中”,让调度中心不再分配新任务,等本地线程池跑完再真正停机。另外,失败重试目前只做了简单的“重置状态重新调度”,对于重试次数和重试间隔还没做配置化。后续计划引入几种策略增强机制:一是适配不同优先级任务,让高优任务可以抢占低优任务;二是接入 Prometheus 指标,把调度延迟、任务成功率、worker 负载情况可视化;三是把 OpenAPI 的异步回调做成可重试的消息队列,避免回调 HTTP 接口临时抖动导致通知丢失。

另一个我强烈建议的扩展方向是执行器端的优雅停机。demo 版本里我把任务丢进线程池后,进程如果杀掉了,执行中的任务直接丢失。可以在项目里维护一个“正在执行任务集合”,在收到停机信号时拒绝新任务,并等待已有任务超时或者上报中断状态。

最后分享一个实际运维中的小技巧:worker 和调度中心的时间一定要做 NTP 同步。调度算法里大量时间戳用的是毫秒,如果机器之间时钟偏差太大,就可能出现任务刚创建出来就立刻触发,或者调度中心判断 worker 心跳超时的误报。别小看这个细节,我第一版上线时唯一一次线上事故就是时钟偏差导致的。

分布式定时任务框架自己完整实现一遍,最大的收获不是代码能跑,而是彻底理解了调度、锁、负载均衡这三者在分布式环境下的协作方式。接下来再碰到任何调度类的业务需求,第一反应就不再是盲目堆中间件,而是能根据业务流量和可靠性要求,随时画出一套适合自己的方案。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询