☰
Cloudflare Workflows 实战指南:Step API、实例管理与多场景触发编排(autoskills cloudflare-deploy 技能参考)
2026/10/10 1:29:48 网站建设 项目流程

【免费下载链接】autoskills

One command. Your entire AI skill stack. Installed.

项目地址:https://gitcode.com/gh_mirrors/au/autoskills
点击查看免费下载

本篇以 autoskills 仓库cloudflare-deploy技能(SKILL.md)中 Workflows 参考文档为主体,系统讲解 Cloudflare Workflows 的完整 API 面:step.do/sleep/waitForEvent等 Step API、实例的创建/查询/控制/事件投递、Worker/Queue/Cron/Workflow 间等多种触发方式、错误处理与幂等设计、Wrangler CLI 与 REST API 操作。读完你将能独立编写、配置、部署、触发与调试一个具备自动重试、状态持久化与长时运行能力的多步骤工作流,并了解其配额限制与常见陷阱。

Cloudflare Workflows 是一种持久化的多步骤应用运行时:步骤可自动重试、步骤间状态持久化(分钟级到数周级)、失败不丢失进度、可等待外部事件/人工审批、可在不消耗资源的情况下休眠。该参考位于 references/workflows/ 目录,由 README.md(核心概念与入门)、configuration.md(配置)、api.md(API 速查,本文主体)、patterns.md(实战模式)与 gotchas.md(陷阱与限额)五份文档组成。Workflows 在 Free 与 Paid Workers 套餐上均可用。

核心概念

  • Workflow:继承WorkflowEntrypoint并实现run方法的类,定义工作流逻辑;
  • Instance(实例):一次独立的执行,拥有唯一 ID 与独立状态;
  • Step(步骤):通过step.do()定义的可独立重试的最小单元——API 调用、数据库查询、AI 调用等;
  • State(状态):由步骤返回值持久化,步骤名即状态缓存键。

最小可运行示例(来自 README.md):

import { WorkflowEntrypoint, WorkflowStep, WorkflowEvent } from 'cloudflare:workers'; type Env = { MY_WORKFLOW: Workflow; DB: D1Database }; type Params = { userId: string }; export class MyWorkflow extends WorkflowEntrypoint<Env, Params> { async run(event: WorkflowEvent<Params>, step: WorkflowStep) { const user = await step.do('fetch user', async () => { return await this.env.DB.prepare('SELECT * FROM users WHERE id = ?') .bind(event.params.userId).first(); }); await step.sleep('wait 7 days', '7 days'); await step.do('send reminder', async () => { await sendEmail(user.email, 'Reminder!'); }); } }

Step API 详解

Step API 是工作流的执行原语,全部来自 api.md:

// step.do() const result = await step.do('step name', async () => { /* logic */ }); const result = await step.do('step name', { retries, timeout }, async () => {}); // step.sleep() await step.sleep('description', '1 hour'); await step.sleep('description', 5000); // ms // step.sleepUntil() await step.sleepUntil('description', Date.parse('2024-12-31')); // step.waitForEvent() const data = await step.waitForEvent<PayloadType>('wait', {event: 'webhook-type', timeout: '24h'}); // Default 24h, max 365d try { const event = await step.waitForEvent('wait', { event: 'approval', timeout: '1h' }); } catch (e) { /* Timeout */ }

step.do 的进阶配置

step.do的第二参数支持重试与超时配置(来自 configuration.md):

