news 2026/10/2 8:18:17

Papermark 中的 Trigger.dev v4 高级任务实战:Tags、批量触发、防抖、队列、重试与幂等设计

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Papermark 中的 Trigger.dev v4 高级任务实战:Tags、批量触发、防抖、队列、重试与幂等设计
  • 后端
  • 前端
  • 企业应用

【免费下载链接】papermark

Papermark is the open-source DocSend alternative and secure data rooms with built-in analytics and custom domains.

项目地址:https://gitcode.com/GitHub_Trending/pa/papermark
点击查看免费下载

本指南系统讲解 Trigger.dev v4 高级任务(Advanced Tasks)模式——包括标签组织、批量触发 v2、防抖合并、并发队列、指数退避重试、机器预设、幂等键、元数据进度追踪与日志追踪等——并以开源数据室项目 Papermark 的真实任务实现(如 bulk-download.ts、convert-pdf-direct.ts、queues.ts)为源码级佐证。读完本文,你将掌握如何在生产项目中编排高可靠、可观测、可横向扩展的后台任务流水线。

本文依据仓库内 advanced-tasks.md 撰写,并配套参考 basic-tasks.md、scheduled-tasks.md 与 SKILL.md 三份技能文档。所有任务都必须基于@trigger.dev/sdk(v4),严禁使用已废弃的client.defineJob(v2 写法会破坏应用)。

1. Tags 与任务组织

标签(Tags)是 Trigger.dev 提供的轻量级运行标记机制,用于在运行列表、日志和订阅流中按业务维度筛选任务运行。

1.1 运行中添加标签

