深入 Dub 后台任务框架:用 defineJob 统一 QStash 负载驱动任务

发布时间:2026/9/11 17:17:07
深入 Dub 后台任务框架:用 defineJob 统一 QStash 负载驱动任务 深入 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)从注册表中懒加载对应 handlerregistry.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) { // 跳过永久性 / 不存在→ returnprocess route 返回 2xxQStash 不重试 // 瞬时失败 → throwprocess route 返回 500QStash 重试 }, });导出的 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.tsdefaults.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/dispatchBatchimport { 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 100constants.ts分块发布请求自带最多 3 次指数退避重试withQStashRetry请求体是{ url: /api/jobs/process/{name}, body: envelope, label }结构send-jobs.tslabel与deduplicationId会自动拼接任务名以避免跨任务冲突。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按四步走搬移逻辑把withCron的函数体移进handle把logAndRespond(skip…)替换为console.info/console.errorreturn切换派发注册新任务并把所有指向该 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 行为工作完成return200停止跳过不存在 / 已完成 / 环境未配置console.*return200停止坏负载抛ZodErrorschema.parse200停止不可重试瞬时失败throw500重试未知任务名和非法 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_ATTEMPTSMAX_JOB_ATTEMPTS 10每次拉取MAX_JOBS_PER_BATCH 100条捞出到期行按 transport 分派重投重投成功即删除行失败则attempts 1并记录lastError截断至 1000 字符超过上限记录jobs.retry_exhausted告警send-jobs.ts的buildReplayRequest会把重投请求构造成 batch 形式并保留原始dispatchedAt与notBefore保证重放语义与首次派发一致。不要用 defineJob 的场景以下情况不要使用defineJobapps/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),仅供参考