深入 Dub 后台任务框架:用 defineJob 统一 QStash 负载驱动任务
2026/9/11 17:17:04 网站建设 项目流程

深入 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.tspartner-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.tsjobLoaders中新增一条静态import()记录。对象键必须与defineJob({ name })完全一致

"unban-partner-job": () => import("./handlers/unban-partner-job").then((m) => m.unbanPartnerJob),

两条硬性约束来自loadJob的实现(registry.ts):

  1. 必须用静态import()——动态拼接路径会让 webpack 无法静态分析而失去代码分割,注册直接失败;
  2. 注册键与任务名必须一致——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按“每次派发优先”合并,支持:delaynotBeforededuplicationIdretriesqueueflowControllabel

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),labeldeduplicationId会自动拼接任务名以避免跨任务冲突。

QStash 发布失败会自动持久化到 jobs outboxpersistBackgroundJobs,见 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,按四步走:

  1. 搬移逻辑:把withCron的函数体移进handle;把logAndRespond("skip…")替换为console.info/console.error+return
  2. 切换派发:注册新任务,并把所有指向该 cron URL 的qstash.publishJSON/enqueueJSON/enqueueBatchJobs全部替换为job.dispatch/dispatchBatch
  3. 保留排水 shim:在旧/api/cron/...URL 上保留一个薄 POST shim,解析旧格式请求体(不是 job envelope)后调用job.execute(payload),用于排空仍在途的 QStash 消息;全新任务不需要 shim
  4. 清理:删除随 handler 一起迁移的 cron-only 辅助函数。

Handle 语义:返回码即重试策略

QStash 的无限重试是后台任务的常见痛点,Dub 通过约定把“是否重试”映射到 process route 的 HTTP 状态码:

结果handle 里怎么做process route 返回QStash 行为
工作完成return200停止
跳过(不存在 / 已完成 / 环境未配置)console.*+return200停止
坏负载ZodError(schema.parse)200停止(不可重试)
瞬时失败throw500重试

未知任务名和非法 envelope 同样返回 2xx,避免 QStash 无限重试(见 route.ts 对 envelope 解析失败、job name 不匹配、未知任务的三个 2xx 分支)。这一设计的反面是:坏负载只会在首次执行时暴露,因此schema 校验务必完整——defineJobexecute会先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通过publishPendingJobsscheduledAt <= now && attempts < MAX_JOB_ATTEMPTSMAX_JOB_ATTEMPTS = 10,每次拉取MAX_JOBS_PER_BATCH = 100条)捞出到期行,按 transport 分派重投;
  • 重投成功即删除行,失败则attempts + 1并记录lastError(截断至 1000 字符),超过上限记录jobs.retry_exhausted告警;
  • send-jobs.tsbuildReplayRequest会把重投请求构造成 batch 形式,并保留原始dispatchedAtnotBefore,保证重放语义与首次派发一致。

不要用 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),仅供参考

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

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

立即咨询