import { task, tags } from "@trigger.dev/sdk"; export const processUser = task({ id: "process-user", run: async (payload: { userId: string; orgId: string }, { ctx }) => { // Add tags during execution await tags.add(`user_${payload.userId}`); await tags.add(`org_${payload.orgId}`); return { processed: true }; }, });

1.2 触发时携带标签

// Trigger with tags await processUser.trigger( { userId: "123", orgId: "abc" }, { tags: ["priority", "user_123", "org_abc"] } // Max 10 tags per run );

1.3 订阅指定标签的运行流

// Subscribe to tagged runs for await (const run of runs.subscribeToRunsWithTag("user_123")) { console.log(`User task ${run.id}: ${run.status}`); }

标签最佳实践:

  • 使用前缀命名:user_123、org_abc、video:456;
  • 每次运行最多 10 个标签,每个标签长度 1~64 字符;
  • 标签不会自动传播到子任务——子任务若需要同样标签,需显式在trigger的 options 中传入。

2. 批量触发 v2(Batch Triggering)

批量触发用于一次拉起多达上千条独立运行,适用于批量文档处理、批量邮件、批量导出等场景。v2 相比 v1 显著提升了单批容量。

2.1 限制与速率

  • 最大批量大小:1,000 项(v1 为 500);
  • 单项 payload:每个最多 3MB(v1 为整体合计 1MB);
  • 单项 payload 超过 512KB 时,会自动卸载(offload)到对象存储,避免请求体积过大。

按环境的速率限制(per environment):

TierBucket SizeRefill Rate
Free1,200 runs100 runs/10 sec
Hobby5,000 runs500 runs/5 sec
Pro5,000 runs500 runs/5 sec

并发批量处理(Concurrent Batches):

TierConcurrent Batches
Free1
Hobby10
Pro10

2.2 基本用法

import { myTask } from "./trigger/myTask"; // Basic batch trigger (up to 1,000 items) const runs = await myTask.batchTrigger([ { payload: { userId: "user-1" } }, { payload: { userId: "user-2" } }, { payload: { userId: "user-3" } }, ]); // Batch trigger with wait const results = await myTask.batchTriggerAndWait([ { payload: { userId: "user-1" } }, { payload: { userId: "user-2" } }, ]); for (const result of results) { if (result.ok) { console.log("Result:", result.output); } }

2.3 逐项 options(幂等键、标签)

const batchHandle = await myTask.batchTrigger([ { payload: { userId: "123" }, options: { idempotencyKey: "user-123-batch", tags: ["priority"], }, }, { payload: { userId: "456" }, options: { idempotencyKey: "user-456-batch", }, }, ]);

2.4 Papermark 中的批量思想:分批而不是一把梭

Papermark 的bulkDownloadTask(lib/trigger/bulk-download.ts)虽然未直接使用batchTrigger,但体现了“单任务内部批次化”的配套策略:它把大量文件按MAX_FILES_PER_BATCH = 500与MAX_ZIP_SIZE_BYTES = 500MB拆成多个 ZIP 批次,然后顺序调用 AWS Lambda 打包(processDownloadBatch),并在每个批次完成后更新 Redis 任务进度。其注释明确说明“Process each batch sequentially (to avoid Lambda concurrency issues)”。这说明:当外部依赖(Lambda、第三方 API)并发能力有限时,即使批量触发容量很大,也仍要在任务内部做分片与串行化控制。从后端代码触发批量任务的另一种入口是tasks.batchTrigger<typeof processData>("process-data", [...])(见 basic-tasks.md),批量上限同样是 1,000 项、单项 3MB。

3. 防抖(Debouncing)

防抖把同一防抖键(key)下、延迟窗口内的多次触发合并为一次执行,用于吸收高频事件。

3.1 适用场景

  • 用户活动更新:把用户的连续快速操作合并为一次运行;
  • Webhook 去重:处理 webhook 突发流量,避免重复处理;
  • 搜索索引更新:把多次文档变更合并为一次索引写入;
  • 通知批处理:合并通知发送,防止骚扰用户。

3.2 基本用法

await myTask.trigger( { userId: "123" }, { debounce: { key: "user-123-update", // Unique identifier for debounce group delay: "5s", // Wait duration ("5s", "1m", or milliseconds) }, } );

delay支持"5s"、"1m"等人类可读格式,也支持毫秒数值。

3.3 执行模式:leading 与 trailing

Leading 模式(默认):使用第一次触发时的 payload 与 options;后续触发只把执行时间往后重新调度。

// First trigger sets the payload await myTask.trigger({ action: "first" }, { debounce: { key: "my-key", delay: "10s" } }); // Second trigger only reschedules - payload remains "first" await myTask.trigger({ action: "second" }, { debounce: { key: "my-key", delay: "10s" } }); // Task executes with { action: "first" }

Trailing 模式:使用最近一次触发时的 payload 与 options。

await myTask.trigger( { data: "latest-value" }, { debounce: { key: "trailing-example", delay: "10s", mode: "trailing", }, } );

在 trailing 模式下,以下选项会随每次触发更新:

  • payload— 任务输入数据
  • metadata— 运行元数据
  • tags— 运行标签(整体替换)
  • maxAttempts— 重试次数
  • maxDuration— 最大计算时长
  • machine— 机器预设

3.4 重要注意事项

  • 幂等键优先于防抖设置:若同时提供idempotencyKey,以幂等键为准;
  • 兼容triggerAndWait():父任务在防抖后的执行上会正确阻塞等待,不会提前返回;
  • 防抖键的作用域是单个任务:不同任务之间的防抖键互不干扰。

4. 并发与队列(Concurrency & Queues)

队列用于限制同一任务的并发执行数,保护下游外部服务不被压垮。Papermark 在 lib/trigger/queues.ts 中定义了一整套按业务与套餐划分的队列,是这一模式的直接落地。

4.1 共享队列、任务级并发、按用户并发

import { task, queue } from "@trigger.dev/sdk"; // Shared queue for related tasks const emailQueue = queue({ name: "email-processing", concurrencyLimit: 5, // Max 5 emails processing simultaneously }); // Task-level concurrency export const oneAtATime = task({ id: "sequential-task", queue: { concurrencyLimit: 1 }, // Process one at a time run: async (payload) => { // Critical section - only one instance runs }, }); // Per-user concurrency export const processUserData = task({ id: "process-user-data", run: async (payload: { userId: string }) => { // Override queue with user-specific concurrency await childTask.trigger(payload, { queue: { name: `user-${payload.userId}`, concurrencyLimit: 2, }, }); }, }); export const emailTask = task({ id: "send-email", queue: emailQueue, // Use shared queue run: async (payload: { to: string }) => { // Send email logic }, });

4.2 Papermark 的队列设计:按任务类型与按套餐隔离

Papermark 的 queues.ts 先为每种重活定义专属队列:

import { queue } from "@trigger.dev/sdk"; // Task-specific queues export const convertFilesToPdfQueue = queue({ name: "convert-files-to-pdf", concurrencyLimit: 10, }); export const convertCadToPdfQueue = queue({ name: "convert-cad-to-pdf", concurrencyLimit: 2, });

它还按套餐(plan)隔离转换任务的并发,防止免费/入门套餐用户的任务挤占高级套餐资源:

// Plan-based conversion queues (used at trigger time) const concurrencyConfig: Record<string, number> = { free: 1, starter: 1, pro: 2, business: 10, datarooms: 10, // ... }; export const conversionFreeQueue = queue({ name: "conversion-free", concurrencyLimit: 1 }); export const conversionProQueue = queue({ name: "conversion-pro", concurrencyLimit: 2 }); // ... /** * Returns the queue name string for the given plan. * The queue must be pre-defined above (v4 requirement). */ export const conversionQueueName = (plan: string): string => { const planName = plan.split("+")[0] as BasePlan; return `conversion-${planName}`; };

注意注释中的关键约束:v4 中队列必须在任务文件中预先定义(The queue must be pre-defined above),触发时只能引用已声明的队列名。这与“按用户动态创建队列”的写法不同——若需要按租户隔离,应像 4.1 示例那样在触发 options 中传入queue.name动态值。同时,不要在任务中把triggerAndWait()或wait.*调用包进Promise.all/Promise.allSettled,这是 Trigger.dev 任务不支持的写法(见 basic-tasks.md)。

5. 错误处理与重试(Error Handling & Retries)

5.1 任务级重试配置 + catchError 钩子

import { task, retry, AbortTaskRunError } from "@trigger.dev/sdk"; export const resilientTask = task({ id: "resilient-task", retry: { maxAttempts: 10, factor: 1.8, // Exponential backoff multiplier minTimeoutInMs: 500, maxTimeoutInMs: 30_000, randomize: false, }, catchError: async ({ error, ctx }) => { // Custom error handling if (error.code === "FATAL_ERROR") { throw new AbortTaskRunError("Cannot retry this error"); } // Log error details console.error(`Task ${ctx.task.id} failed:`, error); // Allow retry by returning nothing return { retryAt: new Date(Date.now() + 60000) }; // Retry in 1 minute }, run: async (payload) => { // Retry specific operations const result = await retry.onThrow( async () => { return await unstableApiCall(payload); }, { maxAttempts: 3 } ); // Conditional HTTP retries const response = await retry.fetch("https://api.example.com", { retry: { maxAttempts: 5, condition: (response, error) => { return response?.status === 429 || response?.status >= 500; }, }, }); return result; }, });

要点:

  • retry.onThrow对单次操作做局部重试;
  • retry.fetch支持条件重试——示例中仅对 429(限流)与 >= 500(服务端错误)重试,避免对 4xx 客户端错误做无意义重试;
  • catchError中抛出AbortTaskRunError表示不可重试的致命错误,立即终止运行;
  • catchError返回{ retryAt }则把下一次重试调度到指定时间(示例为 1 分钟后)。

5.2 区分可重试与不可重试错误

Papermark 的bulkDownloadTask是区分“暂时性失败”与“永久性失败”的教科书案例(lib/trigger/bulk-download.ts):它定义了一个哨兵错误类DownloadNotPermittedError,在catch块中判断——若是权限类失败(链接被归档、过期、删除、禁用下载),不抛错而是直接返回,避免浪费重试次数:

// Permission failures aren't transient: don't waste a retry on them. if (isPermissionFailure) { return { success: false, jobId, downloadUrls: [] }; } throw error;

同时任务本身声明retry: { maxAttempts: 2 },配合“校验失败即返回”的策略,把重试预算留给真正可能因瞬时抖动成功的批次。convertPdfDirectTask(convert-pdf-direct.ts)同样在拉取 PDF 失败或检测到屏蔽关键词时抛出AbortTaskRunError终止重试,例如throw new AbortTaskRunError("Failed to fetch PDF")。

5.3 全局默认重试

Papermark 在 trigger.config.ts 中配置了项目级默认重试策略,所有未显式声明retry的任务都会继承:

retries: { enabledInDev: false, default: { maxAttempts: 3, minTimeoutInMs: 1000, maxTimeoutInMs: 10000, factor: 2, randomize: true, }, },

即最多 3 次尝试、初始退避 1 秒、指数因子 2、上限 10 秒并带随机抖动(randomize),且开发环境下默认不重试(便于本地调试)。

6. 机器预设与性能(Machines & Performance)

任务可以通过machine.preset选择计算资源档位,并通过maxDuration设置超时上限(秒)。

export const heavyTask = task({ id: "heavy-computation", machine: { preset: "large-2x" }, // 8 vCPU, 16 GB RAM maxDuration: 1800, // 30 minutes timeout run: async (payload, { ctx }) => { // Resource-intensive computation if (ctx.machine.preset === "large-2x") { // Use all available cores return await parallelProcessing(payload); } return await standardProcessing(payload); }, }); // Override machine when triggering await heavyTask.trigger(payload, { machine: { preset: "medium-1x" }, // Override for this run });

机器预设表:

PresetvCPURAM
micro0.250.25 GB
small-1x0.50.5 GB(默认)
small-2x11 GB
medium-1x12 GB
medium-2x24 GB
large-1x48 GB
large-2x816 GB

Papermark 的实践:bulkDownloadTask使用默认档small-1x(bulk-download.ts);convertPdfDirectTask使用large-1x(4 vCPU / 8 GB)执行重型的 mupdf WASM 页面渲染(convert-pdf-direct.ts)。此外 trigger.config.ts 设置了maxDuration: timeout.None,即项目全局不限制任务最长运行时间——这意味着单个任务自身的maxDuration或运行时长上限由任务级配置决定,适合耗时长的 PDF 转换与批量导出。

7. 幂等性(Idempotency)

幂等键保证“同键只执行一次”,是支付、扣款、发信等关键操作的必备设施。Trigger.dev 提供idempotencyKeys.create()生成作用域化键,并支持idempotencyKeyTTL设置键的过期时间。

import { task, idempotencyKeys } from "@trigger.dev/sdk"; export const paymentTask = task({ id: "process-payment", retry: { maxAttempts: 3, }, run: async (payload: { orderId: string; amount: number }) => { // Automatically scoped to this task run, so if the task is retried, the idempotency key will be the same const idempotencyKey = await idempotencyKeys.create(`payment-${payload.orderId}`); // Ensure payment is processed only once await chargeCustomer.trigger(payload, { idempotencyKey, idempotencyKeyTTL: "24h", // Key expires in 24 hours }); }, });

关键点:idempotencyKeys.create()的键自动作用域到当前任务运行——如果任务因重试而再次执行,生成的键保持一致,从而保证子任务只被触发一次。

7.1 基于 payload 的幂等

对“同一输入只处理一次”的需求,可对 payload 计算哈希作为幂等键:

// Payload-based idempotency import { createHash } from "node:crypto"; function createPayloadHash(payload: any): string { const hash = createHash("sha256"); hash.update(JSON.stringify(payload)); return hash.digest("hex"); } export const deduplicatedTask = task({ id: "deduplicated-task", run: async (payload) => { const payloadHash = createPayloadHash(payload); const idempotencyKey = await idempotencyKeys.create(payloadHash); await processData.trigger(payload, { idempotencyKey }); }, });

7.2 批量触发中的逐项幂等

回到 2.3 节:batchTrigger的每个 item 都可以通过options.idempotencyKey单独声明幂等键,与单次触发保持一致语义。注意幂等键优先于防抖设置(见 3.4)。

8. 元数据与进度追踪(Metadata & Progress Tracking)

metadataAPI 可以在运行过程中持续记录结构化状态,供 Dashboard 实时展示;子任务还能反向更新父任务的元数据。

8.1 初始化与迭代更新

import { task, metadata } from "@trigger.dev/sdk"; export const batchProcessor = task({ id: "batch-processor", run: async (payload: { items: any[] }, { ctx }) => { const totalItems = payload.items.length; // Initialize progress metadata metadata .set("progress", 0) .set("totalItems", totalItems) .set("processedItems", 0) .set("status", "starting"); const results = []; for (let i = 0; i < payload.items.length; i++) { const item = payload.items[i]; // Process item const result = await processItem(item); results.push(result); // Update progress const progress = ((i + 1) / totalItems) * 100; metadata .set("progress", progress) .increment("processedItems", 1) .append("logs", `Processed item ${i + 1}/${totalItems}`) .set("currentItem", item.id); } // Final status metadata.set("status", "completed"); return { results, totalProcessed: results.length }; }, });

metadata链式方法包括:set(key, value)(赋值)、increment(key, n)(数值自增)、append(key, value)(向列表追加)。

8.2 子任务更新父级/根级元数据

// Update parent metadata from child task export const childTask = task({ id: "child-task", run: async (payload, { ctx }) => { // Update parent task metadata metadata.parent.set("childStatus", "processing"); metadata.root.increment("childrenCompleted", 1); return { processed: true }; }, });

Papermark 的进度追踪实践:convertPdfDirectTask封装了setProgress辅助函数,同时更新自身与父任务(如有)的元数据(convert-pdf-direct.ts):

function setProgress(status: { progress: number; text: string }) { metadata.set("status", status); try { metadata.parent.set("status", status); } catch { // no parent task } }

在逐页渲染 PDF 时,它以20 + (pageNumber / numPages) * 70计算进度并写入元数据(L299-L302),最终在完成时置为 100%。与此同时,bulkDownloadTask把真实进度写入 Redis 任务存储(downloadJobStore.updateJob(jobId, { status, progress, processedFiles, ... }),见 bulk-download.ts),前端据此渲染“Preparing... / ZIPPING / 百分比”等状态——这是“元数据/外部存储双重追踪”的典型组合。

9. 日志与追踪(Logging & Tracing)

logger提供结构化日志;logger.trace创建自定义 span 并附加属性,用于链路追踪。

import { task, logger } from "@trigger.dev/sdk"; export const tracedTask = task({ id: "traced-task", run: async (payload, { ctx }) => { logger.info("Task started", { userId: payload.userId }); // Custom trace with attributes const user = await logger.trace( "fetch-user", async (span) => { span.setAttribute("user.id", payload.userId); span.setAttribute("operation", "database-fetch"); const userData = await database.findUser(payload.userId); span.setAttribute("user.found", !!userData); return userData; }, { userId: payload.userId } ); logger.debug("User fetched", { user: user.id }); try { const result = await processUser(user); logger.info("Processing completed", { result }); return result; } catch (error) { logger.error("Processing failed", { error: error.message, userId: payload.userId, }); throw error; } }, });

用法说明:logger.info/debug/error/warn均接受(message, attributes?);logger.trace(name, fn, attributes?)内通过span.setAttribute(key, value)记录追踪属性,fn的返回值即为trace的返回值。Papermark 在bulkDownloadTask中大量使用logger.info/logger.warn/logger.error记录 jobId、batchNumber、文件数与字节数等结构化上下文(如 bulk-download.ts),使单次运行的全部日志可按 jobId 串联检索。

10. 隐藏任务(Hidden Tasks)

不为export的任务文件内部模块是“隐藏任务”,不会被 CLI 注册成公开任务,但可以被公开任务通过triggerAndWait调用,用于封装内部实现细节。

// Hidden task - not exported, only used internally const internalProcessor = task({ id: "internal-processor", run: async (payload: { data: string }) => { return { processed: payload.data.toUpperCase() }; }, }); // Public task that uses hidden task export const publicWorkflow = task({ id: "public-workflow", run: async (payload: { input: string }) => { // Use hidden task internally const result = await internalProcessor.triggerAndWait({ data: payload.input, }); if (result.ok) { return { output: result.output.processed }; } throw new Error("Internal processing failed"); }, });

两个关键点:

  • 隐藏任务必须定义在同一模块文件内(未导出),通过模块闭包被公开任务引用;
  • 公开任务中通过internalProcessor.triggerAndWait(...)调用,返回值是Result对象——必须先检查result.ok再访问result.output,这是 SKILL.md 中列出的硬性规则之一。

11. 高级任务最佳实践汇总

  1. 并发:用队列防止压垮外部服务——Papermark 按任务类型和套餐各建队列,见 queues.ts;
  2. 重试:为瞬时故障配置指数退避——参考 trigger.config.ts 的全局默认值与retry.onThrow/retry.fetch局部重试;
  3. 幂等:支付/关键操作务必使用幂等键(idempotencyKeys.create+idempotencyKeyTTL);
  4. 元数据:长任务用metadata记录进度,子任务用metadata.parent/root汇总状态;
  5. 机器预设:按计算需求匹配机器档位——轻量任务用默认small-1x,重型渲染用large-1x/large-2x;
  6. 标签:用一致的前缀命名(user_123、org_abc、video:456)便于过滤与订阅;
  7. 防抖:用于用户活动、webhook 突发与通知批处理,注意 leading/trailing 语义差异;
  8. 批量触发:用于最多 1,000 项的大批量操作,单项最大 3MB,超过 512KB 自动卸载到对象存储;
  9. 错误处理:区分可重试错误与致命错误——致命错误抛AbortTaskRunError(如 convert-pdf-direct.ts),永久性失败直接返回不浪费重试(如 bulk-download.ts)。

设计原则:任务应设计为无状态、幂等、对失败有韧性——用元数据做状态追踪,用队列做资源管理。在 v4 中切记:一律使用@trigger.dev/sdk的task/schemaTask/schedules.taskAPI,检查result.ok后再取result.output,不要对triggerAndWait()/wait.*使用Promise.all,并从trigger/目录导出任务。上述技能文档全文可继续查阅 SKILL.md、basic-tasks.md 与 scheduled-tasks.md。

  • 后端
  • 前端
  • 企业应用

【免费下载链接】papermark

Papermark is the open-source DocSend alternative and secure data rooms with built-in analytics and custom domains.

项目地址:https://gitcode.com/GitHub_Trending/pa/papermark
点击查看免费下载

相关推荐

上一篇:Repomix 远程 GitHub 仓库处理实战:从 URL 简写解析到配置信任的安全模型
下一篇:TDengine 通过 taosExplorer 零代码接入 Kafka:数据写入任务完整配置指南

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/2 8:16:14

OpenRig cross-host架构深度解析:在多台机器上运行一个Agent团队

OpenRig cross-host架构深度解析&#xff1a;在多台机器上运行一个Agent团队 【免费下载链接】openrig Multi-agent harness that runs Claude Code and Codex together as one system 项目地址: https://gitcode.com/GitHub_Trending/op/openrig OpenRig 是一个开源的多…

作者头像 李华
网站建设 2026/10/2 8:11:11

Python + Neo4j 知识图谱上传实战:从 CSV 到可查图谱的完整链路

简介&#xff1a;本资源为基于Python与Neo4j的知识图谱上传与处理设计源码&#xff0c;面向希望掌握图数据库应用、数据上传与图查询分析的开发者与研究人员&#xff0c;可作为课程设计、毕业项目或工程实践的参考方案。压缩包共25个文件&#xff0c;约27.84MB&#xff0c;以12…

作者头像 李华