深入 Dub 后台任务框架:用 defineJob 统一 QStash 负载驱动任务
【免费下载链接】dubThe modern link attribution platform. Loved by world-class marketing teams like Framer, Perplexity, Superhuman, Twilio, Buffer and more.项目地址: https://gitcode.com/GitHub_Trending/du/dub
本文档基于 Dub 开源仓库(dub,现代链接归因与联盟营销平台)中.agents/skills/jobs-define-job/SKILL.md的工程规范,结合apps/web/lib/jobs/的实际源码,系统讲解 Dub 后台任务框架的核心约定:如何用defineJob定义负载驱动的 QStash 任务、如何在jobLoaders注册、如何通过job.dispatch/dispatchBatch派发,以及如何把存量/api/cronworker 平滑迁移到新框架。读完本文,你将掌握在 Dub 仓库中新增一个可靠、可重试、可观测的后台任务的完整套路,并理解底层process route、envelope、outbox 重投机制的工作原理。
为什么是 defineJob:一条process route执行所有任务
在 Dub 中,负载驱动(payload-driven)的 QStash 工作一律通过defineJob定义。其核心设计是:不要为每个任务新增一个/api/jobs下的 HTTP 路由——所有任务都由现有的共享执行器apps/web/app/api/jobs/process/[jobName]/route.ts统一处理。
从源码可以看到这条路由是 Dub 后台任务的唯一执行入口:
- 先
verifyQstashSignature校验 QStash 签名(见 route.ts),防止未授权调用; - 解析请求体并用
jobEnvelopeSchema校验信封结构(name/dispatchedAt/payload,见 send-jobs.ts); - 校验 URL 中的
jobName与 envelope 内name一致(route.ts); - 通过
loadJob(jobName)从注册表中懒加载对应 handler(registry.ts),并调用job.execute(envelope.data.payload); - 捕获
z.ZodError返回 2xx(坏负载不重试),其余异常返回 500 交给 QStash 重试。
同时,Vercel GET 定时任务与无 job envelope 的扫描型 cron 仍继续使用withCron(对应cron-use-with-cronskill 的范畴),二者分工明确:defineJob管“消息驱动的工作”,withCron管“定时触发的扫描”。
文件布局:一个任务一个文件
所有后台任务代码集中在apps/web/lib/jobs/目录下:
apps/web/lib/jobs/ ├── index.ts # defineJob 框架本体,非框架变更不修改 ├── registry.ts # jobLoaders,新增任务必须在此注册 ├── send-jobs.ts # envelope + QStash 请求构建器 ├── send-workflows.ts # workflow 传输(outbox 的另一种 transport) ├── outbox.ts # 发布失败的持久化与重投 ├── constants.ts # 批大小 / 重试上限等常量 └── handlers/ └── {name}-job.ts # 每个任务一个文件其中index.ts暴露defineJob工厂:它用jobNameSchema校验任务名,返回带execute/dispatch/dispatchBatch三个方法的对象;registry.ts用静态import()构建jobLoaders映射,让 webpack 将每个 handler 代码分割成独立 chunk,按需加载。仓库现有 12 个 handler,覆盖解封合作方、文件夹/域名/标签删除后的级联清理、合作方搜索索引同步、Shopify 订单处理、欢迎邮件等场景。
第一步:创建 handler 文件
在apps/web/lib/jobs/handlers/{name}-job.ts下新建文件。命名必须是 kebab-case 且以-job结尾,由jobNameSchema强制约束:/^[a-z][a-z0-9]*(-[a-z0-9]+)*-job$/(send-jobs.ts)。
import * as z from "zod/v4"; import { defineJob } from "../index"; const inputSchema = z.object({ programId: z.string(), partnerId: z.string(), }); export const unbanPartnerJob = defineJob({ name: "unban-partner-job", schema: inputSchema, defaults: { retries: 3, // 可选;QStash 会在 5xx 时重试 // queue: "unban-partner", // 可选;命名 QStash 队列 // flowControl: { key: "unban-partner", parallelism: 20 }, }, async handle(input) { // 跳过(永久性 / 不存在)→ return(process route 返回 2xx,QStash 不重试) // 瞬时失败 → throw(process route 返回 500,QStash 重试) }, });导出的 const 需为 camelCase 的{name}Job,与 kebab-case 的name对应,例如unban-partner-job.ts导出unbanPartnerJob。
参考 handler 分类:
- 简单跳过/工作型:
unban-partner-job.ts(解封合作方后恢复被取消的佣金、待发款项与悬赏提交,并清理 fraud 告警)、create-tremendous-campaign-job.ts(为项目创建 Tremendous 奖励活动,已存在则跳过); - 自分页型:
folder-deleted-job.ts(每批处理 500 个链接后通过dispatch({ delay: 1 })续传)、domain-deleted-job.ts、partner-search-sync-job.ts; defaults.flowControl示例:partner-search-sync-job.ts,用key: "partner-search-sync", parallelism: 20限制同一时刻写入搜索提供方的并发数,并用z.discriminatedUnion设计enrollments/partners两种负载形态。
第二步:注册任务
在apps/web/lib/jobs/registry.ts的jobLoaders中新增一条静态import()记录。对象键必须与defineJob({ name })完全一致:
"unban-partner-job": () => import("./handlers/unban-partner-job").then((m) => m.unbanPartnerJob),两条硬性约束来自loadJob的实现(registry.ts):
- 必须用静态
import()——动态拼接路径会让 webpack 无法静态分析而失去代码分割,注册直接失败; - 注册键与任务名必须一致——
loadJob加载后会校验job.name !==注册键,不一致直接抛Job name mismatch。
不要通过修改 process route 来注册任务。注册完成后,registeredJobNames = Object.keys(jobLoaders)会暴露所有已注册任务名,供可观测性与回放基础设施使用。
第三步:从调用点派发任务
在业务代码中直接 import handler(而不是 importqstash),然后调用dispatch/dispatchBatch:
import { unbanPartnerJob } from "@/lib/jobs/handlers/unban-partner-job"; await unbanPartnerJob.dispatch( { workspaceId, programId, partnerId }, { label: partnerId }, ); await folderDeletedJob.dispatchBatch( folderIds.map((folderId) => ({ folderId })), ({ folderId }) => ({ label: folderId }), );派发选项(JobDispatchOptions,见 send-jobs.ts)会与defaults按“每次派发优先”合并,支持:delay、notBefore、deduplicationId、retries、queue、flowControl、label。
从index.ts的实现看,dispatch底层会:
- 单条走
qstash.publishJSON,指定了queue则走qstash.queue({ queueName }).enqueueJSON; - 批量走
qstash.batchJSON,并按QSTASH_BATCH_CHUNK_SIZE = 100(constants.ts)分块; - 发布请求自带最多 3 次指数退避重试(
withQStashRetry); - 请求体是
{ url: "/api/jobs/process/{name}", body: envelope, label }结构(send-jobs.ts),label与deduplicationId会自动拼接任务名以避免跨任务冲突。
QStash 发布失败会自动持久化到 jobs outbox(persistBackgroundJobs,见 outbox.ts):返回deferred状态,由/api/cron/queue/retry定时重投。因此除非连 outbox 持久化失败都不能中断源操作(参考queue-partner-search-sync.ts),否则不要 catch-and-swallow。
自分页:在handle内部用theJob.dispatch(nextPayload, { delay: 1 })续传下一批。以folder-deleted-job.ts为例,当本次取满MAX_LINKS_PER_BATCH = 500时,先派发带 1 秒延迟的下一批再return,否则删除文件夹本体。
不要在应用代码里调用job.execute——它只属于 process route(以及 cron 排水 shim)。
第四步:迁移存量 cron worker
把已有的withCronworker 迁到defineJob,按四步走:
- 搬移逻辑:把
withCron的函数体移进handle;把logAndRespond("skip…")替换为console.info/console.error+return; - 切换派发:注册新任务,并把所有指向该 cron URL 的
qstash.publishJSON/enqueueJSON/enqueueBatchJobs全部替换为job.dispatch/dispatchBatch; - 保留排水 shim:在旧
/api/cron/...URL 上保留一个薄 POST shim,解析旧格式请求体(不是 job envelope)后调用job.execute(payload),用于排空仍在途的 QStash 消息;全新任务不需要 shim; - 清理:删除随 handler 一起迁移的 cron-only 辅助函数。
Handle 语义:返回码即重试策略
QStash 的无限重试是后台任务的常见痛点,Dub 通过约定把“是否重试”映射到 process route 的 HTTP 状态码:
| 结果 | handle 里怎么做 | process route 返回 | QStash 行为 |
|---|---|---|---|
| 工作完成 | return | 200 | 停止 |
| 跳过(不存在 / 已完成 / 环境未配置) | console.*+return | 200 | 停止 |
| 坏负载 | 抛ZodError(schema.parse) | 200 | 停止(不可重试) |
| 瞬时失败 | throw | 500 | 重试 |
未知任务名和非法 envelope 同样返回 2xx,避免 QStash 无限重试(见 route.ts 对 envelope 解析失败、job name 不匹配、未知任务的三个 2xx 分支)。这一设计的反面是:坏负载只会在首次执行时暴露,因此schema 校验务必完整——defineJob的execute会先schema.parse(payload)(index.ts),任何 ZodError 都被 process route 判定为永久性失败。
可靠性设计:outbox 与重投闭环
defineJob框架的可观测性设计是本文档之外最值得关注的部分(outbox.ts):
- 发布失败的 job 以
job_前缀 ID 写入prisma.job表(persistBackgroundJobs),保留scheduledAt(由notBefore/delay换算)与重放选项; /api/cron/queue/retry通过publishPendingJobs按scheduledAt <= now && attempts < MAX_JOB_ATTEMPTS(MAX_JOB_ATTEMPTS = 10,每次拉取MAX_JOBS_PER_BATCH = 100条)捞出到期行,按 transport 分派重投;- 重投成功即删除行,失败则
attempts + 1并记录lastError(截断至 1000 字符),超过上限记录jobs.retry_exhausted告警; send-jobs.ts的buildReplayRequest会把重投请求构造成 batch 形式,并保留原始dispatchedAt与notBefore,保证重放语义与首次派发一致。
不要用 defineJob 的场景
以下情况不要使用defineJob:
apps/web/vercel.json中的 Vercel GET crons;- 只负责扇出(fan out)的扫描器——它 enqueue 的 worker 可以是 job,但扫描器本身不是;
- 需要把续传状态重新发布到同一 URL 的导入器(importer);
/api/cron/queue/retry本身(job 回放基础设施);- 路径参数身份(
/api/cron/links/[linkId]/…)——除非把 id 移入 payload; - 出站 webhook 转发与 postback。
禁区清单(Do not)
- 不为负载驱动工作新增
/api/jobs/...路由或新的/api/cron/...POST worker; - 不要用
qstash.publishJSON/enqueueJSON/enqueueBatchJobs指向/api/jobs/process/...——统一用dispatch; - 不要用 webpack 无法静态分析的动态
import()注册,也不要让注册键与defineJob({ name })不一致; - 任务名不能缺少
-job后缀; - 不要在派发调用点使用
job.execute; - 新增 job 后不要执行
pnpm build(本地构建会被 webpack 动态 chunk 干扰,CI 会负责构建验证)。
小结
Dub 的后台任务体系把“定义—注册—派发—执行—重试—重投”收敛成一条清晰的生产链路:defineJob用 zod schema 保证负载即契约,统一的/api/jobs/process/[jobName]路由把“是否重试”编码进 HTTP 状态码,registry的静态import()实现按需加载,而 outbox 让发布失败也不丢消息。如果你需要在仓库中新增一个后台任务,只需照抄本文四步:建{name}-job.tshandler → 在registry.ts静态注册 → 调用dispatch/dispatchBatch→(如为迁移)保留排水 shim,即可获得与现有 12 个任务同级的可靠性。
【免费下载链接】dubThe modern link attribution platform. Loved by world-class marketing teams like Framer, Perplexity, Superhuman, Twilio, Buffer and more.项目地址: https://gitcode.com/GitHub_Trending/du/dub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考