为什么单体定时任务无法满足业务需求
在业务规模较小的阶段,使用node-cron或node-schedule在应用进程内跑定时任务是最常见的做法。比如每天凌晨清理过期数据、每小时汇总统计数据,几十行代码就能搞定。但当应用从单机扩展到多节点集群后,问题立刻暴露出来:假设服务部署了三个实例,每个实例的定时器都会独立触发同一个任务,数据清理被执行三次,统计任务产生三份重复数据,这在涉及金额或库存的场景下是致命的。
除了重复执行,单体定时任务还有几个明显的短板。其一是无法水平扩展,单个任务的执行能力受限于单个节点的CPU和内存,任务量大时只能干等。其二是没有故障转移能力,承载任务的节点宕机后任务就永久停止,需要人工介入。其三是缺乏统一的任务管理界面,修改执行周期必须改代码发版,运维成本很高。这些问题正是分布式任务调度系统要解决的核心诉求。
一个完整的分布式调度系统通常拆分为三个角色:调度中心负责任务的注册、触发时机的计算和任务的分发;执行器集群负责实际执行任务逻辑并上报结果;存储层负责持久化任务定义、执行日志和分布式锁。Node.js凭借事件驱动和非阻塞IO的特性,非常适合承担调度中心和执行器的角色,下面我们逐一实现各个模块。

调度中心设计与分布式锁防重复执行
调度中心的核心是一个秒级轮询器,它每一秒扫描一次任务表,找出当前时间点需要触发的任务。判断依据通常是任务的cron表达式,可以用cron-parser库计算出下一次执行时间。为了让多节点部署的调度中心不重复触发任务,每次扫描前必须先获取分布式锁,Redis的SET key value NX PX ttl命令是最简单可靠的实现方式。
下面是调度中心核心循环的简化实现:
const cronParser = require('cron-parser');
const Redis = require('ioredis');
const redis = new Redis();
// 调度中心主循环,每秒执行一次
async function scheduleLoop() {
setInterval(async () => {
// 尝试获取本秒的调度锁,只有抢到锁的节点才触发任务
const lockKey = `sched:lock:${Date.now() / 1000 | 0}`;
const locked = await redis.set(lockKey, process.pid, 'NX', 'PX', 900);
if (!locked) return;
const tasks = await loadDueTasks();
for (const task of tasks) {
// 任务级别再加一把锁,防止任务表被多份数据源重复加载
const taskLock = `task:lock:${task.id}:${task.nextRunTime}`;
const ok = await redis.set(taskLock, process.pid, 'NX', 'PX', 60000);
if (!ok) continue;
// 将任务推入执行队列,由空闲执行器消费
await redis.lpush('task:queue', JSON.stringify(task));
// 更新下一次执行时间
const next = cronParser.parseExpression(task.cron).next().getTime();
await updateNextRunTime(task.id, next);
}
}, 1000);
}这里有个容易被忽略的细节:调度锁的过期时间必须大于一轮扫描的最大耗时,否则第一个节点还没处理完,锁就过期了,另一个节点会再次抢到锁并重复分发任务。更稳妥的做法是引入锁续期机制,即在任务处理期间由后台协程定期延长锁的过期时间,Redlock算法就是这个思路的加强版。另外,任务级锁中拼接了nextRunTime,保证同一个任务在不同时间点的多次触发互不干扰。
执行器实现与任务队列消费模型
执行器的职责是监听任务队列,取出任务后调用对应的处理函数。Node.js是单线程模型,如果任务里有CPU密集型计算,会阻塞事件循环导致心跳上报延迟,进而被调度中心误判为宕机。解决方案是把重计算任务丢给worker_threads线程池执行,主线程只负责调度和IO。下面是一个支持并发的执行器骨架:
const { Worker } = require('worker_threads');
class TaskExecutor {
constructor(concurrency = 5) {
this.concurrency = concurrency;
this.running = 0;
}
start() {
setInterval(() => this.consume(), 100);
}
async consume() {
if (this.running >= this.concurrency) return;
const raw = await redis.rpop('task:queue');
if (!raw) return;
this.running++;
const task = JSON.parse(raw);
try {
// CPU密集型任务放入worker线程,避免阻塞主线程
await this.runInWorker(task);
await redis.lpush('task:result', JSON.stringify({
taskId: task.id, status: 'success', finishAt: Date.now()
}));
} catch (err) {
// 失败任务进入重试队列,附带重试次数
const retry = (task.retryCount || 0) + 1;
if (retry <= 3) {
task.retryCount = retry;
await redis.lpush('task:retry', JSON.stringify(task));
} else {
await this.alert(task, err.message);
}
} finally {
this.running--;
}
}
runInWorker(task) {
return new Promise((resolve, reject) => {
const worker = new Worker('./task-worker.js', { workerData: task });
worker.on('message', resolve);
worker.on('error', reject);
});
}
}消费模型的选择也值得讨论。LPUSH配合RPOP构成了简单可靠的消息队列,但如果执行器在RPOP取出任务后、执行完成前突然崩溃,这个任务就永久丢失了。改进方案是用BRPOPLPUSH把任务同时写入一个备份队列,执行成功后再从备份队列删除,启动时先检查备份队列中是否有残留任务并重新入队。这本质上就是Redis官方文档里描述的可靠队列模式,任务量不大时完全可以胜任,任务量达到每秒数千级别再考虑引入RocketMQ或Kafka。
心跳检测、故障转移与可视化管理
执行器集群的健康状态需要调度中心实时感知。每个执行器每隔五秒向Redis写入一条心跳记录,键名为执行器ID,过期时间设为十五秒。调度中心定期扫描所有心跳键,发现某个执行器的心跳消失,就把它标记为离线,并检查它是否有未完成的任务需要重新分发。这个机制保证了单个节点宕机后,它负责的任务会在半分钟内被其他节点接管,实现准实时的故障转移。
// 执行器侧:定时上报心跳
setInterval(async () => {
await redis.set(`node:alive:${NODE_ID}`, Date.now(), 'EX', 15);
}, 5000);
// 调度中心侧:检测离线节点并重新分发其任务
async function checkHealth() {
const nodes = await redis.keys('node:alive:*');
const onlineIds = nodes.map(k => k.split(':')[2]);
const offline = await getNodesNotIn(onlineIds);
for (const node of offline) {
const tasks = await getRunningTasksOf(node);
for (const t of tasks) {
await redis.lpush('task:queue', JSON.stringify(t));
}
await markNodeOffline(node);
}
}最后谈谈可视化管理。调度中心应该暴露一组HTTP接口,配合简单的管理页面实现任务的增删改查、手动触发、暂停恢复和执行日志查询。任务定义建议持久化到MySQL而非只存在Redis里,Redis负责运行态,MySQL负责配置态,两者通过版本号同步。每次任务执行都写入一条日志记录,包含开始时间、结束时间、执行节点、耗时和输出,有了这些数据才能分析任务耗时趋势、发现性能劣化的任务。
总结一下,用Node.js搭建分布式调度系统的关键技术点包括:秒级轮询加分布式锁保证任务不被重复触发,队列加备份队列保证任务不丢失,worker线程隔离CPU密集型逻辑,心跳机制实现故障转移,再加上配置持久化和执行日志形成闭环。这套方案在中等规模下运行稳定,代码总量不到一千行,非常适合作为理解分布式调度原理的实践项目,后续如果要对接更复杂的场景,可以参考XXL-JOB的设计思想逐步演进。