await step.do('api call', { retries: { limit: 10, // 默认:5(或 Infinity) delay: '10 seconds', // 默认:10000ms backoff: 'exponential' // constant | linear | exponential }, timeout: '30 minutes' // 单次尝试超时(默认 10 分钟) }, async () => { const res = await fetch('https://api.example.com/data'); if (!res.ok) throw new Error('Failed'); return res.json(); });

并行、条件与循环步骤

  • 并行步骤:用Promise.all并发执行多个step.do:
const [user, settings] = await Promise.all([ step.do('fetch user', async () => this.env.KV.get(`user:${id}`)), step.do('fetch settings', async () => this.env.KV.get(`settings:${id}`)) ]);
  • 条件步骤:分支必须基于步骤输出等确定性来源,禁止在步骤外使用Date.now():
const config = await step.do('fetch config', async () => this.env.KV.get('flags', { type: 'json' }) ); // ✅ 确定性分支(基于步骤输出) if (config.enableEmail) { await step.do('send email', async () => sendEmail()); } // ❌ 非确定性分支(步骤外使用 Date.now) if (Date.now() > deadline) { /* BAD */ }
  • 动态步骤(循环):遍历列表逐个执行步骤:
const files = await step.do('list files', async () => this.env.BUCKET.list() ); for (const file of files.objects) { await step.do(`process ${file.key}`, async () => { const obj = await this.env.BUCKET.get(file.key); return processData(await obj.arrayBuffer()); }); }

实例管理

每个实例是一次独立执行,支持创建、批量创建、查询状态与生命周期控制(api.md):

// 创建单个实例 const instance = await env.MY_WORKFLOW.create({id: crypto.randomUUID(), params: { userId: 'user123' }}); // id 可选,省略时自动生成 // 自定义保留期创建(默认:免费版 3 天,付费版 30 天) const instance = await env.MY_WORKFLOW.create({ id: crypto.randomUUID(), params: { userId: 'user123' }, retention: '30 days' // 覆盖默认保留期 }); // 批量创建(最多 100 个,幂等:跳过已存在的 ID) const instances = await env.MY_WORKFLOW.createBatch([{id: 'user1', params: {name: 'John'}}, {id: 'user2', params: {name: 'Jane'}}]); // 获取实例与状态 const instance = await env.MY_WORKFLOW.get('instance-id'); const status = await instance.status(); // {status: 'queued' | 'running' | 'paused' | 'errored' | 'terminated' | 'complete' | 'waiting' | 'waitingForPause' | 'unknown', error?, output?} // 生命周期控制 await instance.pause(); await instance.resume(); await instance.terminate(); await instance.restart(); // 发送事件(type 必须与 waitForEvent 中声明的类型匹配) await instance.sendEvent({type: 'approval', payload: { approved: true }});

需要注意:status()返回的状态枚举覆盖排队、运行、暂停、出错、终止、完成、等待等全部生命周期阶段,且可能携带error与output字段;sendEvent投递的事件类型必须与工作流内waitForEvent等待的事件类型一致,否则不会被消费。

触发 Workflows 的方式

Workflows 可从多种入口触发(api.md):

从 Worker(HTTP 入口)触发:

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 }); }};

从 Queue 消费消息触发:

export default { async queue(batch, env) { for (const msg of batch.messages) { await env.MY_WORKFLOW.create({id: `job-${msg.id}`, params: msg.body}); } }};

从 Cron(定时调度)触发:

export default { async scheduled(event, env) { await env.CLEANUP_WORKFLOW.create({id: `cleanup-${Date.now()}`, params: { timestamp: event.scheduledTime }}); }};

从另一个 Workflow 触发(非阻塞):

export class ParentWorkflow extends WorkflowEntrypoint<Env, Params> { async run(event, step) { const child = await step.do('start child', async () => await this.env.CHILD_WORKFLOW.create({id: `child-${event.instanceId}`, params: {}})); } }

从 Pages Functions 触发(通过 service bindings,配置于wrangler.jsonc的service_bindings下,见 configuration.md):

// functions/_middleware.ts export const onRequest: PagesFunction<Env> = async ({ env, request }) => { const instance = await env.MY_WORKFLOW.create({ params: { url: request.url } }); return new Response(`Started ${instance.id}`); };

错误处理与幂等

NonRetryableError:对于校验失败、凭据无效等重试无意义的错误,抛出NonRetryableError可让该步骤直接失败而不消耗重试次数;普通Error则会触发重试:

