分支与并行执行
线性工作流 Step 一个接一个往下走,能解决很多问题,但真实业务往往更复杂。订单金额不同走不同支付通道,一笔交易要同时查风控和库存,用户注册后要等外部回调确认。这些场景需要条件分支、并行执行和长时间等待三种能力。这篇讲清楚怎么在 Cloudflare Workflows 里用标准 JavaScript 控制流实现这三种模式。
条件分支
工作流的 run 方法就是普通 JavaScript,if、else、switch 这些控制流原语都能直接用。分支依据是前面 Step 的返回值,引擎会按实际走到的分支持久化对应 Step。
分支模式对比
| 分支方式 | 写法 | 适用场景 |
|---|---|---|
| if/else | 根据条件走不同 Step | 二选一或少量分支 |
| switch | 多分支匹配 | 枚举值多选一 |
| 循环 | for/while 配合 step.do | 批量处理同类任务 |
| try/catch | 捕获 Step 异常分流 | 降级处理 |
一个按金额分流支付的例子
import { WorkflowEntrypoint, WorkflowStep } from "cloudflare:workers";
import type { WorkflowEvent } from "cloudflare:workers";
type Params = { orderId: string };
export class PaymentWorkflow extends WorkflowEntrypoint<Env, Params> {
async run(event: WorkflowEvent<Params>, step: WorkflowStep) {
const order = await step.do("load order", async () => {
const row = await this.env.DB
.prepare("SELECT * FROM orders WHERE id = ?")
.bind(event.payload.orderId)
.first();
return row;
});
// 按金额分流
let result;
if (order.total < 100) {
// 小额走免密支付
result = await step.do("小额支付", async () => {
const res = await fetch("https://api.payment.com/quick", {
method: "POST",
body: JSON.stringify({ amount: order.total }),
});
return await res.json();
});
} else if (order.total < 10000) {
// 中额走标准支付
result = await step.do("标准支付", async () => {
const res = await fetch("https://api.payment.com/standard", {
method: "POST",
body: JSON.stringify({ amount: order.total }),
});
return await res.json();
});
} else {
// 大额走人工审核
result = await step.do("大额审核", async () => {
return { status: "pending_review", amount: order.total };
});
}
return { orderId: order.id, payment: result };
}
}
分支里每个 Step 的名字是固定的字符串,重启时引擎按实际走到的分支恢复,没走到的分支不会执行。
分支使用的注意事项
| 事项 | 说明 |
|---|---|
| Step 名字确定性 | 不能用运行时随机值当名字,否则存档失效 |
| 分支条件来源 | 必须来自已持久化的 Step 返回值或 event.payload |
| 分支内 Step 数量 | 各分支可以不等长 |
| 嵌套分支 | 可以,但层数多了可读性差,建议拆 Step |
循环处理批量任务也很常见。注意循环里 Step 名字要带稳定索引,不能重复。
export class BatchWorkflow extends WorkflowEntrypoint<Env, Params> {
async run(event: WorkflowEvent<Params>, step: WorkflowStep) {
const tasks = await step.do("load tasks", async () => {
return [{ id: 1 }, { id: 2 }, { id: 3 }];
});
const results = [];
for (const task of tasks) {
// 名字带索引,确保唯一且确定
const r = await step.do(`处理任务 ${task.id}`, async () => {
return { taskId: task.id, done: true };
});
results.push(r);
}
return { count: results.length };
}
}
并行执行 Step
多个互不依赖的 Step 可以并行跑,用 Promise.all 包起来。引擎会同时执行这些 Step,各自独立持久化,全部完成后才继续往下走。
串行与并行对比
| 模式 | 写法 | 总耗时 | 依赖关系 |
|---|---|---|---|
| 串行 | await 一个接一个 | 各 Step 时长之和 | 后者依赖前者 |
| 并行 | Promise.all 包多个 | 最慢那个的时长 | 互不依赖 |
一个并行查风控和库存的例子
export class CheckoutWorkflow extends WorkflowEntrypoint<Env, Params> {
async run(event: WorkflowEvent<Params>, step: WorkflowStep) {
const { userId, productId } = event.payload;
// 两个检查互不依赖,并行跑
const [risk, stock] = await Promise.all([
step.do("风控检查", async () => {
const res = await fetch(`https://risk.example.com/check?user=${userId}`);
return await res.json();
}),
step.do("库存检查", async () => {
const row = await this.env.DB
.prepare("SELECT stock FROM products WHERE id = ?")
.bind(productId)
.first();
return { stock: row?.stock ?? 0 };
}),
]);
// 两个结果都到位后判断
const canProceed = await step.do("综合判断", async () => {
return risk.pass && stock.stock > 0;
});
return { canProceed };
}
}
如果串行跑,风控检查 200 毫秒加库存检查 100 毫秒总共 300 毫秒。并行后只要 200 毫秒,节省三分之一。
并行使用要点
| 要点 | 说明 |
|---|---|
| 仅限互不依赖 | 后一个 Step 用到前一个的返回值就必须串行 |
| Step 名字仍要唯一 | 并行的 Step 名字不能重复 |
| 下游承受能力 | 并行数别超过下游接口或数据库连接上限 |
| 错误处理 | 任一并行 Step 失败,Promise.all 整体拒绝,按重试策略走 |
并行不是越多越好。下游接口有并发限制时,太多并行请求会触发限流反而更慢。一般控制在 3 到 5 个并行 Step。
长时间运行
工作流可以跑很长时间,从几小时到几天甚至更长。这靠的是休眠和等待事件两个能力,休眠期间不占计算资源。
长时间运行能力一览
| 能力 | 方法 | 典型时长 |
|---|---|---|
| 相对休眠 | step.sleep | 几分钟到几天 |
| 定点等待 | step.sleepUntil | 到某个未来时刻 |
| 等待外部事件 | step.waitForEvent | 默认 24 小时 |
休眠用 step.sleep 暂停一段时长,适合固定等待。这里重点讲等待外部事件,它解决的是暂停等人确认的场景,比如审批流、支付回调、用户验证。step.waitForEvent 让工作流暂停,直到收到匹配的事件或超时。
export class ApprovalWorkflow extends WorkflowEntrypoint<Env, Params> {
async run(event: WorkflowEvent<Params>, step: WorkflowStep) {
const { requestId } = event.payload;
await step.do("提交审批", async () => {
return { requestId, status: "pending" };
});
// 暂停等待审批结果,最多等 4 小时
const approval = await step.waitForEvent("等待审批", {
type: "approval-result",
timeout: "4 hours",
});
const result = await step.do("记录结果", async () => {
return {
requestId,
approved: approval.approved,
reviewer: approval.reviewer,
};
});
return result;
}
}
waitForEvent 接收一个配置对象,type 用来匹配外部发来的事件类型,timeout 控制最多等多久,默认 24 小时。事件没在超时前到达会抛超时异常,按重试策略处理。
事件由外部通过实例的 sendEvent 方法发送。通常在另一个 Worker 的 fetch 处理函数里调用。
export default {
async fetch(request: Request, env: Env): Promise<Response> {
const body = await request.json() as {
instanceId: string;
approved: boolean;
reviewer: string;
};
// 取回实例引用
const instance = await env.MY_WORKFLOW.get(body.instanceId);
// 发送事件,type 必须和 waitForEvent 匹配
await instance.sendEvent({
type: "approval-result",
payload: {
approved: body.approved,
reviewer: body.reviewer,
},
});
return Response.json({ ok: true });
},
} satisfies ExportedHandler<Env>;
waitForEvent 关键点
| 要点 | 说明 |
|---|---|
| type 匹配 | sendEvent 的 type 必须和 waitForEvent 完全一致 |
| 返回值 | 收到事件后返回 payload 数据 |
| 超时 | 超时抛异常,走重试或被捕获 |
| 暂停不占资源 | 等待期间实例休眠,不消耗计算 |
| 多事件等待 | 可用 Promise.race 同时等多个事件 |
组合模式
把分支、并行、等待事件组合起来,能处理很复杂的业务。下面是一个完整的订单履约工作流,并行查库存和风控,按结果分支,等待支付回调。
import { WorkflowEntrypoint, WorkflowStep } from "cloudflare:workers";
import type { WorkflowEvent } from "cloudflare:workers";
type Params = { orderId: string; userId: string };
export class FulfillmentWorkflow extends WorkflowEntrypoint<Env, Params> {
async run(event: WorkflowEvent<Params>, step: WorkflowStep) {
const { orderId, userId } = event.payload;
// 并行查库存和风控
const [stock, risk] = await Promise.all([
step.do("查库存", async () => {
const row = await this.env.DB
.prepare("SELECT stock FROM products WHERE order_id = ?")
.bind(orderId)
.first();
return { available: (row?.stock ?? 0) > 0 };
}),
step.do("查风控", async () => {
const res = await fetch(`https://risk.example.com/check?user=${userId}`);
return await res.json();
}),
]);
// 条件分支,风控不通过直接拒绝
if (!risk.pass) {
const rejected = await step.do("记录拒绝", async () => {
return { orderId, status: "rejected", reason: "risk" };
});
return rejected;
}
// 库存不足走补货流程
if (!stock.available) {
await step.do("触发补货", async () => {
return { orderId, status: "restocking" };
});
await step.sleep("等补货", "1 hour");
}
// 等待支付回调,最多 30 分钟
const payment = await step.waitForEvent("等待支付", {
type: "payment-callback",
timeout: "30 minutes",
});
// 支付成功则发货
if (payment.success) {
const shipped = await step.do("发货", async () => {
return { orderId, status: "shipped", tracking: "TRK123" };
});
return shipped;
}
// 支付失败取消订单
const cancelled = await step.do("取消订单", async () => {
return { orderId, status: "cancelled" };
});
return cancelled;
}
}
这个例子串起了全部三种模式。Promise.all 并行检查,if 分流处理不同结果,waitForEvent 暂停等外部回调,最后再分支决定发货还是取消。
复杂编排的拆分建议
| 场景 | 建议做法 |
|---|---|
| 单工作流超过 50 步 | 拆成多个子工作流,各自独立编排 |
| 分支层级超过 3 层 | 把深层逻辑挪进单独 Step,扁平化处理 |
| 并行 Step 超过 5 个 | 分批并行,或评估下游承受能力 |
| 等待事件超过 3 个 | 用 Promise.race 聚合,避免串行等待 |
拆分的核心标准是可读性和可重试性。一个工作流如果失败后从头理解要花很久,就该拆。每个子工作流聚焦一个职责,失败重试的影响范围也小。
小结
条件分支直接用 if/else 和 switch,依据是 Step 返回值,名字保持确定性。并行用 Promise.all 包多个互不依赖的 Step,总耗时取决于最慢那个。长时间运行靠 sleep 和 waitForEvent,休眠和等待期间不占资源,能跑几小时几天。三种模式组合起来能覆盖订单履约、审批流、支付回调这类复杂业务。拆分时以可读性和可重试性为准,单工作流别贪大。
上一篇 状态与错误处理