ARTICLE DETAIL

资讯详情

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

从零实现任务调度系统:状态机、超时重试与延迟队列实战

从零实现任务调度系统:状态机、超时重试与延迟队列实战 TJXT这个项目我原计划五天搞定第一版结果到第五天晚上发现还是乐观了。Day6也就是今天几乎全部花在状态机收尾、超时重试和联调上面。如果你也在做一个带调度的后台系统这天的过程应该很有参考价值。先简单交代一下TJXT是什么。名字是我偷懒起的全称 Task Job eXecution Tracker翻译过来就是“任务工作执行跟踪器”。它解决的问题很具体统一管理后台各种跑批脚本、数据同步、定时报表、脏数据修复任务。我不想给个人项目引入一套重型分布式调度框架但任务数量到了几百个之后你又没办法靠脑子记住哪个任务挂了、哪个任务该重试。TJXT就是用来干这个的——创建任务、设置调度规则、执行任务、记录每次运行结果、失败自动重试、发通知一条链路全包。今天是第六天整个系统已经能跑通“手动建任务 → 定时触发 → 执行结果回写 → 失败重试 → 消息通知”这条完整链路。适合谁看适合想自己写一个任务调度工具的人也适合正在做类似后台系统的后端或者全栈同学。下面把今天的具体实现和踩坑过程都拆开讲。1. 项目定位与这六天干了啥1.1 TJXT到底是个什么东西TJXT的核心设计目标其实只有三个字能自救。所谓“能自救”指的是任务挂了之后系统能自动感知、标记失败、按策略重试、最后通知到人。而不是像很多临时脚本一样日志打到一半进程崩了第二天才发现昨天报表没跑出来。从实现上看TJXT支持两类执行动作HTTP请求任务和本地Shell命令任务。HTTP任务适合调用别人接口或触发远程脚本Shell任务适合在本机执行数据同步或定时清理。这样一个工具基本能覆盖个人服务器和小团队内部的大部分定时需求。技术栈上我没有追新。服务端用Node.js TypeScript Fastify数据库用PostgreSQL缓存和延迟队列用Redis前端是一个简单的管理控制台页面用来手动触发任务和查看运行记录。这套组合的好处是生态成熟、开发效率高对于TJXT这种规模的项目完全够用。Day6之前TJXT还只算一个半成品。今天做完之后我才敢说它勉强能挂到线上长时间跑。1.2 前五天的进展清单先把前五天的进度拉一遍方便你理解第六天为什么紧急。Day1做需求梳理和数据模型初稿。核心表定了三张tasks任务定义、task_runs任务执行记录、notify_configs通知配置。顺便把数据库迁移脚本的结构想清楚了。Day2搭好Fastify TypeScript项目骨架写完所有数据库迁移包括基础表、索引、外键。Day3实现用户登录和任务CRUD接口。任务字段包含任务名、类型、Cron表达式、超时时间、重试次数、通知渠道配置。Day4写好任务执行器。能通过HTTP和Shell两种方式跑任务把标准输出、退出码、耗时写回task_runs。Day5补了定时触发器和最基础的仪表盘页面。但整个状态机只做了“创建、运行、完成”三种状态超时任务会一直卡在RUNNING重试机制也没有。从Day5晚上开始我就意识到第六天肯定要加班了。1.3 为什么第六天是分水岭前五天做出来的东西说白了是一个“看起来能跑”的Demo。你手动点一下任务它能执行能输出结果。但一旦任务卡住、进程崩溃、网络超时、重试次数用完整个系统就瘫痪了。第六天把任务状态从三种扩展成六种并且把调度、超时、重试、通知全部串起来。这一步做完TJXT才从“玩具”变成“工具”。分水岭在于它挂了能自己恢复而不是靠人盯着日志手动处理。所以Day6的主要工作分了四块任务状态机重新设计并把所有状态流转落到代码里。实现基于延迟队列的超时检测和失败重试。把定时调度器和通知模块接入主线流程。修数据层和接口层暴露出来的各种并发问题。下面按顺序讲。2. 任务状态机从0到1的核心建模2.1 六个状态和关键流转规则状态机是整个系统的地基。如果状态设计不合理后面调度、重试、通知都会乱成一锅粥。TJXT最终定义了六种状态PENDING已创建、等待被调度执行。RUNNING正在执行中。SUCCESS执行成功。FAILED执行失败且不再重试。RETRY_WAIT执行失败等待退避后重试。CANCELLED任务被取消。为什么需要一个单独的RETRY_WAIT状态而不是直接在FAILED状态里记一个重试时间因为系统需要一个明确状态来区分“最终失败”和“暂时失败”。最终失败要发通知、要终态展示暂时失败只是中间态它还会被调度器再次拉起来。如果混在一起日志和页面都会很懵。关键流转规则我用一张表说明当前状态允许流转到触发条件PENDINGRUNNING、CANCELLED调度触发或手动执行取消任务RUNNINGSUCCESS、FAILED、RETRY_WAIT执行成功执行失败且可重试执行失败且还有重试次数RETRY_WAITRUNNING、CANCELLED退避时间到用户在等待期间取消任务SUCCESS终态无FAILED终态无CANCELLED终态无注意我没有给SUCCESS和FAILED设计“重置为PENDING”的流转入口。页面如果要重新执行一个终态任务正确做法是复制一份新任务或手动触发一个新run而不是去把旧状态改掉。这样能保证历史记录的不可变性和可审计性。2.2 状态流转的代码实现状态机的代码实现核心是一个流转映射表加一个校验函数type TaskStatus PENDING | RUNNING | SUCCESS | FAILED | RETRY_WAIT | CANCELLED; const TRANSITIONS: RecordTaskStatus, TaskStatus[] { PENDING: [RUNNING, CANCELLED], RUNNING: [SUCCESS, FAILED, RETRY_WAIT], RETRY_WAIT: [RUNNING, CANCELLED], SUCCESS: [], FAILED: [], CANCELLED: [], }; function canTransition(from: TaskStatus, to: TaskStatus): boolean { return TRANSITIONS[from].includes(to); }但仅有这个校验还不够。真正的坑在并发。如果两个请求同时把一个任务从PENDING改成RUNNING数据库层面的状态就错乱了。所以更新状态时必须在SQL层面做条件更新UPDATE tasks SET status RUNNING, updated_at now() WHERE id $1 AND status PENDING RETURNING id;这种写法相当于乐观锁。如果更新影响行数为0说明状态已经被别人改过这次流转直接失败。实测下来用这种方式比“先SELECT后UPDATE”更安全也省掉手动加锁的麻烦。2.3 状态机设计里的三个反模式今天踩了几个坑之后我总结出三个状态机设计时特别容易犯的问题。第一个反模式允许同一任务的多个实例并发RUNNING。有些场景你确实想并行跑但TJXT默认是单执行实例必须避免重复调度。所以我要求从PENDING只能进一个RUNNING并且数据库层用条件更新兜底绝不信任上层代码。第二个反模式用单个字段同时表达状态和子状态。比如用status列存“RUNNING_RETRY_3”这种值查询倒是方便但后续状态扩展和统计会很痛苦。正确的做法是把重试次数等变化信息放到task_runs表里任务主表只保留当前状态。第三个反模式任务成功后随意改回失败。由于网络抖动或回调乱序旧执行结果晚到就可能把新结果覆盖掉。我后面引入run_id来区分每次执行只有最新的run_id才能更新任务状态旧run的回调一律丢弃。这一点在后面的联调部分还会详细讲。3. 超时检测与重试机制今天的主角3.1 方案选型轮询扫表还是延迟队列状态机搞定之后剩下的重头戏是超时检测。一个任务启动后如果一直处于RUNNING状态总得有东西把它捞出来。最朴素的做法是写个定时任务每隔30秒扫一遍task_runs表把所有RUNNING且updated_at超过超时时间的记录找出来统一标记超时。优点是简单缺点也很明显扫描频率越高数据库压力越大扫描频率越低超时判定越迟钝。我最终选择了基于Redis ZSET的延迟队列方案。Redis的ZSET天然支持“按分数排序、按分数范围查询”把超时时间点当作score任务ID当作member就能实现一个精度到秒级的小型延迟队列。考虑到极端情况下Redis数据可能丢失我又保留了一个兜底扫表任务跑得比延迟队列慢得多比如5分钟一次只处理那些Redis里找不到的超时任务。双保险槽点在于多写一套代码但对于调度系统来说这种兜底是值得的。3.2 基于Redis ZSET的延迟队列实现ZSET的用法很直接。任务开始执行时把超时检查点写入Redisawait redis.zadd(task:timeouts, Date.now() task.timeoutMs, task.id);消费逻辑则是一个常驻的循环每秒跑一次取所有已经到期的任务IDconst now Date.now(); const ids await redis.zrangebyscore(task:timeouts, 0, now); for (const id of ids) { const task await findTask(id); if (task task.status RUNNING) { await handleTaskTimeout(task); } await redis.zrem(task:timeouts, id); }看着没什么问题但第一次实现我踩了一个大坑ZSET里面取出来之后任务可能已经被重试流程处理完毕并进入SUCCESS状态但超时检查成员还残留在Redis里。消费到它时又会把任务误判为超时。解决办法是消费端先查任务当前状态只有RUNNING状态才真正执行超时逻辑。另外早期为了简单我用zrangebyscore加zrem两步操作这在多实例部署时有概率重复消费。单个实例问题不大但如果有两个调度器进程就必须用Lua脚本把“取数据删除成员”合并成原子操作local ids redis.call(ZRANGEBYSCORE, KEYS[1], 0, ARGV[1], LIMIT, 0, 100) if #ids 0 then redis.call(ZREM, KEYS[1], unpack(ids)) end return ids让Redis脚本只把这些ids返回给应用层应用层拿到再逐个处理。这样即使有多个实例并发同一个任务也只会被一个实例抢到。3.3 重试退避算法与并发控制超时或失败之后系统不能立刻无脑重试。如果下游服务正在雪崩立刻重试只会加重问题。所以TJXT采用指数退避加随机抖动function calculateRetryDelay(retryCount: number): number { const base 30; // 初始30秒 const maxDelay 3600; // 最长1小时 const exponential base * Math.pow(2, retryCount); const jitter Math.floor(Math.random() * 5000); return Math.min(exponential jitter, maxDelay); }默认每个任务最多重试3次。第一次失败后等30秒第二次失败后等60秒第三次失败后等120秒最多不超过1小时。随机抖动是为了避免大量任务同一时刻集体重试造成“重试风暴”。并发控制这块我专门花了很多时间。多个worker进程同时拉取可重试任务时很容易把同一个任务拉出来执行两遍。我最终用PostgreSQL的行锁定语法来抢任务SELECT * FROM tasks WHERE status RETRY_WAIT AND next_retry_at now() ORDER BY next_retry_at FOR UPDATE SKIP LOCKED LIMIT 10;FOR UPDATE SKIP LOCKED非常关键。它能把正在被其他事务锁定的行直接跳过不会阻塞等待。这样多个worker可以同时扫表但每个任务只会被一个worker拿到。拿到任务后再把状态改成RUNNING释放锁。这个方案的缺点是依赖PostgreSQL换到MySQL就得另想招。但对TJXT这种绑死PostgreSQL的项目来说是最顺手也最不容易出错的选择。4. 调度器与通知模块联调4.1 定时任务触发链路状态机做完后我开始联调调度器。TJXT没有引入Quartz这类重型调度框架而是自己用cron解析器在进程内注册定时器。每个任务创建时解析一次Cron表达式到点后把任务投递到触发队列。触发链路大概是这样的应用启动时从tasks表读取所有启用状态的定时任务。对每个任务注册定时器定时器到点后调用enqueueTask(taskId)。enqueueTask会先查任务当前状态如果是PENDING或者RETRY_WAIT就把状态改为RUNNING并生成一条新的task_runs记录。执行器消费task_runs执行HTTP请求或Shell命令。结果回来后更新task_runs和任务主状态。这里要特别注意定时触发的任务不要拿到就直接跑尤其是同一秒可能有几十个任务同时到期。我给执行器加了一个信号量同时运行的协程数上限设为10。超出上限的任务先在内存队列里排队。实测下来任务量大时不会把服务器CPU打满。今天联调中发现一个特别隐蔽的问题Cron表达式解析时区错了。服务器系统时区是UTC但任务配置者期望的是东八区。结果表现为所有定时任务都比预期提前8小时触发。我查了半天才发现是cron解析库默认用了服务器本地时区。解决方法是解析时显式传入timeZone: Asia/Shanghai并且把时区字段暴露在任务配置页面上让用户明确知道自己在配置哪个时区。4.2 通知模块的抽象与模板任务重试次数用完、最终失败时系统要通知人。今天的第二个大块工作是通知模块。我一开始直接在状态流转代码里写sendWebhook、sendEmail后来发现这样扩展性太差。如果以后要加IM机器人和短信就得改主流程代码。所以我把通知抽成了一个Notifier接口interface Notifier { send(channel: NotifyChannel, payload: NotifyPayload): Promisevoid; }每个渠道实现一个Notifier比如WebhookNotifier、SmtpNotifier。任务配置里可以绑定多个渠道失败时逐一发送。通知内容我用了模板字符串没有引入重型模板引擎。模板大概是这样的任务「{{task.name}}」第 {{retryCount}} 次重试失败。 任务ID{{task.id}} 失败原因{{errorMsg}} 发生时间{{timestamp}}模板的好处是统一文案格式避免不同渠道出现不同风格的通知。另外通知发送不能阻塞主流程。我采用outbox模式在更新任务状态的事务里同时写一条notification_events记录后台有一个消费者负责把事件发出去。这样主流程和通知发送解耦即使通知渠道响应慢也不会拖累任务状态更新。4.3 联调现场记录联调过程中最经典的一个问题是重复通知。任务失败后状态更新和通知发送不在同一个事务里后台消费者可能重复读取同一条通知事件导致同一失败任务发了两遍邮件。解决办法就是唯一约束。在notification_events表里给task_run_id加唯一索引消费者插入事件时利用数据库唯一索引去重。如果插入失败说明事件已经存在直接忽略。另一个问题是outbox消费者删除事件和发送动作之间崩溃会导致事件重新发送。考虑到通知本身允许偶尔重复我选择了“发送后再删”而不是“先删后发送”。宁可重复通知也不能漏通知。这个取舍在任务调度场景里很重要。到傍晚的时候整个链路已经能跑通了。页面手动触发一个Shell任务让脚本故意exit 1系统会在30秒后自动进入RETRY_WAIT等下一次重试3次重试全部失败后任务进入FAILED同时给Webhook地址推一条失败告警。看到这条链路真正跑通的时候我才松了口气。5. 数据层和接口层那些破事5.1 索引、锁、事务配合调度系统对数据层的要求比普通CRUD高很多。今天主要处理了三类问题索引、锁、事务边界。索引方面最核心的查询是“找出所有等待重试的任务”和“找出所有RUNNING中但可能超时的任务”。我给task_runs表加了两个部分索引CREATE INDEX idx_task_runs_retry ON task_runs(task_id, status, next_retry_at) WHERE status IN (RETRY_WAIT, RUNNING);部分索引比普通索引小很多查询时的IO也少。第一次写的时候忘了加task_id导致关联查询每次都要回到主表回表性能差了不少。后来调整成现在的组合才把扫描时间降下来。事务边界方面状态切换和日志插入必须保证原子性。比如任务从RUNNING变成SUCCESS同时要在task_runs里写入执行结果这两件事必须在同一个事务里完成。否则会出现任务状态已经成功但执行记录缺失的情况。我是把“更新任务主表”和“更新task_runs”放在一个事务里通过transaction函数包裹。锁的使用上前面提到的FOR UPDATE SKIP LOCKED解决worker抢任务问题手动触发接口则用应用层分布式锁。对同一个任务手动触发和定时触发可能同时发生所以我加了锁Keytask:lock:{taskId}获取不到锁就直接拒绝避免双执行。5.2 接口幂等与幂等表设计接口层今天的重点是幂等。用户在前端按钮点了两次“立即执行”系统绝不能跑两遍同一个任务。我做了两层保险。第一层是Redis幂等键。手动触发接口要求客户端传一个幂等键const idempotencyKey run:${userId}:${taskId}:${day}; const ok await redis.set(idempotencyKey, 1, EX, 86400, NX); if (!ok) { throw new Error(任务已触发请勿重复提交); }用SETNX命令保证同一任务在同一天只能触发一次。注意这个粒度是“同一天”如果确实需要一天内多次手动执行这个键就不适用。但对于TJXT来说手动触发通常是补偿操作一天一次足够。第二层是task_runs表里的幂等约束。每个任务的每次执行记录会生成一个runId执行器回调时必须以runId作为唯一标识。在task_runs表加unique(task_id, run_id)约束即使回调重复送达数据库也会拒绝重复写入。5.3 排查问题实录从日志到复现今天排查最久的一个问题就是“任务重试成功但最终状态还是FAILED”。从现象看任务第一次失败后进入RETRY_WAIT第二次重试执行成功task_runs里已经有SUCCESS记录但tasks表状态仍然被改成FAILED。翻日志后发现一个典型的乱序覆盖问题第一次失败的异步回调在网络延迟后才到达此时重试已经成功执行任务的当前状态已经是SUCCESS。旧回调却不认识新状态直接按照“执行失败”把任务改成FAILED。复现场景很简单在任务执行器里人为加一个5秒延迟同时让HTTP回调先失败再在5秒后补发一个成功回调。两个回调的到达顺序一乱状态必然出错。我的修复方案就是之前说的run_id机制。每次执行时生成一个新的run_id只有当前task_runs表里最新一条run_id对应的回调才允许更新任务主状态。旧run_id的回调只记录不更新主表。这个方案在今天的测试里非常稳后面又补了一个自动化测试防止回归。6. 测试与部署实录6.1 自动化测试场景设计联调告一段落后我把今天的关键逻辑补了自动化测试。测试框架用的Vitest数据库用真实的PostgreSQL和Redis没有用mock数据库。原因很简单调度和锁的行为只有真实数据库才能暴露出来本地SQLite完全模拟不了SKIP LOCKED。重点覆盖了六个场景非法状态流转必须抛异常。合法状态流转成功后数据库状态正确变更。超时任务能进入RETRY_WAIT而不是直接FAILED。重试次数用完后任务进入FAILED并只发一次通知。并发触发同一任务只允许一个RUNNING。旧run_id回调不能更新新任务状态。每个场景都写成独立的测试文件。其中并发触发的测试写了两个并发请求同时POST到手动触发接口断言task_runs里最终只有一条记录。这个测试在修复幂等逻辑前必挂现在成了回归保障。6.2 部署配置与启动脚本晚上终于开始部署。TJXT用Docker Compose管理三个服务PostgreSQL、Redis、应用本体。docker-compose.yml的核心部分长这样version: 3.8 services: postgres: image: postgres:15-alpine environment: POSTGRES_DB: tjxt POSTGRES_USER: tjxt POSTGRES_PASSWORD: change-me volumes: - pgdata:/var/lib/postgresql/data healthcheck: test: [CMD-SHELL, pg_isready -U tjxt] interval: 5s timeout: 3s retries: 10 redis: image: redis:7-alpine command: [redis-server, --appendonly, yes] volumes: - redisdata:/data app: build: . depends_on: postgres: condition: service_healthy redis: condition: service_started environment: DATABASE_URL: postgres://tjxt:change-mepostgres:5432/tjxt REDIS_URL: redis://redis:6379 TZ: Asia/Shanghai CRON_TIMEZONE: Asia/Shanghai ports: - 8080:8080 volumes: pgdata: redisdata:有个小提示TZ环境变量和CRON_TIMEZONE要区分开。TZ影响容器日志时间CRON_TIMEZONE影响任务调度解析。如果只设置一个另一个地方还是会踩时区坑。初始化时先执行数据库迁移再启动应用npm run migrate up docker compose up -d --build6.3 常见问题速查表今天最后把所有遇到的问题整理成了一个表格后面再排查就直接查表。症状可能原因对策任务一直PENDING不执行cron时区没配对、调度器进程没启动检查CRON_TIMEZONE配置确认应用注册了定时器任务重复执行幂等键缺失、手动接口被双击、回调重复送达Redis SETNX加幂等键task_runs加唯一约束任务卡在RUNNING超时扫描没启动、Redis ZSET key丢失检查延迟队列消费进程兜底扫表任务兜住重试不生效重试次数配置为0、状态机不允许RETRY_WAIT流转检查任务配置确认状态流转函数覆盖RETRY_WAIT通知没收到outbox事件没被消费、Webhook URL配置错误查看notification_events表确认消费者在运行任务状态被旧回调覆盖缺少run_id机制、回调乱序给task_runs加run_id旧回调只记录不更新主状态定时任务提前8小时服务器UTC与东八区混淆cron解析显式传入timeZone最后再分享一个今天最大的体会。状态机这东西画图看起来特别简单无非几个框几条线但真正落地的时候并发、幂等、回调乱序、事务边界每一个细节都能让你加班到半夜。如果重新做一次我会先把整条执行链路的时序图画清楚把所有可能乱序到达的回调都列出来再动手写代码。另外别迷信重型框架。自己写一个最小实现虽然费劲但对整个系统的理解深度完全不是一个层次。
返回列表