Cloudflare Workflows 工作流模式实战:图像处理、用户生命周期、数据管道与人机审批的可靠编排方案

发布时间:2026/9/12 21:45:18
Cloudflare Workflows 工作流模式实战:图像处理、用户生命周期、数据管道与人机审批的可靠编排方案 Cloudflare Workflows 工作流模式实战图像处理、用户生命周期、数据管道与人机审批的可靠编排方案【免费下载链接】skillsSkills Catalog for Codex项目地址: https://gitcode.com/GitHub_Trending/skills4/skills本文基于 skills 仓库中cloudflare-deploy技能包的 Workflows 参考文档系统讲解 Cloudflare Workflows 中最实用的五大工作流模式图像处理管道、用户生命周期、数据管道、人机审批、定时链路与编排技巧并覆盖测试工作流与最佳实践。读完本文你将能直接用step.do()、step.sleep()、step.waitForEvent()组合出可重试、可休眠、状态持久化的长时多步应用并在本地用 Vitest Introspection API 对工作流做确定性测试。Cloudflare Workflows 是一种持久化多步应用运行时每一步独立可重试、步间状态自动持久化、实例可以休眠数天等待外部事件而不占用并发资源。本文对应的参考文档位于 patterns.md与之配套的还有 README.md核心概念、configuration.mdwrangler 配置与步骤配置、api.mdStep API 与实例管理和 gotchas.md陷阱、限额与定价。本文将以 patterns.md 为骨架逐模式展开并交叉引用上述文档中的事实细节。在动手写模式之前先记住三个贯穿全文的核心概念详见 README.mdWorkflow继承WorkflowEntrypoint的类实现run(event, step)方法Instance一次独立执行拥有唯一 ID 与独立状态Step通过step.do()定义的独立可重试单元API 调用、数据库查询、AI 推理等State由每个 step 的返回值持久化步骤名即缓存键。一、图像处理管道AI 描述生成 人工审核后发布第一个模式解决的是典型的内容审核后发布问题对象存储里有用户上传的图片工作流取出图片 → 用 Workers AI 生成文字描述 → 等待人工批准 → 批准后才发布到公开目录。export class ImageProcessingWorkflow extends WorkflowEntrypointEnv, Params { async run(event, step) { // ① 拉取 R2 中的原始图片BUCKET.get 返回对象arrayBuffer() 取出二进制 const imageData await step.do(fetch, async () (await this.env.BUCKET.get(event.params.imageKey)).arrayBuffer() ); // ② 交给 Workers AI 的多模态模型生成描述 // llava-1.5-7b-hf 接收图片字节数组 prompt返回文字描述 const description await step.do(generate description, async () await this.env.AI.run(cf/llava-hf/llava-1.5-7b-hf, { image: Array.from(new Uint8Array(imageData)), prompt: Describe this image, max_tokens: 50 }) ); // ③ 挂起等待人工审核事件24 小时超时 // 超时会抛出异常需要用 try-catch 兜底 await step.waitForEvent(await approval, { event: approved, timeout: 24h }); // ④ 审核通过后把图片写入 R2 的 public/ 前缀下完成发布 await step.do(publish, async () await this.env.BUCKET.put(public/${event.params.imageKey}, imageData) ); } }这个模式的关键点step.do()之间不共享内存变量imageData能跨步骤使用是因为它被持久化进了步骤状态而不是留在内存里。每一步的返回值都会自动持久化详见 api.md。step.waitForEvent()让实例进入waiting状态处于该状态的实例不计入并发配额见 gotchas.md 限额表注释因此你可以同时挂起海量待审核任务。审核动作来自外部调用方通过instance.sendEvent({type: approved, payload})投递事件事件类型必须与waitForEvent中声明的event一致见 api.md。description变量在示例中未被使用实际业务里可将其写入 D1/KV 或作为审核界面的展示内容——这正是Human-in-the-Loop场景的常见做法。二、用户生命周期欢迎邮件、试用期计时与转化检测第二个模式演示了工作流如何扮演业务状态机新用户注册后先发欢迎邮件接着sleep一个 7 天试用期不占用任何计算资源再检查用户是否已升级为付费订阅未转化则发送试用到期提醒邮件。export class UserLifecycleWorkflow extends WorkflowEntrypointEnv, Params { async run(event, step) { // ① 发送欢迎邮件 await step.do(welcome email, async () await sendEmail(event.params.email, Welcome!) ); // ② 睡眠 7 天实例进入 waiting 状态不计并发、不耗 CPU await step.sleep(trial period, 7 days); // ③ 查询 D1 判断用户是否已转化为付费用户 const hasConverted await step.do(check conversion, async () { const user await this.env.DB.prepare( SELECT subscription_status FROM users WHERE id ? ).bind(event.params.userId).first(); return user.subscription_status active; }); // ④ 依据步骤输出而非内存变量做确定性分支 if (!hasConverted) { await step.do(trial expiration email, async () await sendEmail(event.params.email, Trial ending) ); } } }设计要点step.sleep()支持时间字符串与毫秒数两种写法7 days或5000毫秒也支持second/minute/hour/day/week/month/year单位单次最大 365 天见 api.md。若需要绝对时间点可用step.sleepUntil(deadline, Date.parse(2024-12-31))。条件分支必须基于步骤输出hasConverted来自step.do(check conversion, ...)的返回值。这是确定性的若把Date.now()之类的非确定性逻辑放在步骤外面做分支判断重放时会得出不同结果详见最佳实践章节。所有副作用发邮件都要放进step.do()否则工作流休眠后被唤醒重放时游离在步骤之外的副作用可能被重复执行。三、数据管道抽取带指数退避重试→ 转换 → 存储 → 分批入库第三个模式是经典 ETL 管道的云端实现从外部 URL 抓取原始数据失败自动重试归一化转换先写入 R2 暂存再从 R2 读回并以每批 100 条的粒度批量写入 D1。export class DataPipelineWorkflow extends WorkflowEntrypointEnv, Params { async run(event, step) { // ① extract显式配置重试策略 // 最多重试 10 次每次间隔 30s指数退避应对瞬时故障 const rawData await step.do( extract, { retries: { limit: 10, delay: 30s, backoff: exponential } }, async () { const res await fetch(event.params.sourceUrl); if (!res.ok) throw new Error(Fetch failed); // 抛出普通 Error → 触发重试 return res.json(); } ); // ② transform纯函数式归一化无副作用天然幂等 const transformed await step.do(transform, async () rawData.map(item ({ id: item.id, normalized: normalizeData(item) })) ); // ③ store大数据先落 R2只返回引用 { key } // 步骤返回值上限 1 MiB超过就必须存外部存储 const dataRef await step.do(store, async () { const key processed/${Date.now()}.json; await this.env.BUCKET.put(key, JSON.stringify(transformed)); return { key }; }); // ④ load从 R2 读回按 100 条一批批量写入 D1 await step.do(load, async () { const data await (await this.env.BUCKET.get(dataRef.key)).json(); for (let i 0; i data.length; i 100) { await this.env.DB.batch(data.slice(i, i 100).map(item this.env.DB.prepare(INSERT INTO records VALUES (?, ?)) .bind(item.id, item.normalized) )); } }); } }值得展开的三个细节重试配置参数见 configuration.mdretries.limit默认 5或Infinityretries.delay默认10000msretries.backoffconstant | linear | exponential三选一timeout单次尝试的超时时间默认 10 分钟这里场景可设为30 minutes。普通 Error 触发重试NonRetryableError立即失败从cloudflare:workers导入NonRetryableError后像参数缺失401 凭证失效这类重试无意义的错误应该直接抛出NonRetryableError避免浪费重试次数见 api.md。大于 1 MiB 的数据必须存 R2/KV 并只返回引用这是 gotchas.md 中Large Step Returns Exceeding Limit的推荐解法——把{ key: r2-object-key }作为步骤返回值后续步骤再用BUCKET.get(key)取回。此外Date.now()被用在了生成存储键上注意它必须位于step.do()内部才是安全的非确定性逻辑要包进步骤里。四、人机审批工作流挂起等待 超时自动拒绝第四个模式把等待外部事件推向极致工作流先在 D1 创建一条pending状态的审批记录然后通过waitForEvent最长挂起 48 小时等待审批响应收到approval-response事件后按approved字段分流处理如果 48 小时无人响应导致超时则进入 catch 分支把审批记录自动置为auto-rejected。export class ApprovalWorkflow extends WorkflowEntrypointEnv, Params { async run(event, step) { // ① 落库一条待审批记录 await step.do(create approval, async () await this.env.DB.prepare( INSERT INTO approvals (id, user_id, status) VALUES (?, ?, ?) ).bind(event.instanceId, event.params.userId, pending).run() ); try { // ② 挂起等待审批事件超时 48h事件载荷声明为 { approved: boolean } const approval await step.waitForEvent{ approved: boolean }( wait for approval, { event: approval-response, timeout: 48h } ); // ③ 按事件载荷做确定性分流 if (approval.approved) { await step.do(process approval, async () {}); } else { await step.do(handle rejection, async () {}); } } catch (e) { // ④ 超时兜底自动拒绝 await step.do(auto reject, async () await this.env.DB.prepare(UPDATE approvals SET status ? WHERE id ?) .bind(auto-rejected, event.instanceId).run() ); } } }要点waitForEvent默认超时 24 小时最大 365 天见 api.md超时后抛异常所以必须用 try-catch 处理这也是 gotchas.md 中waitForEvent Timeout一节的官方解法。事件通过 REST/Worker 注入外部系统可用instance.sendEvent({type: approval-response, payload: { approved: true }})见 api.md或用 REST APIPOST /accounts/{account_id}/workflows/{workflow_name}/instances/{instance_id}/events投递见 api.md。用event.instanceId作为业务主键审批记录 ID 直接复用工作流实例 ID保证唯一且可回查。实例 ID 在同一保留期内必须唯一见最佳实践。五、测试工作流Vitest Introspection API工作流是长时、有状态、会休眠的程序传统单测难以覆盖。Cloudflare 提供了基于cloudflare/vitest-pool-workers的测试方案并暴露一个Introspection API可以在测试中驱动实例运行、等待指定步骤完成、甚至 mock 步骤返回值。5.1 测试环境配置vitest.config.ts// vitest.config.ts import { defineWorkersConfig } from cloudflare/vitest-pool-workers/config; export default defineWorkersConfig({ test: { poolOptions: { workers: { wrangler: { configPath: ./wrangler.jsonc } // 复用工作流绑定配置 } } } });测试池会读取wrangler.jsonc中的workflows绑定见 configuration.md从而让测试环境里能访问env.MY_WORKFLOW。5.2 使用 Introspection API 驱动与断言import { introspectWorkflowInstance } from cloudflare:test; // 创建实例与生产环境相同的 API const instance await env.MY_WORKFLOW.create({ params: { userId: 123 } }); // 取得该实例的内省句柄 const introspector await introspectWorkflowInstance(env.MY_WORKFLOW, instance.id); // 等待指定步骤完成按步骤名 序号定位 const result await introspector.waitForStepResult({ name: fetch user, index: 0 }); // Mock 步骤行为让名为 api call 的步骤直接返回 { mocked: true } await introspector.modify(async (m) { await m.mockStepResult({ name: api call }, { mocked: true }); });使用技巧waitForStepResult适合等待真实执行到某个步骤并检查其输出modifymockStepResult适合跳过外部依赖如第三方 API、支付网关让测试快速且确定与 configuration.md 中的并行步骤、条件步骤、动态循环步骤见该文 Parallel / Conditional / Dynamic Steps 小节配合可以在测试里逐段验证整条管道。六、编排模式Fan-Out、父子工作流、竞速与定时链路除了四种业务工作流patterns.md 还给出四种通用编排手法。6.1 Fan-Out 并行处理把一个 R2 桶里的所有文件并行分派处理。每个文件用独立的步骤名process ${i}序号保证确定性用Promise.all并发执行const files await step.do(list, async () this.env.BUCKET.list()); await Promise.all(files.objects.map((file, i) step.do(process ${i}, async () processFile(await (await this.env.BUCKET.get(file.key)).arrayBuffer()) ) ));注意这里的索引i来自files步骤输出因此步骤名是确定性的不会破坏重放缓存。6.2 父-子工作流Parent-Child父工作流在某个步骤里启动子工作流实例然后继续做自己的事——子工作流异步独立运行二者互不阻塞const child await step.do(start child, async () await this.env.CHILD_WORKFLOW.create({ id: child-${event.instanceId}, // 用父实例 ID 派生保证唯一 params: { data: result.data } }) ); await step.do(other work, async () console.log(Child started: ${child.id}) );在 api.md 中可以看到官方同样采用这种从另一个工作流非阻塞启动子工作流的写法。跨脚本调用时需要在调用方wrangler.jsonc中为绑定增加script_name: other-worker指向被调方见 configuration.md。6.3 竞速模式Race对同一目标发起多条路径谁先完成用谁的结果例如快路径 API与慢路径 API并行取先返回者const winner await Promise.race([ step.do(option A, async () slowOperation()), step.do(option B, async () fastOperation()) ]);6.4 定时调度链Scheduled Workflow Chain通过scheduled入口Cron Trigger启动工作流工作流内部再叠加sleep形成日任务 → 周跟进的多级定时链路export default { async scheduled(event, env) { await env.DAILY_WORKFLOW.create({ id: daily-${event.scheduledTime}, // 确定性 ID同一时间点幂等 params: { timestamp: event.scheduledTime } }); } }; export class DailyWorkflow extends WorkflowEntrypointEnv, Params { async run(event, step) { await step.do(daily task, async () {}); await step.sleep(wait 7 days, 7 days); // 挂起 7 天不占资源 await step.do(weekly followup, async () {}); } }api.md 中列出的四种触发入口Worker fetch、Queue 消费、Cron scheduled、另一工作流都适用此模式本模式选用 Cron 入口配合event.scheduledTime生成幂等实例 ID避免重复触发产生重复实例。七、最佳实践DO 与 DONTpatterns.md 给出了 8 条要做与 8 条不要做下面逐条结合仓库文档给出依据与落地说明。✅ 应该这样做步骤保持细粒度一次step.do()只做一次 API 调用除非你能证明幂等。细粒度让重试只影响失败的那一步而不是整个阶段gotchas.md 的 Idempotency Violation 正与此相关。保证幂等采用先检查再执行或使用幂等键。典型写法见 api.md 的charge示例先查sub.charged是否已扣款已扣款直接返回旧结果。步骤名必须确定使用静态字符串或基于步骤输出的名字如process ${i}、child-${event.instanceId}。步骤名是状态缓存键动态名字如拼上Date.now()会破坏重放缓存见 gotchas.md Non-Deterministic Step Names。状态靠步骤返回值持久化不要依赖模块级或局部变量——实例休眠后内存会被回收。const total await step.do(step 1, async () 10)返回值自动持久化。永远await step.do()漏掉await会产生 fire-and-forget 行为导致步骤状态与执行结果不一致gotchas.md Missing await on step.do。条件判断要确定性只基于event.payload或步骤输出做分支。if (Date.now() deadline)这类写法应改写为const isLate await step.do(check, async () Date.now() deadline)。大数据存外部超过 1 MiB 的步骤返回值存 R2/KV只返回引用{ key: r2-object-key }。同时注意实例总状态上限为 100 MB免费/ 1 GB付费见 gotchas.md 限额表。批量创建实例多个实例用createBatch()一次创建上限 100 个幂等——已存在的 ID 会被跳过详见 api.md。❌ 不要这样做一个巨型步骤包所有事会丧失持久化与重试的粒度控制任何一步失败都要整体重来。把状态放在步骤之外模块级/局部变量在休眠后丢失必须通过步骤返回值传递。修改事件对象事件event是不可变的需要新状态就基于事件重新计算并返回而不是原地修改。在步骤外放非确定性逻辑Math.random()、Date.now()必须放进step.do()内部否则重放结果不一致。在步骤外产生副作用游离的副作用在实例重启时可能被重复执行。使用非确定性步骤名破坏缓存与去重见上文 DO #3。忽略waitForEvent超时超时会抛异常必须用 try-catch 处理否则实例直接进入errored状态。复用实例 ID同一保留期内实例 ID 必须唯一否则产生冲突推荐crypto.randomUUID()或userId-${Date.now()}之类带时间戳的组合gotchas.md Instance ID Collision。八、把模式接进真实环境配置与触发速查要让上述模式真正跑起来需要三件事wrangler.jsonc里声明 workflow 绑定、用step.do访问this.env下的其他绑定KV/D1/R2/AI、以及从外部入口触发实例。完整细节见 configuration.md 与 api.md这里给出最小可运行骨架// wrangler.jsonc摘自 configuration.md { name: my-worker, main: src/index.ts, compatibility_date: 2025-01-01, observability: { enabled: true }, workflows: [ { name: my-workflow, // Workflow 名称 binding: MY_WORKFLOW, // Env 绑定名 class_name: MyWorkflow // TS 类名 } ], limits: { cpu_ms: 300000 } // CPU 上限 5 分钟默认 30s }// 从 Worker 触发摘自 api.md export default { async fetch(req, env) { const instance await env.MY_WORKFLOW.create({ id: crypto.randomUUID(), params: { userId: user123 } }); return Response.json({ id: instance.id }); } };常用的运维命令摘自 api.mdnpx wrangler deploy npx wrangler workflows list npx wrangler workflows trigger my-workflow {userId:user123} npx wrangler workflows instances list my-workflow npx wrangler workflows instances describe my-workflow instance-id npx wrangler workflows instances pause/resume/terminate my-workflow instance-id延伸阅读配置指南 configuration.mdwrangler.jsonc 完整配置、步骤级重试/超时参数、并行/条件/动态步骤、跨脚本绑定、Pages Functions 触发API 参考 api.mdStep API 全量签名、实例创建/控制/事件注入、错误处理与NonRetryableError、类型约束、CLI 与 REST API陷阱与调试 gotchas.md11 类常见错误定位、免费/付费限额对照表、计费说明Cloudflare Workflows 总览 README.md核心概念与快速上手示例。本文是cloudflare-deploy技能包中 Workflows 参考文档 的实战延伸在 skills 仓库的技能决策树中Workflows 对应长时多步任务场景见 SKILL.md 中Long-running multi-step jobs → workflows/。当你需要状态协调/多步骤编排之外的备选方案时可对比阅读 durable-objects有状态实体与 queues消息驱动两类相邻参考。【免费下载链接】skillsSkills Catalog for Codex项目地址: https://gitcode.com/GitHub_Trending/skills4/skills创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考