import { NonRetryableError } from 'cloudflare:workers'; await step.do('validate', async () => { if (!event.params.paymentMethod) throw new NonRetryableError('Payment method required'); const res = await fetch('https://api.example.com/charge', { method: 'POST' }); if (res.status === 401) throw new NonRetryableError('Invalid credentials'); // 不重试 if (!res.ok) throw new Error('Retryable failure'); // 会重试 return res.json(); });

捕获错误并执行兜底:waitForEvent超时、非重试错误等都会抛出,可用 try-catch 包裹后执行清理步骤:

try { await step.do('risky op', async () => { throw new NonRetryableError('Failed'); }); } catch (e) { await step.do('cleanup', async () => {}); }

幂等性(check-then-execute):重试可能导致重复执行副作用操作,因此在执行前先检查是否已完成:

await step.do('charge', async () => { const sub = await fetch(`https://api/subscriptions/${id}`).then(r => r.json()); if (sub.charged) return sub; // 已完成,直接返回 return await fetch(`https://api/subscriptions/${id}`, {method: 'POST', body: JSON.stringify({ amount: 10.0 })}).then(r => r.json()); });

类型约束:Rpc.Serializable

Params 与步骤返回值必须是Rpc.Serializable<T>类型(api.md):

// ✅ 合法类型 type ValidParams = { userId: string; count: number; tags: string[]; metadata: Record<string, unknown>; }; // ❌ 非法类型 type InvalidParams = { callback: () => void; // 函数不可序列化 symbol: symbol; // Symbol 不可序列化 circular: any; // 不允许循环引用 }; // 步骤返回值遵循同样规则 const result = await step.do('fetch', async () => { return { userId: '123', data: [1, 2, 3] }; // ✅ 纯对象 });

休眠与调度

支持相对时间(sleep)与绝对时间(sleepUntil)两种调度方式:

// 相对时间 await step.sleep('wait 1 hour', '1 hour'); await step.sleep('wait 30 days', '30 days'); await step.sleep('wait 5s', 5000); // 毫秒 // 绝对时间 await step.sleepUntil('launch date', Date.parse('24 Oct 2024 13:00:00 UTC')); await step.sleepUntil('deadline', new Date('2024-12-31T23:59:59Z'));

时间单位支持:second、minute、hour、day、week、month、year,单次休眠最长 365 天。处于休眠状态的实例不占用并发额度,因此可以支撑数百万个同时"沉睡"的工作流实例(详见 gotchas.md 的并发限制说明)。

参数传递

从 Worker 传入参数:

const instance = await env.MY_WORKFLOW.create({ id: crypto.randomUUID(), params: { userId: 'user123', email: 'user@example.com' } });

在工作流内访问:

async run(event: WorkflowEvent<Params>, step: WorkflowStep) { const userId = event.params.userId; const instanceId = event.instanceId; const createdAt = event.timestamp; }

WorkflowEvent除params外还暴露instanceId与timestamp,可用于生成确定性步骤名与审计。

CLI 触发并传参:

npx wrangler workflows trigger my-workflow '{"userId":"user123"}'

Wrangler CLI 命令

完整命令链(api.md):

npm create cloudflare@latest my-workflow -- --template "cloudflare/workflows-starter" npx 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

部署前建议先执行npx wrangler whoami确认已认证(交互式使用wrangler login,CI/CD 使用CLOUDFLARE_API_TOKEN环境变量,详见 SKILL.md)。

REST API

工作流实例也可通过 Cloudflare API v4 管理(api.md):

# 创建实例 curl -X POST "https://api.cloudflare.com/client/v4/accounts/{account_id}/workflows/{workflow_name}/instances" -H "Authorization: Bearer {token}" -d '{"id":"custom-id","params":{"userId":"user123"}}' # 查询状态 curl "https://api.cloudflare.com/client/v4/accounts/{account_id}/workflows/{workflow_name}/instances/{instance_id}/status" -H "Authorization: Bearer {token}" # 发送事件 curl -X POST "https://api.cloudflare.com/client/v4/accounts/{account_id}/workflows/{workflow_name}/instances/{instance_id}/events" -H "Authorization: Bearer {token}" -d '{"type":"approval","payload":{"approved":true}}'

配置详解:wrangler.jsonc 与绑定

wrangler.jsonc 基础配置

在wrangler.jsonc中声明工作流(configuration.md):

{ "name": "my-worker", "main": "src/index.ts", "compatibility_date": "2025-01-01", // 新项目使用当前日期 "observability": { "enabled": true // 启用 Workflows 仪表盘与结构化日志 }, "workflows": [ { "name": "my-workflow", // 工作流名称 "binding": "MY_WORKFLOW", // Env 绑定名 "class_name": "MyWorkflow" // TS 类名 // "script_name": "other-worker" // 跨脚本调用时使用 } ], "limits": { "cpu_ms": 300000 // CPU 上限 5 分钟(默认 30s) } }

多个工作流与跨脚本绑定

一个 Worker 可声明多个工作流,每个类继承WorkflowEntrypoint并拥有自己的Params类型:

{ "workflows": [ {"name": "user-onboarding", "binding": "USER_ONBOARDING", "class_name": "UserOnboarding"}, {"name": "data-processing", "binding": "DATA_PROCESSING", "class_name": "DataProcessing"} ] }

跨脚本调用:Worker B 通过script_name指向定义工作流的 Worker A:

// Worker B(调用方) { "workflows": [{ "name": "billing-workflow", "binding": "BILLING", "script_name": "billing-worker" // 指向 Worker A }] }

通过 this.env 访问绑定

工作流内通过this.env访问 KV、D1、R2、AI、Vectorize 等 Cloudflare 绑定:

type Env = { MY_WORKFLOW: Workflow; KV: KVNamespace; DB: D1Database; BUCKET: R2Bucket; AI: Ai; VECTORIZE: VectorizeIndex; }; await step.do('use bindings', async () => { const kv = await this.env.KV.get('key'); const db = await this.env.DB.prepare('SELECT * FROM users').first(); const file = await this.env.BUCKET.get('file.txt'); const ai = await this.env.AI.run('@cf/meta/llama-2-7b-chat-int8', { prompt: 'Hi' }); });

实战模式:四类典型工作流

patterns.md 给出了可直接套用的完整示例:

图片处理流水线(获取 → AI 生成描述 → 等待审批 → 发布):

export class ImageProcessingWorkflow extends WorkflowEntrypoint<Env, Params> { async run(event, step) { const imageData = await step.do('fetch', async () => (await this.env.BUCKET.get(event.params.imageKey)).arrayBuffer()); 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}) ); await step.waitForEvent('await approval', { event: 'approved', timeout: '24h' }); await step.do('publish', async () => await this.env.BUCKET.put(`public/${event.params.imageKey}`, imageData)); } }

用户生命周期(欢迎邮件 → 试用期 7 天 → 转化检查 → 到期提醒):

export class UserLifecycleWorkflow extends WorkflowEntrypoint<Env, Params> { async run(event, step) { await step.do('welcome email', async () => await sendEmail(event.params.email, 'Welcome!')); await step.sleep('trial period', '7 days'); 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')); } }

数据管道(提取 → 转换 → 存储 → 分批装载):

export class DataPipelineWorkflow extends WorkflowEntrypoint<Env, Params> { async run(event, step) { 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'); return res.json(); }); const transformed = await step.do('transform', async () => rawData.map(item => ({ id: item.id, normalized: normalizeData(item) })) ); const dataRef = await step.do('store', async () => { const key = `processed/${Date.now()}.json`; await this.env.BUCKET.put(key, JSON.stringify(transformed)); return { key }; }); 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) )); } }); } }

人工审批(Human-in-the-Loop):创建审批记录 → 等待审批事件(48h)→ 通过/拒绝/超时自动拒绝:

export class ApprovalWorkflow extends WorkflowEntrypoint<Env, 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 { 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()); } } }

编排模式

  • Fan-Out(扇出并行):Promise.all并行处理多个对象;
  • Parent-Child(父子工作流):父工作流在step.do中启动子工作流(非阻塞),随后继续其他工作;
  • Race(竞速):Promise.race取最先完成的步骤结果;
  • 调度链(Scheduled Chain):Cron 触发后,工作流内用step.sleep串联"日任务 → 7 天后的周任务"。

测试工作流

使用@cloudflare/vitest-pool-workers编写测试(patterns.md):

// vitest.config.ts import { defineWorkersConfig } from '@cloudflare/vitest-pool-workers/config'; export default defineWorkersConfig({ test: { poolOptions: { workers: { wrangler: { configPath: './wrangler.jsonc' } } } } });

配合cloudflare:test的Introspection API可以等待步骤完成甚至 mock 步骤行为:

import { introspectWorkflowInstance } from 'cloudflare:test'; 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 步骤行为 await introspector.modify(async (m) => { await m.mockStepResult({ name: 'api call' }, { mocked: true }); });

最佳实践清单

✅ 应该做(来自 patterns.md):

  1. 步骤粒度化:一次 API 调用一个步骤(除非要证明幂等);
  2. 幂等设计:先检查再执行,使用幂等键;
  3. 确定性步骤名:使用静态名或基于步骤输出的名字;
  4. 用返回值持久化状态:通过步骤返回值保存状态,而不是局部变量;
  5. 始终 await:await step.do(),避免悬空 Promise;
  6. 确定性条件分支:基于event.payload或步骤输出判断;
  7. 大对象外置存储:超过 1 MiB 的数据放入 R2/KV,仅返回引用;
  8. 批量创建:多个实例用createBatch()。

❌ 不要做:

  1. 写一个巨型步骤(破坏持久化与重试控制);
  2. 在步骤外保存状态(休眠时会丢失);
  3. 修改事件(事件不可变,应返回新状态);
  4. 在步骤外使用非确定性逻辑(Math.random()、Date.now()必须在步骤内);
  5. 在步骤外执行副作用(重启后可能重复执行);
  6. 使用非确定性步骤名(破坏缓存);
  7. 忽略超时(waitForEvent会抛出,需 try-catch);
  8. 复用实例 ID(保留期内必须唯一)。

常见陷阱与限制

gotchas.md 总结了高频错误与对策:

错误原因解决方案
Step Timeout步骤超过默认 10 分钟超时用{timeout: '30 minutes'}自定义,或在 wrangler.jsonc 提高 CPU 上限(最大 5 分钟 CPU 时间)
waitForEvent Timeout超时周期内未收到事件(默认 24h,最大 365d)try-catch 包裹,超时后走默认行为
非确定性步骤名步骤名中使用Date.now()等动态值使用event.instanceId等确定性值
状态丢失用模块级/局部变量存状态,休眠后丢失从step.do()返回值,会自动持久化
非确定性条件步骤外用Date.now()判断把非确定性操作移入步骤:const isLate = await step.do('check', async () => Date.now() > deadline)
步骤返回值超限步骤返回超过 1 MiB 的数据大数据存 R2,只返回{ key: 'r2-object-key' }引用
CPU 超限但运行 <30s混淆 CPU 时间与墙钟时间网络请求、数据库查询、休眠不计入 CPU;30s 指活跃计算时间
幂等性违规重试导致重复扣款/重复动作执行前先检查是否已完成
实例 ID 冲突复用实例 ID用时间戳保证唯一:${userId}-${Date.now()}
完成后实例数据消失完成/出错实例在保留期后被自动删除(免费 3 天/付费 30 天)完成前将关键数据导出到 KV/R2/D1
漏写 awaitstep.do()未 await 造成 fire-and-forget始终await step.do('task', ...)

限额速查表

限制项FreePaid说明
每步骤 CPU10ms30s(默认),5min(最大)通过 wrangler.jsonc 的limits.cpu_ms设置
步骤状态1 MiB1 MiB单步返回值
实例状态100 MB1 GB单实例全部状态
每工作流步骤数1,0241,024step.sleep()不计入
每日执行数100k无限每日执行上限
并发实例2510k等待状态的实例不计入
排队实例100k1M最大排队实例数
每步子请求501,000每步骤最大出站请求
状态保留期3 天30 天已完成实例保留时长
步骤默认超时10 min10 min单次尝试
waitForEvent 默认超时24h24h最大 365 天
waitForEvent 最大超时365 天365 天最长等待时间

关键提示:处于waiting状态(step.sleep或step.waitForEvent期间)的实例不计入并发实例上限,因此可以支撑数百万个休眠中的工作流。

计费参考

指标FreePaid说明
请求100k/天10M/月 + $0.30/M工作流调用
CPU 时间10ms/次30M CPU-ms/月 + $0.02/M CPU-ms实际 CPU 用量
存储1 GB1 GB/月 + $0.20/GB-月所有实例(运行/出错/休眠/完成)

阅读建议

  • 入门顺序:先读 configuration.md(配置)→ api.md(API)→ patterns.md(模式);
  • 排障:查阅 gotchas.md;
  • 相关技术选型:有状态协调/实时场景可对比 durable-objects,消息驱动场景参考 queues,工作流入口本身运行在 workers 之上;部署与认证流程见 SKILL.md。

【免费下载链接】autoskills

One command. Your entire AI skill stack. Installed.

项目地址:https://gitcode.com/gh_mirrors/au/autoskills
点击查看免费下载

相关推荐

上一篇:GraphiQL合作伙伴:生态系统建设与商业合作
下一篇:5分钟部署游戏串流服务器:Sunshine自动化工具链全攻略

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

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

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

立即咨询