后端前端企业应用【免费下载链接】papermarkPapermark 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/sdkv4严禁使用已废弃的client.defineJobv2 写法会破坏应用。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. 批量触发 v2Batch Triggering批量触发用于一次拉起多达上千条独立运行适用于批量文档处理、批量邮件、批量导出等场景。v2 相比 v1 显著提升了单批容量。2.1 限制与速率最大批量大小1,000 项v1 为 500单项 payload每个最多 3MBv1 为整体合计 1MB单项 payload 超过 512KB 时会自动卸载offload到对象存储避免请求体积过大。按环境的速率限制per environmentTierBucket SizeRefill RateFree1,200 runs100 runs/10 secHobby5,000 runs500 runs/5 secPro5,000 runs500 runs/5 sec并发批量处理Concurrent BatchesTierConcurrent BatchesFree1Hobby10Pro102.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 的bulkDownloadTasklib/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.batchTriggertypeof 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 与 trailingLeading 模式默认使用第一次触发时的 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: Recordstring, 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 Retries5.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 arent transient: dont waste a retry on them. if (isPermissionFailure) { return { success: false, jobId, downloadUrls: [] }; } throw error;同时任务本身声明retry: { maxAttempts: 2 }配合“校验失败即返回”的策略把重试预算留给真正可能因瞬时抖动成功的批次。convertPdfDirectTaskconvert-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 });机器预设表PresetvCPURAMmicro0.250.25 GBsmall-1x0.50.5 GB默认small-2x11 GBmedium-1x12 GBmedium-2x24 GBlarge-1x48 GBlarge-2x816 GBPapermark 的实践bulkDownloadTask使用默认档small-1xbulk-download.tsconvertPdfDirectTask使用large-1x4 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 TrackingmetadataAPI 可以在运行过程中持续记录结构化状态供 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.tsfunction 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 Tracinglogger提供结构化日志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. 高级任务最佳实践汇总并发用队列防止压垮外部服务——Papermark 按任务类型和套餐各建队列见 queues.ts重试为瞬时故障配置指数退避——参考 trigger.config.ts 的全局默认值与retry.onThrow/retry.fetch局部重试幂等支付/关键操作务必使用幂等键idempotencyKeys.createidempotencyKeyTTL元数据长任务用metadata记录进度子任务用metadata.parent/root汇总状态机器预设按计算需求匹配机器档位——轻量任务用默认small-1x重型渲染用large-1x/large-2x标签用一致的前缀命名user_123、org_abc、video:456便于过滤与订阅防抖用于用户活动、webhook 突发与通知批处理注意 leading/trailing 语义差异批量触发用于最多 1,000 项的大批量操作单项最大 3MB超过 512KB 自动卸载到对象存储错误处理区分可重试错误与致命错误——致命错误抛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。赞分享后端前端企业应用【免费下载链接】papermarkPapermark 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 高级任务开发指南Tags、批量触发、防抖、并发队列与幂等实战Trigger.dev v4 高级任务开发指南Tags、批量触发、防抖、并发队列与幂等实战 本篇技术指南以 Trigger.dev v4 的 AdvancedAI Agent后端任务调度开发工具可观测性AI 应用Trigger.dev v4 高级任务开发指南Tags、批量触发、防抖、并发、重试与幂等全解析Trigger.dev v4 高级任务开发指南Tags、批量触发、防抖、并发、重试与幂等全解析 Trigger.dev 是一个用于构建和部署持久化durabAI Agent后端任务调度开发工具可观测性AI 应用Trigger.dev v4 高级任务编写指南标签、并发、重试、幂等与性能调优实战Trigger.dev v4 高级任务编写指南标签、并发、重试、幂等与性能调优实战 本文聚焦 Trigger.dev v4当前仓库中编写生产级任务的进阶模AI Agent后端任务调度开发工具可观测性AI 应用上一篇Repomix 远程 GitHub 仓库处理实战从 URL 简写解析到配置信任的安全模型下一篇TDengine 通过 taosExplorer 零代码接入 Kafka数据写入任务完整配置指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考